multiprocessing --- 基于進(jìn)程的并行?
概述?
multiprocessing 是一個用與 threading 模塊相似API的支持產(chǎn)生進(jìn)程的包。 multiprocessing 包同時提供本地和遠(yuǎn)程并發(fā),使用子進(jìn)程代替線程,有效避免 Global Interpreter Lock 帶來的影響。因此, multiprocessing 模塊允許程序員充分利用機(jī)器上的多個核心。Unix 和 Windows 上都可以運行。
multiprocessing 模塊還引入了在 threading 模塊中沒有類似物的API。這方面的一個主要例子是 Pool 對象,它提供了一種方便的方法,可以跨多個輸入值并行化函數(shù)的執(zhí)行,跨進(jìn)程分配輸入數(shù)據(jù)(數(shù)據(jù)并行)。以下示例演示了在模塊中定義此類函數(shù)的常見做法,以便子進(jìn)程可以成功導(dǎo)入該模塊。這個數(shù)據(jù)并行的基本例子使用 Pool ,
from multiprocessing import Pool
def f(x):
return x*x
if __name__ == '__main__':
with Pool(5) as p:
print(p.map(f, [1, 2, 3]))
將打印到標(biāo)準(zhǔn)輸出
[1, 4, 9]
Process 類?
在 multiprocessing 中,通過創(chuàng)建一個 Process 對象然后調(diào)用它的 start() 方法來生成進(jìn)程。 Process 和 threading.Thread API 相同。 一個簡單的多進(jìn)程程序示例是:
from multiprocessing import Process
def f(name):
print('hello', name)
if __name__ == '__main__':
p = Process(target=f, args=('bob',))
p.start()
p.join()
要顯示所涉及的各個進(jìn)程ID,這是一個擴(kuò)展示例:
from multiprocessing import Process
import os
def info(title):
print(title)
print('module name:', __name__)
print('parent process:', os.getppid())
print('process id:', os.getpid())
def f(name):
info('function f')
print('hello', name)
if __name__ == '__main__':
info('main line')
p = Process(target=f, args=('bob',))
p.start()
p.join()
為了解釋為什么 if __name__ == '__main__' 部分是必需的,請參見 編程指導(dǎo)。
上下文和啟動方法?
根據(jù)不同的平臺, multiprocessing 支持三種啟動進(jìn)程的方法。這些 啟動方法 有
- spawn
父進(jìn)程啟動一個新的Python解釋器進(jìn)程。子進(jìn)程只會繼承那些運行進(jìn)程對象的
run()方法所需的資源。特別是父進(jìn)程中非必須的文件描述符和句柄不會被繼承。相對于使用 fork 或者 forkserver,使用這個方法啟動進(jìn)程相當(dāng)慢。可在Unix和Windows上使用。 Windows上的默認(rèn)設(shè)置。
- fork
父進(jìn)程使用
os.fork()來產(chǎn)生 Python 解釋器分叉。子進(jìn)程在開始時實際上與父進(jìn)程相同。父進(jìn)程的所有資源都由子進(jìn)程繼承。請注意,安全分叉多線程進(jìn)程是棘手的。只存在于Unix。Unix中的默認(rèn)值。
- forkserver
程序啟動并選擇* forkserver * 啟動方法時,將啟動服務(wù)器進(jìn)程。從那時起,每當(dāng)需要一個新進(jìn)程時,父進(jìn)程就會連接到服務(wù)器并請求它分叉一個新進(jìn)程。分叉服務(wù)器進(jìn)程是單線程的,因此使用
os.fork()是安全的。沒有不必要的資源被繼承。可在Unix平臺上使用,支持通過Unix管道傳遞文件描述符。
在 3.4 版更改: spawn 在所有unix平臺上添加,并且為一些unix平臺添加了 forkserver 。子進(jìn)程不再繼承Windows上的所有上級進(jìn)程可繼承的句柄。
在Unix上使用 spawn 或 forkserver 啟動方法也將啟動一個 信號量跟蹤器 進(jìn)程,該進(jìn)程跟蹤由程序進(jìn)程創(chuàng)建的未鏈接的命名信號量。當(dāng)所有進(jìn)程退出時,信號量跟蹤器取消鏈接任何剩余的信號量。通常不應(yīng)該有,但如果一個進(jìn)程被信號殺死,可能會有一些“泄露”的信號量。(取消鏈接命名的信號量是一個嚴(yán)重的問題,因為系統(tǒng)只允許有限的數(shù)量,并且在下次重新啟動之前它們不會自動取消鏈接。)
要選擇一個啟動方法,你應(yīng)該在主模塊的 if __name__ == '__main__' 子句中調(diào)用 set_start_method() 。例如:
import multiprocessing as mp
def foo(q):
q.put('hello')
if __name__ == '__main__':
mp.set_start_method('spawn')
q = mp.Queue()
p = mp.Process(target=foo, args=(q,))
p.start()
print(q.get())
p.join()
在程序中 set_start_method() 不應(yīng)該被多次調(diào)用。
或者,你可以使用 get_context() 來獲取上下文對象。上下文對象與多處理模塊具有相同的API,并允許在同一程序中使用多個啟動方法。:
import multiprocessing as mp
def foo(q):
q.put('hello')
if __name__ == '__main__':
ctx = mp.get_context('spawn')
q = ctx.Queue()
p = ctx.Process(target=foo, args=(q,))
p.start()
print(q.get())
p.join()
請注意,對象在不同上下文創(chuàng)建的進(jìn)程間可能并不兼容。 特別是,使用 fork 上下文創(chuàng)建的鎖不能傳遞給使用 spawn 或 forkserver 啟動方法啟動的進(jìn)程。
想要使用特定啟動方法的庫應(yīng)該使用 get_context() 以避免干擾庫用戶的選擇。
警告
'spawn' 和 'forkserver' 啟動方法當(dāng)前不能在Unix上和“凍結(jié)的”可執(zhí)行內(nèi)容一同使用(例如,有類似 PyInstaller 和 cx_Freeze 的包產(chǎn)生的二進(jìn)制文件)。 'fork' 啟動方法可以使用。
在進(jìn)程之間交換對象?
multiprocessing 支持進(jìn)程之間的兩種通信通道:
隊列
Queue類是一個近似queue.Queue的克隆。 例如:from multiprocessing import Process, Queue def f(q): q.put([42, None, 'hello']) if __name__ == '__main__': q = Queue() p = Process(target=f, args=(q,)) p.start() print(q.get()) # prints "[42, None, 'hello']" p.join()隊列是線程和進(jìn)程安全的。
管道
Pipe()函數(shù)返回一個由管道連接的連接對象,默認(rèn)情況下是雙工(雙向)。例如:from multiprocessing import Process, Pipe def f(conn): conn.send([42, None, 'hello']) conn.close() if __name__ == '__main__': parent_conn, child_conn = Pipe() p = Process(target=f, args=(child_conn,)) p.start() print(parent_conn.recv()) # prints "[42, None, 'hello']" p.join()返回的兩個連接對象
Pipe()表示管道的兩端。每個連接對象都有send()和recv()方法(相互之間的)。請注意,如果兩個進(jìn)程(或線程)同時嘗試讀取或?qū)懭牍艿赖?同一 端,則管道中的數(shù)據(jù)可能會損壞。當(dāng)然,同時使用管道的不同端的進(jìn)程不存在損壞的風(fēng)險。
進(jìn)程之間的同步?
multiprocessing 包含來自 threading 的所有同步原語的等價物。例如,可以使用鎖來確保一次只有一個進(jìn)程打印到標(biāo)準(zhǔn)輸出:
from multiprocessing import Process, Lock
def f(l, i):
l.acquire()
try:
print('hello world', i)
finally:
l.release()
if __name__ == '__main__':
lock = Lock()
for num in range(10):
Process(target=f, args=(lock, num)).start()
不使用來自不同進(jìn)程的鎖輸出容易產(chǎn)生混淆。
在進(jìn)程之間共享狀態(tài)?
如上所述,在進(jìn)行并發(fā)編程時,通常最好盡量避免使用共享狀態(tài)。使用多個進(jìn)程時尤其如此。
但是,如果你真的需要使用一些共享數(shù)據(jù),那么 multiprocessing 提供了兩種方法。
共享內(nèi)存
可以使用
Value或Array將數(shù)據(jù)存儲在共享內(nèi)存映射中。例如,以下代碼:from multiprocessing import Process, Value, Array def f(n, a): n.value = 3.1415927 for i in range(len(a)): a[i] = -a[i] if __name__ == '__main__': num = Value('d', 0.0) arr = Array('i', range(10)) p = Process(target=f, args=(num, arr)) p.start() p.join() print(num.value) print(arr[:])將打印
3.1415927 [0, -1, -2, -3, -4, -5, -6, -7, -8, -9]創(chuàng)建
num和arr時使用的'd'和'i'參數(shù)是array模塊使用的類型的 typecode :'d'表示雙精度浮點數(shù),'i'表示有符號整數(shù)。這些共享對象將是進(jìn)程和線程安全的。為了更靈活地使用共享內(nèi)存,可以使用
multiprocessing.sharedctypes模塊,該模塊支持創(chuàng)建從共享內(nèi)存分配的任意ctypes對象。
服務(wù)器進(jìn)程
由
Manager()返回的管理器對象控制一個服務(wù)器進(jìn)程,該進(jìn)程保存Python對象并允許其他進(jìn)程使用代理操作它們。
Manager()返回的管理器支持類型:list、dict、Namespace、Lock、RLock、Semaphore、BoundedSemaphore、Condition、Event、Barrier、Queue、Value和Array。例如from multiprocessing import Process, Manager def f(d, l): d[1] = '1' d['2'] = 2 d[0.25] = None l.reverse() if __name__ == '__main__': with Manager() as manager: d = manager.dict() l = manager.list(range(10)) p = Process(target=f, args=(d, l)) p.start() p.join() print(d) print(l)將打印
{0.25: None, 1: '1', '2': 2} [9, 8, 7, 6, 5, 4, 3, 2, 1, 0]服務(wù)器進(jìn)程管理器比使用共享內(nèi)存對象更靈活,因為它們可以支持任意對象類型。此外,單個管理器可以通過網(wǎng)絡(luò)由不同計算機(jī)上的進(jìn)程共享。但是,它們比使用共享內(nèi)存慢。
使用工作進(jìn)程?
Pool 類表示一個工作進(jìn)程池。它具有允許以幾種不同方式將任務(wù)分配到工作進(jìn)程的方法。
例如
from multiprocessing import Pool, TimeoutError
import time
import os
def f(x):
return x*x
if __name__ == '__main__':
# start 4 worker processes
with Pool(processes=4) as pool:
# print "[0, 1, 4,..., 81]"
print(pool.map(f, range(10)))
# print same numbers in arbitrary order
for i in pool.imap_unordered(f, range(10)):
print(i)
# evaluate "f(20)" asynchronously
res = pool.apply_async(f, (20,)) # runs in *only* one process
print(res.get(timeout=1)) # prints "400"
# evaluate "os.getpid()" asynchronously
res = pool.apply_async(os.getpid, ()) # runs in *only* one process
print(res.get(timeout=1)) # prints the PID of that process
# launching multiple evaluations asynchronously *may* use more processes
multiple_results = [pool.apply_async(os.getpid, ()) for i in range(4)]
print([res.get(timeout=1) for res in multiple_results])
# make a single worker sleep for 10 secs
res = pool.apply_async(time.sleep, (10,))
try:
print(res.get(timeout=1))
except TimeoutError:
print("We lacked patience and got a multiprocessing.TimeoutError")
print("For the moment, the pool remains available for more work")
# exiting the 'with'-block has stopped the pool
print("Now the pool is closed and no longer available")
請注意,池的方法只能由創(chuàng)建它的進(jìn)程使用。
注解
該軟件包中的功能要求子項可以導(dǎo)入 __main__ 模塊。這包含在 編程指導(dǎo) 中,但值得指出。這意味著一些示例,例如 multiprocessing.pool.Pool 示例在交互式解釋器中不起作用。例如:
>>> from multiprocessing import Pool
>>> p = Pool(5)
>>> def f(x):
... return x*x
...
>>> with p:
... p.map(f, [1,2,3])
Process PoolWorker-1:
Process PoolWorker-2:
Process PoolWorker-3:
Traceback (most recent call last):
AttributeError: 'module' object has no attribute 'f'
AttributeError: 'module' object has no attribute 'f'
AttributeError: 'module' object has no attribute 'f'
(如果你嘗試這個,它實際上會以半隨機(jī)的方式輸出三個完整的回溯,然后你可能不得不以某種方式停止主進(jìn)程。)
參考?
multiprocessing 包大部分復(fù)制了 threading 模塊的API。
Process 和異常?
-
class
multiprocessing.Process(group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)? 進(jìn)程對象表示在單獨進(jìn)程中運行的活動。
Process類等價于threading.Thread。應(yīng)始終使用關(guān)鍵字參數(shù)調(diào)用構(gòu)造函數(shù)。 group 應(yīng)該始終是
None;它僅用于兼容threading.Thread。 target 是由run()方法調(diào)用的可調(diào)用對象。它默認(rèn)為None,意味著什么都沒有被調(diào)用。 name 是進(jìn)程名稱(有關(guān)詳細(xì)信息,請參閱name)。 args 是目標(biāo)調(diào)用的參數(shù)元組。 kwargs 是目標(biāo)調(diào)用的關(guān)鍵字參數(shù)字典。如果提供,則鍵參數(shù) daemon 將進(jìn)程daemon標(biāo)志設(shè)置為True或False。如果是None(默認(rèn)值),則該標(biāo)志將從創(chuàng)建的進(jìn)程繼承。默認(rèn)情況下,不會將任何參數(shù)傳遞給 target 。
如果子類重寫構(gòu)造函數(shù),它必須確保它在對進(jìn)程執(zhí)行任何其他操作之前調(diào)用基類構(gòu)造函數(shù)(
Process.__init__())。在 3.3 版更改: 加入 daemon 參數(shù)。
-
run()? 表示進(jìn)程活動的方法。
你可以在子類中重載此方法。標(biāo)準(zhǔn)
run()方法調(diào)用傳遞給對象構(gòu)造函數(shù)的可調(diào)用對象作為目標(biāo)參數(shù)(如果有),分別從 args 和 kwargs 參數(shù)中獲取順序和關(guān)鍵字參數(shù)。
-
join([timeout])? 如果可選參數(shù) timeout 是
None(默認(rèn)值),則該方法將阻塞,直到調(diào)用join()方法的進(jìn)程終止。如果 timeout 是一個正數(shù),它最多會阻塞 timeout 秒。請注意,如果進(jìn)程終止或方法超時,則該方法返回None。檢查進(jìn)程的exitcode以確定它是否終止。一個進(jìn)程可以合并多次。
進(jìn)程無法并入自身,因為這會導(dǎo)致死鎖。嘗試在啟動進(jìn)程之前合并進(jìn)程是錯誤的。
-
name? 進(jìn)程的名稱。該名稱是一個字符串,僅用于識別目的。它沒有語義??梢詾槎鄠€進(jìn)程指定相同的名稱。
初始名稱由構(gòu)造器設(shè)定。 如果沒有為構(gòu)造器提供顯式名稱,則會構(gòu)造一個形式為 'Process-N1:N2:...:Nk' 的名稱,其中每個 Nk 是其父親的第 N 個孩子。
-
daemon? 進(jìn)程的守護(hù)標(biāo)志,一個布爾值。這必須在
start()被調(diào)用之前設(shè)置。初始值繼承自創(chuàng)建進(jìn)程。
當(dāng)進(jìn)程退出時,它會嘗試終止其所有守護(hù)進(jìn)程子進(jìn)程。
請注意,不允許在守護(hù)進(jìn)程中創(chuàng)建子進(jìn)程。這是因為當(dāng)守護(hù)進(jìn)程由于父進(jìn)程退出而中斷時,其子進(jìn)程會變成孤兒進(jìn)程。 另外,這些 不是 Unix 守護(hù)進(jìn)程或服務(wù),它們是正常進(jìn)程,如果非守護(hù)進(jìn)程已經(jīng)退出,它們將被終止(并且不被合并)。
除了
threading.ThreadAPI ,Process對象還支持以下屬性和方法:-
pid? 返回進(jìn)程ID。在生成該進(jìn)程之前,這將是
None。
-
exitcode? 的退子進(jìn)程出代碼。如果進(jìn)程尚未終止,這將是
None。負(fù)值 -N 表示孩子被信號 N 終止。
-
authkey? 進(jìn)程的身份驗證密鑰(字節(jié)字符串)。
當(dāng)
multiprocessing初始化時,主進(jìn)程使用os.urandom()分配一個隨機(jī)字符串。當(dāng)創(chuàng)建
Process對象時,它將繼承其父進(jìn)程的身份驗證密鑰,盡管可以通過將authkey設(shè)置為另一個字節(jié)字符串來更改。參見 認(rèn)證密碼 。
-
sentinel? 系統(tǒng)對象的數(shù)字句柄,當(dāng)進(jìn)程結(jié)束時將變?yōu)?"ready" 。
如果要使用
multiprocessing.connection.wait()一次等待多個事件,可以使用此值。否則調(diào)用join()更簡單。在Windows上,這是一個操作系統(tǒng)句柄,可以與
WaitForSingleObject和WaitForMultipleObjects系列API調(diào)用一起使用。在Unix上,這是一個文件描述符,可以使用來自select模塊的原語。3.3 新版功能.
-
terminate()? 終止進(jìn)程。 在Unix上,這是使用
SIGTERM信號完成的;在Windows上使用TerminateProcess()。 請注意,不會執(zhí)行退出處理程序和finally子句等。請注意,進(jìn)程的后代進(jìn)程將不會被終止 —— 它們將簡單地變成孤立的。
警告
如果在關(guān)聯(lián)進(jìn)程使用管道或隊列時使用此方法,則管道或隊列可能會損壞,并可能無法被其他進(jìn)程使用。類似地,如果進(jìn)程已獲得鎖或信號量等,則終止它可能導(dǎo)致其他進(jìn)程死鎖。
-
kill()? 與
terminate()相同,但在Unix上使用SIGKILL信號。3.7 新版功能.
-
close()? 關(guān)閉
Process對象,釋放與之關(guān)聯(lián)的所有資源。如果底層進(jìn)程仍在運行,則會引發(fā)ValueError。一旦close()成功返回,Process對象的大多數(shù)其他方法和屬性將引發(fā)ValueError。3.7 新版功能.
注意
start()、join()、is_alive()、terminate()和exitcode方法只能由創(chuàng)建進(jìn)程對象的進(jìn)程調(diào)用。Process一些方法的示例用法:>>> import multiprocessing, time, signal >>> p = multiprocessing.Process(target=time.sleep, args=(1000,)) >>> print(p, p.is_alive()) <Process(Process-1, initial)> False >>> p.start() >>> print(p, p.is_alive()) <Process(Process-1, started)> True >>> p.terminate() >>> time.sleep(0.1) >>> print(p, p.is_alive()) <Process(Process-1, stopped[SIGTERM])> False >>> p.exitcode == -signal.SIGTERM True
-
-
exception
multiprocessing.ProcessError? 所有
multiprocessing異常的基類。
-
exception
multiprocessing.BufferTooShort? 當(dāng)提供的緩沖區(qū)對象太小而無法讀取消息時,
Connection.recv_bytes_into()引發(fā)的異常。如果
e是一個BufferTooShort實例,那么e.args[0]將把消息作為字節(jié)字符串給出。
-
exception
multiprocessing.AuthenticationError? 出現(xiàn)身份驗證錯誤時引發(fā)。
-
exception
multiprocessing.TimeoutError? 有超時的方法超時時引發(fā)。
管道和隊列?
使用多進(jìn)程時,一般使用消息機(jī)制實現(xiàn)進(jìn)程間通信,盡可能避免使用同步原語,例如鎖。
消息機(jī)制包含: Pipe() (可以用于在兩個進(jìn)程間傳遞消息),以及隊列(能夠在多個生產(chǎn)者和消費者之間通信)。
Queue, SimpleQueue 以及 JoinableQueue 都是多生產(chǎn)者,多消費者,并且實現(xiàn)了 FIFO 的隊列類型,其表現(xiàn)與標(biāo)準(zhǔn)庫中的 queue.Queue 類相似。 不同之處在于 Queue 缺少標(biāo)準(zhǔn)庫的 queue.Queue 從 Python 2.5 開始引入的 task_done() 和 join() 方法。
如果你使用了 JoinableQueue ,那么你 必須 對每個已經(jīng)移出隊列的任務(wù)調(diào)用 JoinableQueue.task_done()。 不然的話用于統(tǒng)計未完成任務(wù)的信號量最終會溢出并拋出異常。
另外還可以通過使用一個管理器對象創(chuàng)建一個共享隊列,詳見 數(shù)據(jù)管理器 。
注解
multiprocessing 使用了普通的 queue.Empty 和 queue.Full 異常去表示超時。 你需要從 queue 中導(dǎo)入它們,因為它們并不在 multiprocessing 的命名空間中。
注解
當(dāng)一個對象被放入一個隊列中時,這個對象首先會被一個后臺線程用 pickle 序列化,并將序列化后的數(shù)據(jù)通過一個底層管道的管道傳遞到隊列中。 這種做法會有點讓人驚訝,但一般不會出現(xiàn)什么問題。 如果它們確實妨礙了你,你可以使用一個由管理器 manager 創(chuàng)建的隊列替換它。
將一個對象放入一個空隊列后,可能需要極小的延遲,隊列的方法
empty()? 才會返回False。而get_nowait()可以不拋出queue.Empty直接返回。如果有多個進(jìn)程同時將對象放入隊列,那么在隊列的另一端接受到的對象可能是無序的。但是由同一個進(jìn)程放入的多個對象的順序在另一端輸出時總是一樣的。
警告
如果一個進(jìn)程通過調(diào)用 Process.terminate() 或 os.kill() 在嘗試使用 Queue 期間被終止了,那么隊列中的數(shù)據(jù)很可能被破壞。 這可能導(dǎo)致其他進(jìn)程在嘗試使用該隊列時遇到異常。
警告
正如剛才提到的,如果一個子進(jìn)程將一些對象放進(jìn)隊列中 (并且它沒有用 JoinableQueue.cancel_join_thread 方法),那么這個進(jìn)程在所有緩沖區(qū)的對象被刷新進(jìn)管道之前,是不會終止的。
這意味著,除非你確定所有放入隊列中的對象都已經(jīng)被消費了,否則如果你試圖等待這個進(jìn)程,你可能會陷入死鎖中。相似地,如果該子進(jìn)程不是后臺進(jìn)程,那么父進(jìn)程可能在試圖等待所有非后臺進(jìn)程退出時掛起。
注意用管理器創(chuàng)建的隊列不存在這個問題,詳見 編程指導(dǎo) 。
該 例子 展示了如何使用隊列實現(xiàn)進(jìn)程間通信。
-
multiprocessing.Pipe([duplex])? 返回一對
Connection`對象 ``(conn1, conn2)`, 分別表示管道的兩端。如果 duplex 被置為
True(默認(rèn)值),那么該管道是雙向的。如果 duplex 被置為False,那么該管道是單向的,即conn1只能用于接收消息,而conn2僅能用于發(fā)送消息。
-
class
multiprocessing.Queue([maxsize])? 返回一個使用一個管道和少量鎖和信號量實現(xiàn)的共享隊列實例。當(dāng)一個進(jìn)程將一個對象放進(jìn)隊列中時,一個寫入線程會啟動并將對象從緩沖區(qū)寫入管道中。
一旦超時,將拋出標(biāo)準(zhǔn)庫
queue模塊中常見的異常queue.Empty和queue.Full。除了
task_done()和join()之外,Queue? 實現(xiàn)了標(biāo)準(zhǔn)庫類queue.Queue中所有的方法。-
qsize()? 返回隊列的大致長度。由于多線程或者多進(jìn)程的上下文,這個數(shù)字是不可靠的。
注意,在 Unix 平臺上,例如 Mac OS X ,這個方法可能會拋出
NotImplementedError? 異常,因為該平臺沒有實現(xiàn)sem_getvalue()。
-
empty()? 如果隊列是空的,返回
True,反之返回False。 由于多線程或多進(jìn)程的環(huán)境,該狀態(tài)是不可靠的。
-
full()? 如果隊列是滿的,返回
True,反之返回False。 由于多線程或多進(jìn)程的環(huán)境,該狀態(tài)是不可靠的。
-
put(obj[, block[, timeout]])? 將 obj 放入隊列。如果可選參數(shù) block 是
True(默認(rèn)值) 而且 timeout 是None(默認(rèn)值), 將會阻塞當(dāng)前進(jìn)程,直到有空的緩沖槽。如果 timeout 是正數(shù),將會在阻塞了最多 timeout 秒之后還是沒有可用的緩沖槽時拋出queue.Full? 異常。反之 (block 是False時),僅當(dāng)有可用緩沖槽時才放入對象,否則拋出queue.Full異常 (在這種情形下 timeout 參數(shù)會被忽略)。
-
put_nowait(obj)? 相當(dāng)于
put(obj, False)。
-
get([block[, timeout]])? 從隊列中取出并返回對象。如果可選參數(shù) block 是
True(默認(rèn)值) 而且 timeout 是None(默認(rèn)值), 將會阻塞當(dāng)前進(jìn)程,直到隊列中出現(xiàn)可用的對象。如果 timeout 是正數(shù),將會在阻塞了最多 timeout 秒之后還是沒有可用的對象時拋出queue.Empty異常。反之 (block 是False時),僅當(dāng)有可用對象能夠取出時返回,否則拋出queue.Empty異常 (在這種情形下 timeout 參數(shù)會被忽略)。
-
get_nowait()? 相當(dāng)于
get(False)。
multiprocessing.Queue類有一些在queue.Queue類中沒有出現(xiàn)的方法。這些方法在大多數(shù)情形下并不是必須的。-
close()? 指示當(dāng)前進(jìn)程將不會再往隊列中放入對象。一旦所有緩沖區(qū)中的數(shù)據(jù)被寫入管道之后,后臺的線程會退出。這個方法在隊列被gc回收時會自動調(diào)用。
-
join_thread()? 等待后臺線程。這個方法僅在調(diào)用了
close()方法之后可用。這會阻塞當(dāng)前進(jìn)程,直到后臺線程退出,確保所有緩沖區(qū)中的數(shù)據(jù)都被寫入管道中。默認(rèn)情況下,如果一個不是隊列創(chuàng)建者的進(jìn)程試圖退出,它會嘗試等待這個隊列的后臺線程。這個進(jìn)程可以使用
cancel_join_thread()讓join_thread()方法什么都不做直接跳過。
-
cancel_join_thread()? 防止
join_thread()方法阻塞當(dāng)前進(jìn)程。具體而言,這防止進(jìn)程退出時自動等待后臺線程退出。詳見join_thread()。可能這個方法稱為”
allow_exit_without_flush()“ 會更好。這有可能會導(dǎo)致正在排隊進(jìn)入隊列的數(shù)據(jù)丟失,大多數(shù)情況下你不需要用到這個方法,僅當(dāng)你不關(guān)心底層管道中可能丟失的數(shù)據(jù),只是希望進(jìn)程能夠馬上退出時使用。
注解
該類的功能依賴于宿主操作系統(tǒng)具有可用的共享信號量實現(xiàn)。否則該類將被禁用,任何試圖實例化一個
Queue對象的操作都會拋出ImportError異常,更多信息詳見 bpo-3770 。后續(xù)說明的任何專用隊列對象亦如此。-
-
class
multiprocessing.SimpleQueue? 這是一個簡化的
Queue類的實現(xiàn),很像帶鎖的Pipe。-
empty()? 如果隊列為空返回
True,否則返回False。
-
get()? 從隊列中移出并返回一個對象。
-
put(item)? 將 item 放入隊列。
-
-
class
multiprocessing.JoinableQueue([maxsize])? JoinableQueue類是Queue的子類,額外添加了task_done()和join()方法。-
task_done()? 指出之前進(jìn)入隊列的任務(wù)已經(jīng)完成。由隊列的消費者進(jìn)程使用。對于每次調(diào)用
get()獲取的任務(wù),執(zhí)行完成后調(diào)用task_done()告訴隊列該任務(wù)已經(jīng)處理完成。如果
join()方法正在阻塞之中,該方法會在所有對象都被處理完的時候返回 (即對之前使用put()放進(jìn)隊列中的所有對象都已經(jīng)返回了對應(yīng)的task_done()) 。如果被調(diào)用的次數(shù)多于放入隊列中的項目數(shù)量,將引發(fā)
ValueError異常 。
-
join()? 阻塞至隊列中所有的元素都被接收和處理完畢。
當(dāng)條目添加到隊列的時候,未完成任務(wù)的計數(shù)就會增加。每當(dāng)消費者進(jìn)程調(diào)用
task_done()表示這個條目已經(jīng)被回收,該條目所有工作已經(jīng)完成,未完成計數(shù)就會減少。當(dāng)未完成計數(shù)降到零的時候,join()阻塞被解除。
-
雜項?
-
multiprocessing.active_children()? 返回當(dāng)前進(jìn)程存活的子進(jìn)程的列表。
調(diào)用該方法有“等待”已經(jīng)結(jié)束的進(jìn)程的副作用。
-
multiprocessing.cpu_count()? 返回系統(tǒng)的CPU數(shù)量。
該數(shù)量不同于當(dāng)前進(jìn)程可以使用的CPU數(shù)量??捎玫腃PU數(shù)量可以由
len(os.sched_getaffinity(0))方法獲得。可能引發(fā)
NotImplementedError。
-
multiprocessing.freeze_support()? 為使用了
multiprocessing? 的程序,提供凍結(jié)以產(chǎn)生 Windows 可執(zhí)行文件的支持。(在 py2exe, PyInstaller 和 cx_Freeze 上測試通過)需要在 main 模塊的
if __name__ == '__main__'該行之后馬上調(diào)用該函數(shù)。例如:from multiprocessing import Process, freeze_support def f(): print('hello world!') if __name__ == '__main__': freeze_support() Process(target=f).start()
如果沒有調(diào)用
freeze_support()在嘗試運行被凍結(jié)的可執(zhí)行文件時會拋出RuntimeError異常。對
freeze_support()的調(diào)用在非 Windows 平臺上是無效的。如果該模塊在 Windows 平臺的 Python 解釋器中正常運行 (該程序沒有被凍結(jié)), 調(diào)用``freeze_support()`` 也是無效的。
-
multiprocessing.get_all_start_methods()? 返回支持的啟動方法的列表,該列表的首項即為默認(rèn)選項??赡艿膯臃椒ㄓ?
'fork','spawn'和``'forkserver'。在 Windows 中,只有 ``'spawn'是可用的。Unix平臺總是支持``'fork'`` 和``'spawn',且 ``'fork'是默認(rèn)值。3.4 新版功能.
-
multiprocessing.get_context(method=None)? 返回一個 Context 對象。該對象具有和
multiprocessing模塊相同的API。如果 method 設(shè)置成
None那么將返回默認(rèn)上下文對象。否則 method? 應(yīng)該是'fork','spawn','forkserver'。 如果指定的啟動方法不存在,將拋出ValueError? 異常。3.4 新版功能.
-
multiprocessing.get_start_method(allow_none=False)? 返回啟動進(jìn)程時使用的啟動方法名。
如果啟動方法已經(jīng)固定,并且 allow_none 被設(shè)置成 False ,那么啟動方法將被固定為默認(rèn)的啟動方法,并且返回其方法名。如果啟動方法沒有設(shè)定,并且 allow_none 被設(shè)置成 True ,那么將返回
None? 。返回值將為
'fork','spawn','forkserver'或者None。'fork'``是 Unix 的默認(rèn)值, ``'spawn'是 Windows 的默認(rèn)值。3.4 新版功能.
-
multiprocessing.set_executable()? 設(shè)置在啟動子進(jìn)程時使用的 Python 解釋器路徑。 ( 默認(rèn)使用
sys.executable) 嵌入式編程人員可能需要這樣做:set_executable(os.path.join(sys.exec_prefix, 'pythonw.exe'))
以使他們可以創(chuàng)建子進(jìn)程。
在 3.4 版更改: 現(xiàn)在在 Unix 平臺上使用
'spawn'? 啟動方法時支持調(diào)用該方法。
-
multiprocessing.set_start_method(method)? 設(shè)置啟動子進(jìn)程的方法。 method 可以是
'fork','spawn'或者'forkserver'。注意這最多只能調(diào)用一次,并且需要藏在 main 模塊中,由
if __name__ == '__main__'? 保護(hù)著。3.4 新版功能.
注解
multiprocessing 并沒有包含類似 threading.active_count() , threading.enumerate() , threading.settrace() , threading.setprofile(), threading.Timer , 或者 threading.local 的方法和類。
連接對象(Connection)?
Connection 對象允許收發(fā)可以序列化的對象或字符串。它們可以看作面向消息的連接套接字。
通常使用 Pipe 創(chuàng)建 Connection 對象。詳見 : 監(jiān)聽者及客戶端.
-
class
multiprocessing.connection.Connection? -
send(obj)? 將一個對象發(fā)送到連接的另一端,可以用
recv()讀取。發(fā)送的對象必須是可以序列化的,過大的對象 ( 接近 32MiB+ ,這個值取決于操作系統(tǒng) ) 有可能引發(fā)
ValueError? 異常。
-
fileno()? 返回由連接對象使用的描述符或者句柄。
-
close()? 關(guān)閉連接對象。
當(dāng)連接對象被垃圾回收時會自動調(diào)用。
-
poll([timeout])? 返回連接對象中是否有可以讀取的數(shù)據(jù)。
如果未指定 timeout ,此方法會馬上返回。如果 timeout 是一個數(shù)字,則指定了最大阻塞的秒數(shù)。如果 timeout 是
None? ,那么將一直等待,不會超時。注意通過使用
multiprocessing.connection.wait()可以一次輪詢多個連接對象。
-
send_bytes(buffer[, offset[, size]])? 從一個 bytes-like object? (字節(jié)類對象)對象中取出字節(jié)數(shù)組并作為一條完整消息發(fā)送。
如果由 offset? 給定了在 buffer 中讀取數(shù)據(jù)的位置。 如果給定了 size ,那么將會從緩沖區(qū)中讀取多個字節(jié)。 過大的緩沖區(qū) ( 接近 32MiB+ ,此值依賴于操作系統(tǒng) ) 有可能引發(fā)
ValueError? 異常。
-
recv_bytes([maxlength])? 以字符串形式返回一條從連接對象另一端發(fā)送過來的字節(jié)數(shù)據(jù)。此方法在接收到數(shù)據(jù)前將一直阻塞。 如果連接對象被對端關(guān)閉或者沒有數(shù)據(jù)可讀取,將拋出
EOFError? 異常。如果給定了 maxlength 并且消息長于 maxlength 那么將拋出
OSError并且該連接對象將不再可讀。
-
recv_bytes_into(buffer[, offset])? 將一條完整的字節(jié)數(shù)據(jù)消息讀入 buffer 中并返回消息的字節(jié)數(shù)。 此方法在接收到數(shù)據(jù)前將一直阻塞。 如果連接對象被對端關(guān)閉或者沒有數(shù)據(jù)可讀取,將拋出
EOFError? 異常。buffer must be a writable bytes-like object. If offset is given then the message will be written into the buffer from that position. Offset must be a non-negative integer less than the length of buffer (in bytes).
如果緩沖區(qū)太小,則將引發(fā)
BufferTooShort? 異常,并且完整的消息將會存放在異常實例e的e.args[0]中。
在 3.3 版更改: 現(xiàn)在連接對象自身可以通過
Connection.send()和Connection.recv()在進(jìn)程之間傳遞。3.3 新版功能: 連接對象現(xiàn)已支持上下文管理協(xié)議 -- 參見 see 上下文管理器類型 。
__enter__()返回連接對象,__exit__()會調(diào)用close()。-
例如:
>>> from multiprocessing import Pipe
>>> a, b = Pipe()
>>> a.send([1, 'hello', None])
>>> b.recv()
[1, 'hello', None]
>>> b.send_bytes(b'thank you')
>>> a.recv_bytes()
b'thank you'
>>> import array
>>> arr1 = array.array('i', range(5))
>>> arr2 = array.array('i', [0] * 10)
>>> a.send_bytes(arr1)
>>> count = b.recv_bytes_into(arr2)
>>> assert count == len(arr1) * arr1.itemsize
>>> arr2
array('i', [0, 1, 2, 3, 4, 0, 0, 0, 0, 0])
警告
Connection.recv() 方法會自動解封它收到的數(shù)據(jù),除非你能夠信任發(fā)送消息的進(jìn)程,否則此處可能有安全風(fēng)險。
因此, 除非連接對象是由 Pipe() 產(chǎn)生的,否則你應(yīng)該僅在使用了某種認(rèn)證手段之后才使用 recv() 和 send() 方法。 參考 認(rèn)證密碼。
警告
如果一個進(jìn)程在試圖讀寫管道時被終止了,那么管道中的數(shù)據(jù)很可能是不完整的,因為此時可能無法確定消息的邊界。
同步原語?
通常來說同步原語在多進(jìn)程環(huán)境中并不像它們在多線程環(huán)境中那么必要。參考 threading 模塊的文檔。
注意可以使用管理器對象創(chuàng)建同步原語,參考 數(shù)據(jù)管理器 。
-
class
multiprocessing.Barrier(parties[, action[, timeout]])? 類似
threading.Barrier的柵欄對象。3.3 新版功能.
-
class
multiprocessing.BoundedSemaphore([value])? 非常類似
threading.BoundedSemaphore的有界信號量對象。一個小小的不同在于,它的
acquire? 方法的第一個參數(shù)名是和Lock.acquire()一樣的 block 。注解
在 Mac OS X 平臺上, 該對象于
Semaphore? 不同在于sem_getvalue()方法并沒有在該平臺上實現(xiàn)。
-
class
multiprocessing.Condition([lock])? 條件變量:
threading.Condition的別名。指定的 lock 參數(shù)應(yīng)該是
multiprocessing模塊中的Lock或者RLock對象。在 3.3 版更改: 新增了
wait_for()? 方法。
-
class
multiprocessing.Event? A clone of
threading.Event.
-
class
multiprocessing.Lock? 原始鎖(非遞歸鎖)對象,類似于
threading.Lock。一旦一個進(jìn)程或者線程拿到了鎖,后續(xù)的任何其他進(jìn)程或線程的其他請求都會被阻塞直到鎖被釋放。任何進(jìn)程或線程都可以釋放鎖。除非另有說明,否則multiprocessing.Lock? 用于進(jìn)程或者線程的概念和行為都和threading.Lock? 一致。注意
Lock實際上是一個工廠函數(shù)。它返回由默認(rèn)上下文初始化的multiprocessing.synchronize.Lock? 對象。Locksupports the context manager protocol and thus may be used inwithstatements.-
acquire(block=True, timeout=None)? 可以阻塞或非阻塞地獲得鎖。
如果 block 參數(shù)被設(shè)為
True( 默認(rèn)值 ) , 對該方法的調(diào)用在鎖處于釋放狀態(tài)之前都會阻塞,然后將鎖設(shè)置為鎖住狀態(tài)并返回True。需要注意的是第一個參數(shù)名與threading.Lock.acquire()的不同。如果 block 參數(shù)被設(shè)置成
False,方法的調(diào)用將不會阻塞。 如果鎖當(dāng)前處于鎖住狀態(tài),將返回False; 否則將鎖設(shè)置成鎖住狀態(tài),并返回True。當(dāng) timeout 是一個正浮點數(shù)時,會在等待鎖的過程中最多阻塞等待 timeout 秒,當(dāng) timeout 是負(fù)數(shù)時,效果和 timeout 為0時一樣,當(dāng) timeout 是
None(默認(rèn)值)時,等待時間是無限長。需要注意的是,對于 timeout 是負(fù)數(shù)和None的情況, 其行為與threading.Lock.acquire()是不一樣的。當(dāng) block 參數(shù)為False時, timeout? 并沒有實際用處,會直接忽略, 當(dāng) block 參數(shù)為True時,函數(shù)會在拿到鎖后返回True或者 超時沒拿到鎖后返回False。
-
release()? 釋放鎖,可以在任何進(jìn)程、線程使用,并不限于鎖的擁有者。
其行為與
threading.Lock.release()一樣,只不過當(dāng)嘗試釋放一個沒有被持有的鎖時,會拋出ValueError異常。
-
-
class
multiprocessing.RLock? 遞歸鎖對象: 類似于
threading.RLock。遞歸鎖必須由持有線程、進(jìn)程親自釋放。如果某個進(jìn)程或者線程拿到了遞歸鎖,這個進(jìn)程或者線程可以再次拿到這個鎖而不需要等待。但是這個進(jìn)程或者線程的拿鎖操作和釋放鎖操作的次數(shù)必須相同。注意
RLock是一個工廠函數(shù),調(diào)用后返回一個使用默認(rèn) context 初始化的multiprocessing.synchronize.RLock對象。RLock支持 context manager 協(xié)議,因此可在with語句內(nèi)使用。-
acquire(block=True, timeout=None)? 可以阻塞或非阻塞地獲得鎖。
當(dāng) block 設(shè)置為
True時,會一直阻塞直到鎖處于空閑狀態(tài)(沒有被任何進(jìn)程、線程擁有),除非當(dāng)前進(jìn)程/線程已經(jīng)擁有了這把鎖。然后當(dāng)前進(jìn)程會持有這把鎖(在鎖沒有被持有者的情況下),鎖內(nèi)的遞歸等級加一,并返回True. 注意, 這個函數(shù)第一個參數(shù)和threading.RLock.acquire()有幾個不同點,包括參數(shù)名本身。當(dāng) block 參數(shù)是
False, 將不會阻塞,如果此時鎖被其他進(jìn)程或者線程持有,當(dāng)前進(jìn)程、線程獲取鎖操作失敗,鎖的遞歸等級也不會改變,函數(shù)返回False, 如果當(dāng)前鎖已經(jīng)處于釋放狀態(tài),則當(dāng)前進(jìn)程、線程則會拿到鎖,并且鎖內(nèi)的遞歸等級加一,函數(shù)返回True。timeout 參數(shù)的使用方法及行為與
Lock.acquire()一樣。但是要注意 timeout 的其中一些行為和threading.RLock.acquire()中實現(xiàn)的行為是不同的。
-
release()? 釋放鎖,使鎖內(nèi)的遞歸等級減一。如果釋放后鎖內(nèi)的遞歸等級降低為0,則會重置鎖的狀態(tài)為釋放狀態(tài)(即沒有被任何進(jìn)程、線程持有),重置后如果有有其他進(jìn)程和線程在等待這把鎖,他們中的一個會獲得這個鎖而繼續(xù)運行。如果釋放后鎖內(nèi)的遞歸等級還沒到達(dá)0,則這個鎖仍將保持未釋放狀態(tài)且當(dāng)前進(jìn)程和線程仍然是持有者。
只有當(dāng)前進(jìn)程或線程是鎖的持有者時,才允許調(diào)用這個方法。如果當(dāng)前進(jìn)程或線程不是這個鎖的擁有者,或者這個鎖處于已釋放的狀態(tài)(即沒有任何擁有者),調(diào)用這個方法會拋出
AssertionError異常。注意這里拋出的異常類型和threading.RLock.release()中實現(xiàn)的行為不一樣。
-
-
class
multiprocessing.Semaphore([value])? 一種信號量對象: 類似于
threading.Semaphore.一個小小的不同在于,它的
acquire? 方法的第一個參數(shù)名是和Lock.acquire()一樣的 block 。
注解
在 Mac OS X 上,不支持 sem_timedwait ,所以,使用調(diào)用 acquire() 時如果使用 timeout 參數(shù),會通過循環(huán)sleep來模擬timeout的行為。
注解
假如信號 SIGINT 是來自于 Ctrl-C ,并且主線程被 BoundedSemaphore.acquire(), Lock.acquire(), RLock.acquire(), Semaphore.acquire(), Condition.acquire() 或 Condition.wait() 阻塞,則調(diào)用會立即中斷同時拋出 KeyboardInterrupt 異常。
這和 threading 的行為不同,此模塊中當(dāng)執(zhí)行對應(yīng)的阻塞式調(diào)用時,SIGINT 會被忽略。
注解
這個庫的某些功能依賴于宿主機(jī)系統(tǒng)的共享信號量,如果系統(tǒng)沒有這個特性, multiprocessing.synchronize 會被禁用,嘗試導(dǎo)入這個包會引發(fā) ImportError 異常,詳細(xì)信息請查看 bpo-3770 。
共享 ctypes 對象?
在共享內(nèi)存上創(chuàng)建可被子進(jìn)程繼承的共享對象時是可行的。
-
multiprocessing.Value(typecode_or_type, *args, lock=True)? 返回一個從共享內(nèi)存上創(chuàng)建的
ctypes對象。默認(rèn)情況下返回的實際上是經(jīng)過了同步包裝器包裝過的。可以通過Value的 value 屬性訪問這個對象本身。typecode_or_type 指明了返回的對象類型: 它可能是一個 ctypes 類型或者
array? 模塊中每個類型對應(yīng)的單字符長度的字符串。 *args 會透傳給這個類的構(gòu)造函數(shù)。如果 lock 參數(shù)是
True(默認(rèn)值), 將會新建一個遞歸鎖用于同步對于此值的訪問操作。 如果 lock 是Lock或者RLock對象,那么這個傳入的鎖將會用于同步對這個值的訪問操作,如果 lock 是False, 那么對這個對象的訪問將沒有鎖保護(hù),也就是說這個變量不是進(jìn)程安全的。諸如
+=這類的操作會引發(fā)獨立的讀操作和寫操作,也就是說這類操作符并不具有原子性。所以,如果你想讓遞增操作具有原子性,這樣的方式并不能達(dá)到要求:counter.value += 1
假設(shè)共享對象內(nèi)部關(guān)聯(lián)的鎖時遞歸鎖(默認(rèn)情況下就是), 那么你可以采用這種方式
with counter.get_lock(): counter.value += 1
注意 locl 只能是命名參數(shù)。
-
multiprocessing.Array(typecode_or_type, size_or_initializer, *, lock=True)? 從共享內(nèi)存中申請并返回一個具有ctypes類型的數(shù)組對象。默認(rèn)情況下返回值實際上是被同步器包裝過的數(shù)組對象。
typecode_or_type 指明了返回的數(shù)組中的元素類型: 它可能是一個 ctypes 類型或者
array模塊中每個類型對應(yīng)的單字符長度的字符串。 如果 size_or_initializer 是一個整數(shù),那就會當(dāng)做數(shù)組的長度,并且整個數(shù)組的內(nèi)存會初始化為0。否則,如果 size_or_initializer 會被當(dāng)成一個序列用于初始化數(shù)組中的內(nèi)一個元素,并且會根據(jù)元素個數(shù)自動判斷數(shù)組的長度。如果 lock 為
True(默認(rèn)值) 則將創(chuàng)建一個新的鎖對象用于同步對值的訪問。 如果 lock 為一個Lock或RLock對象則該對象將被用于同步對值的訪問。 如果 lock 為False則對返回對象的訪問將不會自動得到鎖的保護(hù),也就是說它不是“進(jìn)程安全的”。請注意 lock 是一個僅限關(guān)鍵字參數(shù)。
請注意
ctypes.c_char的數(shù)組具有 value 和 raw 屬性,允許被用來保存和提取字符串。
multiprocessing.sharedctypes 模塊?
multiprocessing.sharedctypes 模塊提供了一些函數(shù),用于分配來自共享內(nèi)存的、可被子進(jìn)程繼承的 ctypes 對象。
注解
雖然可以將指針存儲在共享內(nèi)存中,但請記住它所引用的是特定進(jìn)程地址空間中的位置。 而且,指針很可能在第二個進(jìn)程的上下文中無效,嘗試從第二個進(jìn)程對指針進(jìn)行解引用可能會導(dǎo)致崩潰。
從共享內(nèi)存中申請并返回一個 ctypes 數(shù)組。
typecode_or_type 指明了返回的數(shù)組中的元素類型: 它可能是一個 ctypes 類型或者
array模塊中使用的類型字符。 如果 size_or_initializer 是一個整數(shù),那就會當(dāng)做數(shù)組的長度,并且整個數(shù)組的內(nèi)存會初始化為0。否則,如果 size_or_initializer 會被當(dāng)成一個序列用于初始化數(shù)組中的每一個元素,并且會根據(jù)元素個數(shù)自動判斷數(shù)組的長度。注意對元素的訪問、賦值操作可能是非原子操作 - 使用
Array()來借助其中的鎖保證操作的原子性。
從共享內(nèi)存中申請并返回一個 ctypes 對象。
typecode_or_type 指明了返回的對象類型: 它可能是一個 ctypes 類型或者
array? 模塊中每個類型對應(yīng)的單字符長度的字符串。 *args 會透傳給這個類的構(gòu)造函數(shù)。注意對 value 的訪問、賦值操作可能是非原子操作 - 使用
Value()來借助其中的鎖保證操作的原子性。請注意
ctypes.c_char的數(shù)組具有 value 和 raw 屬性,允許被用來保存和提取字符串 - 請查看ctypes文檔。
返回一個純 ctypes 數(shù)組, 或者在此之上經(jīng)過同步器包裝過的對象,這取決于 lock 參數(shù)的值,除此之外,和
RawArray()一樣。如果 lock 為
True(默認(rèn)值) 則將創(chuàng)建一個新的鎖對象用于同步對值的訪問。 如果 lock 為一個Lock或RLock對象則該對象將被用于同步對值的訪問。 如果 lock 為False則對所返回對象的訪問將不會自動得到鎖的保護(hù),也就是說它將不是“進(jìn)程安全的”。注意 locl 只能是命名參數(shù)。
返回一個純 ctypes 數(shù)組, 或者在此之上經(jīng)過同步器包裝過的進(jìn)程安全的對象,這取決于 lock 參數(shù)的值,除此之外,和
RawArray()一樣。如果 lock 為
True(默認(rèn)值) 則將創(chuàng)建一個新的鎖對象用于同步對值的訪問。 如果 lock 為一個Lock或RLock對象則該對象將被用于同步對值的訪問。 如果 lock 為False則對所返回對象的訪問將不會自動得到鎖的保護(hù),也就是說它將不是“進(jìn)程安全的”。注意 locl 只能是命名參數(shù)。
從共享內(nèi)存中申請一片空間將 ctypes 對象 obj 過來,然后返回一個新的 ctypes 對象。
將一個 ctypes 對象包裝為進(jìn)程安全的對象并返回,使用 lock 同步對于它的操作。如果 lock 是
None(默認(rèn)值) ,則會自動創(chuàng)建一個multiprocessing.RLock對象。同步器包裝后的對象會在原有對象基礎(chǔ)上額外增加兩個方法:
get_obj()返回被包裝的對象,get_lock()返回內(nèi)部用于同步的鎖。需要注意的是,訪問包裝后的ctypes對象會比直接訪問原來的純 ctypes 對象慢得多。
在 3.5 版更改: 同步器包裝后的對象支持 context manager 協(xié)議。
下面的表格對比了創(chuàng)建普通ctypes對象和基于共享內(nèi)存上創(chuàng)建共享ctypes對象的語法。(表格中的 MyStruct 是 ctypes.Structure 的子類)
ctypes |
使用類型的共享ctypes |
使用 typecode 的共享 ctypes |
|---|---|---|
c_double(2.4) |
RawValue(c_double, 2.4) |
RawValue('d', 2.4) |
MyStruct(4, 6) |
RawValue(MyStruct, 4, 6) |
|
(c_short * 7)() |
RawArray(c_short, 7) |
RawArray('h', 7) |
(c_int * 3)(9, 2, 8) |
RawArray(c_int, (9, 2, 8)) |
RawArray('i', (9, 2, 8)) |
下面是一個在子進(jìn)程中修改多個ctypes對象的例子。
from multiprocessing import Process, Lock
from multiprocessing.sharedctypes import Value, Array
from ctypes import Structure, c_double
class Point(Structure):
_fields_ = [('x', c_double), ('y', c_double)]
def modify(n, x, s, A):
n.value **= 2
x.value **= 2
s.value = s.value.upper()
for a in A:
a.x **= 2
a.y **= 2
if __name__ == '__main__':
lock = Lock()
n = Value('i', 7)
x = Value(c_double, 1.0/3.0, lock=False)
s = Array('c', b'hello world', lock=lock)
A = Array(Point, [(1.875,-6.25), (-5.75,2.0), (2.375,9.5)], lock=lock)
p = Process(target=modify, args=(n, x, s, A))
p.start()
p.join()
print(n.value)
print(x.value)
print(s.value)
print([(a.x, a.y) for a in A])
輸出如下
49
0.1111111111111111
HELLO WORLD
[(3.515625, 39.0625), (33.0625, 4.0), (5.640625, 90.25)]
數(shù)據(jù)管理器?
管理器提供了一種創(chuàng)建共享數(shù)據(jù)的方法,從而可以在不同進(jìn)程中共享,甚至可以通過網(wǎng)絡(luò)跨機(jī)器共享數(shù)據(jù)。管理器維護(hù)一個用于管理 共享對象 的服務(wù)。其他進(jìn)程可以通過代理訪問這些共享對象。
返回一個已啟動的
SyncManager管理器對象,這個對象可以用于在不同進(jìn)程中共享數(shù)據(jù)。返回的管理器對象對應(yīng)了一個 spawned 方式啟動的子進(jìn)程,并且擁有一系列方法可以用于創(chuàng)建共享對象、返回對應(yīng)的代理。
當(dāng)管理器被垃圾回收或者父進(jìn)程退出時,管理器進(jìn)程會立即退出。管理器類定義在 multiprocessing.managers 模塊:
-
class
multiprocessing.managers.BaseManager([address[, authkey]])? 創(chuàng)建一個 BaseManager 對象。
一旦創(chuàng)建,應(yīng)該及時調(diào)用
start()或者get_server().serve_forever()以確保管理器對象對應(yīng)的管理進(jìn)程已經(jīng)啟動。address 是管理器對象監(jiān)聽的地址。如果 address 是
None,則允許和任意主機(jī)的請求建立連接。authkey 是認(rèn)證標(biāo)識,用于檢查連接服務(wù)進(jìn)程的請求合法性。如果 authkey 是
None, 則會使用current_process().authkey, 否則,就使用 authkey , 需要保證它必須是 byte 類型的字符串。-
start([initializer[, initargs]])? 為管理器開啟一個子進(jìn)程,如果 initializer 不是
None, 子進(jìn)程在啟動時將會調(diào)用initializer(*initargs)。
-
get_server()? 返回一個
Server? 對象,它是管理器在后臺控制的真實的服務(wù)。Server? 對象擁有serve_forever()方法。>>> from multiprocessing.managers import BaseManager >>> manager = BaseManager(address=('', 50000), authkey=b'abc') >>> server = manager.get_server() >>> server.serve_forever()
Server額外擁有一個address屬性。
-
connect()? 將本地管理器對象連接到一個遠(yuǎn)程管理器進(jìn)程:
>>> from multiprocessing.managers import BaseManager >>> m = BaseManager(address=('127.0.0.1', 50000), authkey=b'abc') >>> m.connect()
-
register(typeid[, callable[, proxytype[, exposed[, method_to_typeid[, create_method]]]]])? 一個 classmethod,可以將一個類型或者可調(diào)用對象注冊到管理器類。
typeid 是一種 "類型標(biāo)識",用于唯一表示某種共享類型,必須是一個字符串。
callable 是一個用來為此類型標(biāo)識符創(chuàng)建對象的可調(diào)用對象。如果一個管理器實例將使用
connect()方法連接到服務(wù)器,或者 create_method 參數(shù)為False,那么這里可留下None。proxytype 是
BaseProxy? 的子類,可以根據(jù) typeid 為共享對象創(chuàng)建一個代理,如果是None, 則會自動創(chuàng)建一個代理類。exposed 是一個函數(shù)名組成的序列,用來指明只有這些方法可以使用
BaseProxy._callmethod()代理。(如果 exposed 是None, 則會在proxytype._exposed_存在的情況下轉(zhuǎn)而使用它) 當(dāng)暴露的方法列表沒有指定的時候,共享對象的所有 “公共方法” 都會被代理。(這里的“公共方法”是指所有擁有__call__()方法并且不是以'_'開頭的屬性)method_to_typeid 是一個映射,用來指定那些應(yīng)該返回代理對象的暴露方法所返回的類型。(如果 method_to_typeid 是
None, 則proxytype._method_to_typeid_會在存在的情況下被使用)如果方法名稱不在這個映射中或者映射是None,則方法返回的對象會是一個值拷貝。create_method 指明,是否要創(chuàng)建一個以 typeid 命名并返回一個代理對象的函數(shù),這個函數(shù)會被服務(wù)進(jìn)程用于創(chuàng)建共享對象,默認(rèn)為
True。
BaseManager實例也有一個只讀屬性。-
address? 管理器所用的地址。
在 3.3 版更改: 管理器對象支持上下文管理協(xié)議 - 查看 上下文管理器類型 。
__enter__()啟動服務(wù)進(jìn)程(如果它還沒有啟動)并且返回管理器對象,__exit__()會調(diào)用shutdown()。在之前的版本中,如果管理器服務(wù)進(jìn)程沒有啟動,
__enter__()不會負(fù)責(zé)啟動它。-
-
class
multiprocessing.managers.SyncManager? BaseManager的子類,可用于進(jìn)程的同步。這個類型的對象使用multiprocessing.Manager()創(chuàng)建。它擁有一系列方法,可以為大部分常用數(shù)據(jù)類型創(chuàng)建并返回 代理對象 代理,用于進(jìn)程間同步。甚至包括共享列表和字典。
-
Barrier(parties[, action[, timeout]])? 創(chuàng)建一個共享的
threading.Barrier對象并返回它的代理。3.3 新版功能.
-
BoundedSemaphore([value])? 創(chuàng)建一個共享的
threading.BoundedSemaphore對象并返回它的代理。
-
Condition([lock])? 創(chuàng)建一個共享的
threading.Condition對象并返回它的代理。如果提供了 lock 參數(shù),那它必須是
threading.Lock或threading.RLock的代理對象。在 3.3 版更改: 新增了
wait_for()? 方法。
-
Event()? 創(chuàng)建一個共享的
threading.Event對象并返回它的代理。
-
Lock()? 創(chuàng)建一個共享的
threading.Lock對象并返回它的代理。
-
Namespace()? 創(chuàng)建一個共享的
Namespace`對象并返回它的代理。
-
Queue([maxsize])? 創(chuàng)建一個共享的
queue.Queue對象并返回它的代理。
-
RLock()? 創(chuàng)建一個共享的
threading.RLock對象并返回它的代理。
-
Semaphore([value])? 創(chuàng)建一個共享的
threading.Semaphore對象并返回它的代理。
-
Array(typecode, sequence)? 創(chuàng)建一個數(shù)組并返回它的代理。
-
Value(typecode, value)? 創(chuàng)建一個具有可寫
value屬性的對象并返回它的代理。
在 3.6 版更改: 共享對象能夠嵌套。例如, 共享的容器對象如共享列表,可以包含另一個共享對象,他們?nèi)紩?
SyncManager中進(jìn)行管理和同步。-
-
class
multiprocessing.managers.Namespace? 一個可以注冊到
SyncManager的類型。命名空間對象沒有公共方法,但是擁有可寫的屬性。它的表示(repr)會顯示所有屬性的值。
值得一提的是,當(dāng)對命名空間對象使用代理的時候,訪問所有名稱以
'_'開頭的屬性都只是代理器上的屬性,而不是命名空間對象的屬性。>>> manager = multiprocessing.Manager() >>> Global = manager.Namespace() >>> Global.x = 10 >>> Global.y = 'hello' >>> Global._z = 12.3 # this is an attribute of the proxy >>> print(Global) Namespace(x=10, y='hello')
自定義管理器?
要創(chuàng)建一個自定義的管理器,需要新建一個 BaseManager 的子類,然后使用這個管理器類上的 register() 類方法將新類型或者可調(diào)用方法注冊上去。例如:
from multiprocessing.managers import BaseManager
class MathsClass:
def add(self, x, y):
return x + y
def mul(self, x, y):
return x * y
class MyManager(BaseManager):
pass
MyManager.register('Maths', MathsClass)
if __name__ == '__main__':
with MyManager() as manager:
maths = manager.Maths()
print(maths.add(4, 3)) # prints 7
print(maths.mul(7, 8)) # prints 56
使用遠(yuǎn)程管理器?
可以將管理器服務(wù)運行在一臺機(jī)器上,然后使用客戶端從其他機(jī)器上訪問。(假設(shè)它們的防火墻允許這樣的網(wǎng)絡(luò)通信)
運行下面的代碼可以啟動一個服務(wù),此付包含了一個共享隊列,允許遠(yuǎn)程客戶端訪問:
>>> from multiprocessing.managers import BaseManager
>>> from queue import Queue
>>> queue = Queue()
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue', callable=lambda:queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()
遠(yuǎn)程客戶端可以通過下面的方式訪問服務(wù):
>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.put('hello')
也可以通過下面的方式:
>>> from multiprocessing.managers import BaseManager
>>> class QueueManager(BaseManager): pass
>>> QueueManager.register('get_queue')
>>> m = QueueManager(address=('foo.bar.org', 50000), authkey=b'abracadabra')
>>> m.connect()
>>> queue = m.get_queue()
>>> queue.get()
'hello'
本地進(jìn)程也可以訪問這個隊列,利用上面的客戶端代碼通過遠(yuǎn)程方式訪問:
>>> from multiprocessing import Process, Queue
>>> from multiprocessing.managers import BaseManager
>>> class Worker(Process):
... def __init__(self, q):
... self.q = q
... super(Worker, self).__init__()
... def run(self):
... self.q.put('local hello')
...
>>> queue = Queue()
>>> w = Worker(queue)
>>> w.start()
>>> class QueueManager(BaseManager): pass
...
>>> QueueManager.register('get_queue', callable=lambda: queue)
>>> m = QueueManager(address=('', 50000), authkey=b'abracadabra')
>>> s = m.get_server()
>>> s.serve_forever()
代理對象?
代理是一個 指向 其他共享對象的對象,這個對象(很可能)在另外一個進(jìn)程中。共享對象也可以說是代理 指涉 的對象。多個代理對象可能指向同一個指涉對象。
代理對象代理了指涉對象的一系列方法調(diào)用(雖然并不是指涉對象的每個方法都有必要被代理)。通過這種方式,代理的使用方法可以和它的指涉對象一樣:
>>> from multiprocessing import Manager
>>> manager = Manager()
>>> l = manager.list([i*i for i in range(10)])
>>> print(l)
[0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
>>> print(repr(l))
<ListProxy object, typeid 'list' at 0x...>
>>> l[4]
16
>>> l[2:5]
[4, 9, 16]
注意,對代理使用 str() 函數(shù)會返回指涉對象的字符串表示,但是 repr() 卻會返回代理本身的內(nèi)部字符串表示。
被代理的對象很重要的一點是必須可以被序列化,這樣才能允許他們在進(jìn)程間傳遞。因此,指涉對象可以包含 代理對象 。這允許管理器中列表、字典或者其他 代理對象 對象之間的嵌套。
>>> a = manager.list()
>>> b = manager.list()
>>> a.append(b) # referent of a now contains referent of b
>>> print(a, b)
[<ListProxy object, typeid 'list' at ...>] []
>>> b.append('hello')
>>> print(a[0], b)
['hello'] ['hello']
類似地,字典和列表代理也可以相互嵌套:
>>> l_outer = manager.list([ manager.dict() for i in range(2) ])
>>> d_first_inner = l_outer[0]
>>> d_first_inner['a'] = 1
>>> d_first_inner['b'] = 2
>>> l_outer[1]['c'] = 3
>>> l_outer[1]['z'] = 26
>>> print(l_outer[0])
{'a': 1, 'b': 2}
>>> print(l_outer[1])
{'c': 3, 'z': 26}
如果指涉對象包含了普通 list 或 dict 對象,對這些內(nèi)部可變對象的修改不會通過管理器傳播,因為代理無法得知被包含的值什么時候被修改了。但是把存放在容器代理中的值本身是會通過管理器傳播的(會觸發(fā)代理對象中的 __setitem__ )從而有效修改這些對象,所以可以把修改過的值重新賦值給容器代理:
# create a list proxy and append a mutable object (a dictionary)
lproxy = manager.list()
lproxy.append({})
# now mutate the dictionary
d = lproxy[0]
d['a'] = 1
d['b'] = 2
# at this point, the changes to d are not yet synced, but by
# updating the dictionary, the proxy is notified of the change
lproxy[0] = d
在大多是使用情形下,這種實現(xiàn)方式并不比嵌套 代理對象 方便,但是依然演示了對于同步的一種控制級別。
注解
multiprocessing 中的代理類并沒有提供任何對于代理值比較的支持。所以,我們會得到如下結(jié)果:
>>> manager.list([1,2,3]) == [1,2,3]
False
當(dāng)需要比較值的時候,應(yīng)該替換為使用指涉對象的拷貝。
-
class
multiprocessing.managers.BaseProxy? 代理對象是
BaseProxy派生類的實例。-
_callmethod(methodname[, args[, kwds]])? 調(diào)用指涉對象的方法并返回結(jié)果。
如果
proxy是一個代理且其指涉的是obj, 那么下面的表達(dá)式:proxy._callmethod(methodname, args, kwds)
相當(dāng)于求取以下表達(dá)式的值:
getattr(obj, methodname)(*args, **kwds)
于管理器進(jìn)程。
返回結(jié)果會是一個值拷貝或者一個新的共享對象的代理 - 見函數(shù)
BaseManager.register()中關(guān)于參數(shù) method_to_typeid 的文檔。如果這個調(diào)用熬出了異常,則這個異常會被
_callmethod()透傳出來。如果是管理器進(jìn)程本身拋出的一些其他異常,則會被_callmethod()轉(zhuǎn)換為RemoteError異常重新拋出。特別注意,如果 methodname 沒有 暴露 出來,將會引發(fā)一個異常。
_callmethod()的一個使用示例:>>> l = manager.list(range(10)) >>> l._callmethod('__len__') 10 >>> l._callmethod('__getitem__', (slice(2, 7),)) # equivalent to l[2:7] [2, 3, 4, 5, 6] >>> l._callmethod('__getitem__', (20,)) # equivalent to l[20] Traceback (most recent call last): ... IndexError: list index out of range
-
_getvalue()? 返回指涉對象的一份拷貝。
如果指涉對象無法序列化,則會拋出一個異常。
-
__repr__()? 返回代理對象的字符串表示。
-
__str__()? 返回指涉對象的字符串表示。
-
清理?
代理對象使用了一個弱引用回調(diào),當(dāng)它被垃圾回收時,會將自己從擁有此指涉對象的管理器上反注冊,
當(dāng)共享對象沒有被任何代理器引用時,會被管理器進(jìn)程刪除。
進(jìn)程池?
可以創(chuàng)建一個進(jìn)程池,它將使用 Pool 類執(zhí)行提交給它的任務(wù)。
-
class
multiprocessing.pool.Pool([processes[, initializer[, initargs[, maxtasksperchild[, context]]]]])? 一個進(jìn)程池對象,它控制可以提交作業(yè)的工作進(jìn)程池。它支持帶有超時和回調(diào)的異步結(jié)果,以及一個并行的 map 實現(xiàn)。
processes 是要使用的工作進(jìn)程數(shù)目。如果 processes 為
None,則使用os.cpu_count()返回的值。如果 initializer 不為
None,則每個工作進(jìn)程將會在啟動時調(diào)用initializer(*initargs)。maxtasksperchild 是一個工作進(jìn)程在它退出或被一個新的工作進(jìn)程代替之前能完成的任務(wù)數(shù)量,為了釋放未使用的資源。默認(rèn)的 maxtasksperchild 是
None,意味著工作進(jìn)程壽與池齊。context 可被用于指定啟動的工作進(jìn)程的上下文。通常一個進(jìn)程池是使用函數(shù)
multiprocessing.Pool()或者一個上下文對象的Pool()方法創(chuàng)建的。在這兩種情況下, context 都是適當(dāng)設(shè)置的。注意,進(jìn)程池對象的方法只有創(chuàng)建它的進(jìn)程能夠調(diào)用。
警告
multiprocessing.pool對象具有需要正確管理的內(nèi)部資源 (像任何其他資源一樣),具體方式是將進(jìn)程池用作上下文管理器,或者手動調(diào)用close()和terminate()。 未做此類操作將導(dǎo)致進(jìn)程在終結(jié)階段掛起。請注意依賴?yán)厥掌鱽礓N毀進(jìn)程池是 不正確 的做法,因為 CPython 并不保證會調(diào)用進(jìn)程池終結(jié)程序(請參閱
object.__del__()來了解詳情)。3.2 新版功能: maxtasksperchild
3.4 新版功能: context
注解
通常來說,
Pool中的 Worker 進(jìn)程的生命周期和進(jìn)程池的工作隊列一樣長。一些其他系統(tǒng)中(如 Apache, mod_wsgi 等)也可以發(fā)現(xiàn)另一種模式,他們會讓工作進(jìn)程在完成一些任務(wù)后退出,清理、釋放資源,然后啟動一個新的進(jìn)程代替舊的工作進(jìn)程。Pool的 maxtasksperchild 參數(shù)給用戶提供了這種能力。-
apply(func[, args[, kwds]])? 使用 args 參數(shù)以及 kwds 命名參數(shù)調(diào)用 func , 它會返回結(jié)果前阻塞。這種情況下,
apply_async()更適合并行化工作。另外 func 只會在一個進(jìn)程池中的一個工作進(jìn)程中執(zhí)行。
-
apply_async(func[, args[, kwds[, callback[, error_callback]]]])? apply()方法的一個變種,返回一個結(jié)果對象。如果指定了 callback , 它必須是一個接受單個參數(shù)的可調(diào)用對象。當(dāng)執(zhí)行成功時, callback 會被用于處理執(zhí)行后的返回結(jié)果,否則,調(diào)用 error_callback 。
如果指定了 error_callback , 它必須是一個接受單個參數(shù)的可調(diào)用對象。當(dāng)目標(biāo)函數(shù)執(zhí)行失敗時, 會將拋出的異常對象作為參數(shù)傳遞給 error_callback 執(zhí)行。
回調(diào)函數(shù)應(yīng)該立即執(zhí)行完成,否則會阻塞負(fù)責(zé)處理結(jié)果的線程。
-
map(func, iterable[, chunksize])? 內(nèi)置
map()函數(shù)的并行版本 (但它只支持一個 iterable 參數(shù),對于多個可迭代對象請參閱starmap())。 它會保持阻塞直到獲得結(jié)果。這個方法會將可迭代對象分割為許多塊,然后提交給進(jìn)程池。可以將 chunksize 設(shè)置為一個正整數(shù)從而(近似)指定每個塊的大小可以。
注意對于很長的迭代對象,可能消耗很多內(nèi)存。可以考慮使用
imap()或imap_unordered()并且顯示指定 chunksize 以提升效率。
-
map_async(func, iterable[, chunksize[, callback[, error_callback]]])? 和
map()方法類似,但是返回一個結(jié)果對象。如果指定了 callback , 它必須是一個接受單個參數(shù)的可調(diào)用對象。當(dāng)執(zhí)行成功時, callback 會被用于處理執(zhí)行后的返回結(jié)果,否則,調(diào)用 error_callback 。
如果指定了 error_callback , 它必須是一個接受單個參數(shù)的可調(diào)用對象。當(dāng)目標(biāo)函數(shù)執(zhí)行失敗時, 會將拋出的異常對象作為參數(shù)傳遞給 error_callback 執(zhí)行。
回調(diào)函數(shù)應(yīng)該立即執(zhí)行完成,否則會阻塞負(fù)責(zé)處理結(jié)果的線程。
-
imap(func, iterable[, chunksize])? map()的延遲執(zhí)行版本。chunksize 參數(shù)的作用和
map()方法的一樣。對于很長的迭代器,給 chunksize 設(shè)置一個很大的值會比默認(rèn)值1極大 地加快執(zhí)行速度。同樣,如果 chunksize 是
1, 那么imap()方法所返回的迭代器的next()方法擁有一個可選的 timeout 參數(shù): 如果無法在 timeout 秒內(nèi)執(zhí)行得到結(jié)果,則``next(timeout)`` 會拋出multiprocessing.TimeoutError異常。
-
imap_unordered(func, iterable[, chunksize])? 和
imap()相同,只不過通過迭代器返回的結(jié)果是任意的。(當(dāng)進(jìn)程池中只有一個工作進(jìn)程的時候,返回結(jié)果的順序才能認(rèn)為是"有序"的)
-
starmap(func, iterable[, chunksize])? 和
map()類似,不過 iterable 中的每一項會被解包再作為函數(shù)參數(shù)。比如可迭代對象
[(1,2), (3, 4)]會轉(zhuǎn)化為等價于[func(1,2), func(3,4)]的調(diào)用。3.3 新版功能.
-
starmap_async(func, iterable[, chunksize[, callback[, error_callback]]])? 相當(dāng)于
starmap()與map_async()的結(jié)合,迭代 iterable 的每一項,解包作為 func 的參數(shù)并執(zhí)行,返回用于獲取結(jié)果的對象。3.3 新版功能.
-
close()? 阻止后續(xù)任務(wù)提交到進(jìn)程池,當(dāng)所有任務(wù)執(zhí)行完成后,工作進(jìn)程會退出。
-
terminate()? 不必等待未完成的任務(wù),立即停止工作進(jìn)程。當(dāng)進(jìn)程池對象被垃圾回收時,會立即調(diào)用
terminate()。
-
join()? 等待工作進(jìn)程結(jié)束。調(diào)用
join()前必須先調(diào)用close()或者terminate()。
3.3 新版功能: 進(jìn)程池對象現(xiàn)在支持上下文管理器協(xié)議 - 參見 上下文管理器類型 。
__enter__()返回進(jìn)程池對象,__exit__()會調(diào)用terminate()。-
-
class
multiprocessing.pool.AsyncResult? Pool.apply_async()和Pool.map_async()返回對象所屬的類。-
get([timeout])? 用于獲取執(zhí)行結(jié)果。如果 timeout 不是
None并且在 timeout 秒內(nèi)仍然沒有執(zhí)行完得到結(jié)果,則拋出multiprocessing.TimeoutError異常。如果遠(yuǎn)程調(diào)用發(fā)生異常,這個異常會通過get()重新拋出。
-
wait([timeout])? 阻塞,直到返回結(jié)果,或者 timeout 秒后超時。
-
ready()? 用于判斷執(zhí)行狀態(tài),是否已經(jīng)完成。
-
successful()? 判斷調(diào)用是否已經(jīng)完成并且未引發(fā)異常。 如果還未獲得結(jié)果則將引發(fā)
ValueError。
-
下面的例子演示了進(jìn)程池的用法:
from multiprocessing import Pool
import time
def f(x):
return x*x
if __name__ == '__main__':
with Pool(processes=4) as pool: # start 4 worker processes
result = pool.apply_async(f, (10,)) # evaluate "f(10)" asynchronously in a single process
print(result.get(timeout=1)) # prints "100" unless your computer is *very* slow
print(pool.map(f, range(10))) # prints "[0, 1, 4,..., 81]"
it = pool.imap(f, range(10))
print(next(it)) # prints "0"
print(next(it)) # prints "1"
print(it.next(timeout=1)) # prints "4" unless your computer is *very* slow
result = pool.apply_async(time.sleep, (10,))
print(result.get(timeout=1)) # raises multiprocessing.TimeoutError
監(jiān)聽者及客戶端?
通常情況下,進(jìn)程間通過隊列或者 Pipe() 返回的 Connection 傳遞消息。
不過,multiprocessing.connection 其實提供了一些更靈活的特性。最基礎(chǔ)的用法是通過它抽象出來的高級來操作socket或者Windows命名管道。也提供一些高級用法,如通過 hmac 模塊來支持 摘要認(rèn)證,以及同時監(jiān)聽多個管道連接。
-
multiprocessing.connection.deliver_challenge(connection, authkey)? 發(fā)送一個隨機(jī)生成的消息到另一端,并等待回復(fù)。
如果收到的回復(fù)與使用 authkey 生成的信息摘要匹配成功,就會發(fā)送一個歡迎信息給管道另一端。否則拋出
AuthenticationError異常。
-
multiprocessing.connection.answer_challenge(connection, authkey)? 接收一條信息,使用 authkey 作為鍵計算信息摘要,然后將摘要發(fā)送回去。
如果沒有收到歡迎消息,就拋出
AuthenticationError異常。
-
multiprocessing.connection.Client(address[, family[, authkey]])? 嘗試在監(jiān)聽者上使用 address 地址初始化一個連接,返回
Connection。連接的類型取決于 family 參數(shù),但是通??梢允÷裕驗榭梢酝ㄟ^ address 的格式推導(dǎo)出來。(查看 地址格式 )
如果提供了 authkey 參數(shù)并且不是 None,那它必須是一個字符串并且會被當(dāng)做基于 HMAC 認(rèn)證的密鑰。如果 authkey 是None 則不會有認(rèn)證行為。認(rèn)證失敗拋出
AuthenticationError異常,請查看 See 認(rèn)證密碼 。
-
class
multiprocessing.connection.Listener([address[, family[, backlog[, authkey]]]])? 可以監(jiān)聽連接請求,是對于綁定套接字或者 Windows 命名管道的封裝。
address 是監(jiān)聽器對象中的綁定套接字或命名管道使用的地址。
注解
如果使用 '0.0.0.0' 作為監(jiān)聽地址,那么在Windows上這個地址無法建立連接。想要建立一個可連接的端點,應(yīng)該使用 '127.0.0.1' 。
family 是套接字(或者命名管道)使用的類型。它可以是以下一種:
'AF_INET'( TCP 套接字類型),'AF_UNIX'( Unix 域套接字) 或者'AF_PIPE'( Windows 命名管道)。其中只有第一個保證各平臺可用。如果 family 是None,那么 family 會根據(jù) address 的格式自動推導(dǎo)出來。如果 address 也是None, 則取默認(rèn)值。默認(rèn)值為可用類型中速度最快的。見 地址格式 。注意,如果 family 是'AF_UNIX'而address是``None`` ,套接字會在一個tempfile.mkstemp()創(chuàng)建的私有臨時目錄中創(chuàng)建。如果監(jiān)聽器對象使用了套接字,backlog (默認(rèn)值為1) 會在套接字綁定后傳遞給它的
listen()方法。如果提供了 authkey 參數(shù)并且不是 None,那它必須是一個字符串并且會被當(dāng)做基于 HMAC 認(rèn)證的密鑰。如果 authkey 是None 則不會有認(rèn)證行為。認(rèn)證失敗拋出
AuthenticationError異常,請查看 See 認(rèn)證密碼 。-
accept()? 接受一個連接并返回一個
Connection對象,其連接到的監(jiān)聽器對象已綁定套接字或者命名管道。如果已經(jīng)嘗試過認(rèn)證并且失敗了,則會拋出AuthenticationError異常。
-
close()? 關(guān)閉監(jiān)聽器上的綁定套接字或者命名管道。此函數(shù)會在監(jiān)聽器被垃圾回收后自動調(diào)用。不過仍然建議顯式調(diào)用函數(shù)關(guān)閉。
監(jiān)聽器對象擁有下列只讀屬性:
-
address? 被監(jiān)聽器對象使用的地址。
-
last_accepted? 最后一個連接所使用的地址。如果沒有的話就是
None。
3.3 新版功能: 監(jiān)聽器對象現(xiàn)在支持了上下文管理協(xié)議 - 見 上下文管理器類型 。
__enter__()返回一個監(jiān)聽器對象,__exit__()會調(diào)用close()。-
-
multiprocessing.connection.wait(object_list, timeout=None)? 一直等待直到 object_list 中某個對象處于就緒狀態(tài)。返回 object_list 中處于就緒狀態(tài)的對象。如果 timeout 是一個浮點型,該方法會最多阻塞這么多秒。如果 timeout 是
None,則會允許阻塞的事件沒有限制。timeout為負(fù)數(shù)的情況下和為0的情況相同。對于 Unix 和 Windows ,下列對象都可以出現(xiàn)在 object_list 中
可讀的
Connection對象;一個已連接并且可讀的
socket.socket對象;或者
當(dāng)一個連接或者套接字對象擁有有效的數(shù)據(jù)可被讀取的時候,或者另一端關(guān)閉后,這個對象就處于就緒狀態(tài)。
Unix:
wait(object_list, timeout)和select.select(object_list, [], [], timeout)幾乎相同。差別在于,如果select.select()被信號中斷,它會拋出一個附帶錯誤號為EINTR的OSError異常,而wait()不會。Windows: object_list 中的元素必須是一個表示為整數(shù)的可等待的句柄(按照 Win32 函數(shù)
WaitForMultipleObjects()的文檔中所定義) 或者一個擁有fileno()方法的對象,這個對象返回一個套接字句柄或者管道句柄。(注意管道和套接字兩種句柄 不是 可等待的句柄)3.3 新版功能.
示例
下面的服務(wù)代碼創(chuàng)建了一個使用 'secret password' 作為認(rèn)證密碼的監(jiān)聽器。它會等待連接然后發(fā)送一些數(shù)據(jù)給客戶端:
from multiprocessing.connection import Listener
from array import array
address = ('localhost', 6000) # family is deduced to be 'AF_INET'
with Listener(address, authkey=b'secret password') as listener:
with listener.accept() as conn:
print('connection accepted from', listener.last_accepted)
conn.send([2.25, None, 'junk', float])
conn.send_bytes(b'hello')
conn.send_bytes(array('i', [42, 1729]))
下面的代碼連接到服務(wù)然后從服務(wù)器上j接收一些數(shù)據(jù):
from multiprocessing.connection import Client
from array import array
address = ('localhost', 6000)
with Client(address, authkey=b'secret password') as conn:
print(conn.recv()) # => [2.25, None, 'junk', float]
print(conn.recv_bytes()) # => 'hello'
arr = array('i', [0, 0, 0, 0, 0])
print(conn.recv_bytes_into(arr)) # => 8
print(arr) # => array('i', [42, 1729, 0, 0, 0])
下面的代碼使用了 wait() ,以便在同時等待多個進(jìn)程發(fā)來消息。
import time, random
from multiprocessing import Process, Pipe, current_process
from multiprocessing.connection import wait
def foo(w):
for i in range(10):
w.send((i, current_process().name))
w.close()
if __name__ == '__main__':
readers = []
for i in range(4):
r, w = Pipe(duplex=False)
readers.append(r)
p = Process(target=foo, args=(w,))
p.start()
# We close the writable end of the pipe now to be sure that
# p is the only process which owns a handle for it. This
# ensures that when p closes its handle for the writable end,
# wait() will promptly report the readable end as being ready.
w.close()
while readers:
for r in wait(readers):
try:
msg = r.recv()
except EOFError:
readers.remove(r)
else:
print(msg)
地址格式?
'AF_INET'地址是(主機(jī), 端口)? 形式的元組類型,其中 主機(jī) 是一個字符串,端口 是整數(shù)。'AF_UNIX'地址是文件系統(tǒng)上文件名的字符串。'AF_PIPE'是這種格式的字符串r'\.\pipe{PipeName}'。如果要用Client()連接到一個名為 ServerName 的遠(yuǎn)程命名管道,應(yīng)該替換為使用r'\ServerName\pipe{PipeName}'這種格式。
注意,使用兩個反斜線開頭的字符串默認(rèn)被當(dāng)做 'AF_PIPE' 地址而不是 'AF_UNIX'? 。
認(rèn)證密碼?
當(dāng)使用 Connection.recv 接收數(shù)據(jù)時,數(shù)據(jù)會自動被反序列化。不幸的是,對于一個不可信的數(shù)據(jù)源發(fā)來的數(shù)據(jù),反序列化是存在安全風(fēng)險的。所以 Listener 和 Client() 之間使用 hmac 模塊進(jìn)行摘要認(rèn)證。
認(rèn)證密鑰是一個 byte 類型的字符串,可以認(rèn)為是和密碼一樣的東西,連接建立好后,雙方都會要求另一方證明知道認(rèn)證密鑰。(這個證明過程不會通過連接發(fā)送密鑰)
如果要求認(rèn)證但是沒有指定認(rèn)證密鑰,則會使用 current_process().authkey 的返回值 (參見 Process)。 這個值將被當(dāng)前進(jìn)程所創(chuàng)建的任何 Process 對象自動繼承。 這意味著 (在默認(rèn)情況下) 一個包含多進(jìn)程的程序中的所有進(jìn)程會在相互間建立連接的時候共享單個認(rèn)證密鑰。
os.urandom() 也可以用來生成合適的認(rèn)證密鑰。
日志記錄?
當(dāng)前模塊也提供了一些對 logging 的支持。注意, logging 模塊本身并沒有使用進(jìn)程間共享的鎖,所以來自于多個進(jìn)程的日志可能(具體取決于使用的日志 handler)相互覆蓋或者混雜。
-
multiprocessing.get_logger()? 返回
multiprocessing使用的 logger,必要的話會創(chuàng)建一個新的。如果創(chuàng)建的首個 logger 日志級別為
logging.NOTSET并且沒有默認(rèn) handler。通過這個 logger 打印的消息不會傳遞到根 logger。注意在 Windows 上,子進(jìn)程只會繼承父進(jìn)程 logger 的日志級別 - 對于logger的其他自定義項不會繼承。
-
multiprocessing.log_to_stderr()? 此函數(shù)會調(diào)用
get_logger()但是會在返回的 logger 上增加一個 handler,將所有輸出都使用'[%(levelname)s/%(processName)s] %(message)s'的格式發(fā)送到sys.stderr。
下面是一個在交互式解釋器中打開日志功能的例子:
>>> import multiprocessing, logging
>>> logger = multiprocessing.log_to_stderr()
>>> logger.setLevel(logging.INFO)
>>> logger.warning('doomed')
[WARNING/MainProcess] doomed
>>> m = multiprocessing.Manager()
[INFO/SyncManager-...] child process calling self.run()
[INFO/SyncManager-...] created temp directory /.../pymp-...
[INFO/SyncManager-...] manager serving at '/.../listener-...'
>>> del m
[INFO/MainProcess] sending shutdown message to manager
[INFO/SyncManager-...] manager exiting with exitcode 0
要查看日志等級的完整列表,見 logging 模塊。
multiprocessing.dummy 模塊?
multiprocessing.dummy 復(fù)制了 multiprocessing 的 API,不過是在 threading 模塊之上包裝了一層。
特別地,multiprocessing.dummy 所提供的 Pool 函數(shù)會返回一個 ThreadPool 的實例,該類是 Pool 的子類,它支持所有相同的方法調(diào)用但會使用一個工作線程池而非工作進(jìn)程池。
-
class
multiprocessing.pool.ThreadPool([processes[, initializer[, initargs]]])? 一個線程池對象,用來控制可向其提交任務(wù)的工作線程池。
ThreadPool實例與Pool實例是完全接口兼容的,并且它們的資源也必須被正確地管理,或者是將線程池作為上下文管理器來使用,或者是通過手動調(diào)用close()和terminate()。processes 是要使用的工作線程數(shù)目。 如果 processes 為
None,則使用os.cpu_count()返回的值。如果 initializer 不為
None,則每個工作進(jìn)程將會在啟動時調(diào)用initializer(*initargs)。不同于
Pool,maxtasksperchild 和 context 不可被提供。注解
ThreadPool具有與Pool相同的接口,它圍繞一個進(jìn)程池進(jìn)行設(shè)計并且先于concurrent.futures模塊的引入。 因此,它繼承了一些對于基于線程的池來說沒有意義的操作,并且它具有自己的用于表示異步任務(wù)狀態(tài)的類型AsyncResult,該類型不為任何其他庫所知。用戶通常應(yīng)該傾向于使用
concurrent.futures.ThreadPoolExecutor,它擁有從一開始就圍繞線程進(jìn)行設(shè)計的更簡單接口,并且返回與許多其他庫相兼容的concurrent.futures.Future實例,包括asyncio庫。
編程指導(dǎo)?
使用 multiprocessing 時,應(yīng)遵循一些指導(dǎo)原則和習(xí)慣用法。
所有啟動方法?
下面這些使用于所有啟動方法。
避免共享狀態(tài)
應(yīng)該盡可能避免在進(jìn)程間傳遞大量數(shù)據(jù),越少越好。
最好堅持使用隊列或者管道進(jìn)行進(jìn)程間通信,而不是底層的同步原語。
可序列化
保證所代理的方法的參數(shù)是可以序列化的。
代理的線程安全性
不要在多線程中同時使用一個代理對象,除非你用鎖保護(hù)它。
(而在不同進(jìn)程中使用 相同 的代理對象從不會發(fā)生問題。)
使用 Join 避免僵尸進(jìn)程
在 Unix 上,如果一個進(jìn)程執(zhí)行完成但是沒有被 join,就會變成僵尸進(jìn)程。一般來說,僵尸進(jìn)程不會很多,因為每次新啟動進(jìn)程(或者
active_children()被調(diào)用)時,所有已執(zhí)行完成且沒有被 join 的進(jìn)程都會自動被 join,而且對一個執(zhí)行完的進(jìn)程調(diào)用Process.is_alive也會 join 這個進(jìn)程。盡管如此,對自己啟動的進(jìn)程顯式調(diào)用 join 依然是最佳實踐。
繼承優(yōu)于序列化、反序列化
當(dāng)使用 spawn 或者 forkserver 的啟動方式時,
multiprocessing中的許多類型都必須是可序列化的,這樣子進(jìn)程才能使用它們。但是通常我們都應(yīng)該避免使用管道和隊列發(fā)送共享對象到另外一個進(jìn)程,而是重新組織代碼,對于其他進(jìn)程創(chuàng)建出來的共享對象,讓那些需要訪問這些對象的子進(jìn)程可以直接將這些對象從父進(jìn)程繼承過來。
避免殺死進(jìn)程
聽過
Process.terminate? 停止一個進(jìn)程很容易導(dǎo)致這個進(jìn)程正在使用的共享資源(如鎖、信號量、管道和隊列)損壞或者變得不可用,無法在其他進(jìn)程中繼續(xù)使用。所以,最好只對那些從來不使用共享資源的進(jìn)程調(diào)用
Process.terminate。
Join 使用隊列的進(jìn)程
記住,往隊列放入數(shù)據(jù)的進(jìn)程會一直等待直到所有緩存項被"feeder" 線程傳給底層管道。(子進(jìn)程可以調(diào)用隊列的
Queue.cancel_join_thread方法禁止這種行為)這意味著,任何使用隊列的時候,你都要確保在進(jìn)程join之前,所有存放到隊列中的項將會被其他進(jìn)程、線程完全消費。否則不能保證這個寫過隊列的進(jìn)程可以正常終止。記住非精靈進(jìn)程會自動 join 。
下面是一個會導(dǎo)致死鎖的例子:
from multiprocessing import Process, Queue def f(q): q.put('X' * 1000000) if __name__ == '__main__': queue = Queue() p = Process(target=f, args=(queue,)) p.start() p.join() # this deadlocks obj = queue.get()交換最后兩行可以修復(fù)這個問題(或者直接刪掉
p.join())。
顯式傳遞資源給子進(jìn)程
在Unix上,使用 fork 方式啟動的子進(jìn)程可以使用父進(jìn)程中全局創(chuàng)建的共享資源。
除了(部分原因)讓代碼兼容 Windows 以及其他的進(jìn)程啟動方式外,這種形式還保證了在子進(jìn)程生命期這個對象是不會被父進(jìn)程垃圾回收的。如果父進(jìn)程中的某些對象被垃圾回收會導(dǎo)致資源釋放,這就變得很重要。
所以對于實例:
from multiprocessing import Process, Lock def f(): ... do something using "lock" ... if __name__ == '__main__': lock = Lock() for i in range(10): Process(target=f).start()應(yīng)當(dāng)重寫成這樣:
from multiprocessing import Process, Lock def f(l): ... do something using "l" ... if __name__ == '__main__': lock = Lock() for i in range(10): Process(target=f, args=(lock,)).start()
謹(jǐn)防將 sys.stdin 數(shù)據(jù)替換為 “類似文件的對象”
multiprocessing內(nèi)部會無條件地這樣調(diào)用:os.close(sys.stdin.fileno())在
multiprocessing.Process._bootstrap()? 方法中 —— 這會導(dǎo)致與"進(jìn)程中的進(jìn)程"相關(guān)的一些問題。這已經(jīng)被修改成了:sys.stdin.close() sys.stdin = open(os.open(os.devnull, os.O_RDONLY), closefd=False)它解決了進(jìn)程相互沖突導(dǎo)致文件描述符錯誤的根本問題,但是對使用帶緩沖的“文件類對象”替換
sys.stdin()作為輸出的應(yīng)用程序造成了潛在的危險。如果多個進(jìn)程調(diào)用了此文件類對象的close()方法,會導(dǎo)致相同的數(shù)據(jù)多次刷寫到此對象,損壞數(shù)據(jù)。如果你寫入文件類對象并實現(xiàn)了自己的緩存,可以在每次追加緩存數(shù)據(jù)時記錄當(dāng)前進(jìn)程id,從而將其變成 fork 安全的,當(dāng)發(fā)現(xiàn)進(jìn)程id變化后舍棄之前的緩存,例如:
@property def cache(self): pid = os.getpid() if pid != self._pid: self._pid = pid self._cache = [] return self._cache
spawn 和 forkserver 啟動方式?
相對于 fork 啟動方式,有一些額外的限制。
更依賴序列化
Process.__init__()的所有參數(shù)都必須可序列化。同樣的,當(dāng)你繼承Process時,需要保證當(dāng)調(diào)用Process.start方法時,實例可以被序列化。
全局變量
記住,如果子進(jìn)程中的代碼嘗試訪問一個全局變量,它所看到的值可能和父進(jìn)程中執(zhí)行
Process.start那一刻的值不一樣。當(dāng)全局變量知識模塊級別的常量時,是不會有問題的。
安全導(dǎo)入主模塊
確保主模塊可以被新啟動的Python解釋器安全導(dǎo)入而不會引發(fā)什么副作用(比如又啟動了一個子進(jìn)程)
例如,使用 spawn 或 forkserver 啟動方式執(zhí)行下面的模塊,會引發(fā)
RuntimeError異常而失敗。from multiprocessing import Process def foo(): print('hello') p = Process(target=foo) p.start()應(yīng)該通過下面的方法使用
if __name__ == '__main__':,從而保護(hù)程序"入口點":from multiprocessing import Process, freeze_support, set_start_method def foo(): print('hello') if __name__ == '__main__': freeze_support() set_start_method('spawn') p = Process(target=foo) p.start()(如果程序?qū)⒄_\行而不是凍結(jié),則可以省略
freeze_support()行)這允許新啟動的 Python 解釋器安全導(dǎo)入模塊然后運行模塊中的
foo()函數(shù)。如果主模塊中創(chuàng)建了進(jìn)程池或者管理器,這個規(guī)則也適用。
例子?
創(chuàng)建和使用自定義管理器、代理的示例:
from multiprocessing import freeze_support
from multiprocessing.managers import BaseManager, BaseProxy
import operator
##
class Foo:
def f(self):
print('you called Foo.f()')
def g(self):
print('you called Foo.g()')
def _h(self):
print('you called Foo._h()')
# A simple generator function
def baz():
for i in range(10):
yield i*i
# Proxy type for generator objects
class GeneratorProxy(BaseProxy):
_exposed_ = ['__next__']
def __iter__(self):
return self
def __next__(self):
return self._callmethod('__next__')
# Function to return the operator module
def get_operator_module():
return operator
##
class MyManager(BaseManager):
pass
# register the Foo class; make `f()` and `g()` accessible via proxy
MyManager.register('Foo1', Foo)
# register the Foo class; make `g()` and `_h()` accessible via proxy
MyManager.register('Foo2', Foo, exposed=('g', '_h'))
# register the generator function baz; use `GeneratorProxy` to make proxies
MyManager.register('baz', baz, proxytype=GeneratorProxy)
# register get_operator_module(); make public functions accessible via proxy
MyManager.register('operator', get_operator_module)
##
def test():
manager = MyManager()
manager.start()
print('-' * 20)
f1 = manager.Foo1()
f1.f()
f1.g()
assert not hasattr(f1, '_h')
assert sorted(f1._exposed_) == sorted(['f', 'g'])
print('-' * 20)
f2 = manager.Foo2()
f2.g()
f2._h()
assert not hasattr(f2, 'f')
assert sorted(f2._exposed_) == sorted(['g', '_h'])
print('-' * 20)
it = manager.baz()
for i in it:
print('<%d>' % i, end=' ')
print()
print('-' * 20)
op = manager.operator()
print('op.add(23, 45) =', op.add(23, 45))
print('op.pow(2, 94) =', op.pow(2, 94))
print('op._exposed_ =', op._exposed_)
##
if __name__ == '__main__':
freeze_support()
test()
使用 Pool:
import multiprocessing
import time
import random
import sys
#
# Functions used by test code
#
def calculate(func, args):
result = func(*args)
return '%s says that %s%s = %s' % (
multiprocessing.current_process().name,
func.__name__, args, result
)
def calculatestar(args):
return calculate(*args)
def mul(a, b):
time.sleep(0.5 * random.random())
return a * b
def plus(a, b):
time.sleep(0.5 * random.random())
return a + b
def f(x):
return 1.0 / (x - 5.0)
def pow3(x):
return x ** 3
def noop(x):
pass
#
# Test code
#
def test():
PROCESSES = 4
print('Creating pool with %d processes\n' % PROCESSES)
with multiprocessing.Pool(PROCESSES) as pool:
#
# Tests
#
TASKS = [(mul, (i, 7)) for i in range(10)] + \
[(plus, (i, 8)) for i in range(10)]
results = [pool.apply_async(calculate, t) for t in TASKS]
imap_it = pool.imap(calculatestar, TASKS)
imap_unordered_it = pool.imap_unordered(calculatestar, TASKS)
print('Ordered results using pool.apply_async():')
for r in results:
print('\t', r.get())
print()
print('Ordered results using pool.imap():')
for x in imap_it:
print('\t', x)
print()
print('Unordered results using pool.imap_unordered():')
for x in imap_unordered_it:
print('\t', x)
print()
print('Ordered results using pool.map() --- will block till complete:')
for x in pool.map(calculatestar, TASKS):
print('\t', x)
print()
#
# Test error handling
#
print('Testing error handling:')
try:
print(pool.apply(f, (5,)))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from pool.apply()')
else:
raise AssertionError('expected ZeroDivisionError')
try:
print(pool.map(f, list(range(10))))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from pool.map()')
else:
raise AssertionError('expected ZeroDivisionError')
try:
print(list(pool.imap(f, list(range(10)))))
except ZeroDivisionError:
print('\tGot ZeroDivisionError as expected from list(pool.imap())')
else:
raise AssertionError('expected ZeroDivisionError')
it = pool.imap(f, list(range(10)))
for i in range(10):
try:
x = next(it)
except ZeroDivisionError:
if i == 5:
pass
except StopIteration:
break
else:
if i == 5:
raise AssertionError('expected ZeroDivisionError')
assert i == 9
print('\tGot ZeroDivisionError as expected from IMapIterator.next()')
print()
#
# Testing timeouts
#
print('Testing ApplyResult.get() with timeout:', end=' ')
res = pool.apply_async(calculate, TASKS[0])
while 1:
sys.stdout.flush()
try:
sys.stdout.write('\n\t%s' % res.get(0.02))
break
except multiprocessing.TimeoutError:
sys.stdout.write('.')
print()
print()
print('Testing IMapIterator.next() with timeout:', end=' ')
it = pool.imap(calculatestar, TASKS)
while 1:
sys.stdout.flush()
try:
sys.stdout.write('\n\t%s' % it.next(0.02))
except StopIteration:
break
except multiprocessing.TimeoutError:
sys.stdout.write('.')
print()
print()
if __name__ == '__main__':
multiprocessing.freeze_support()
test()
一個演示如何使用隊列來向一組工作進(jìn)程提供任務(wù)并收集結(jié)果的例子:
import time
import random
from multiprocessing import Process, Queue, current_process, freeze_support
#
# Function run by worker processes
#
def worker(input, output):
for func, args in iter(input.get, 'STOP'):
result = calculate(func, args)
output.put(result)
#
# Function used to calculate result
#
def calculate(func, args):
result = func(*args)
return '%s says that %s%s = %s' % \
(current_process().name, func.__name__, args, result)
#
# Functions referenced by tasks
#
def mul(a, b):
time.sleep(0.5*random.random())
return a * b
def plus(a, b):
time.sleep(0.5*random.random())
return a + b
#
#
#
def test():
NUMBER_OF_PROCESSES = 4
TASKS1 = [(mul, (i, 7)) for i in range(20)]
TASKS2 = [(plus, (i, 8)) for i in range(10)]
# Create queues
task_queue = Queue()
done_queue = Queue()
# Submit tasks
for task in TASKS1:
task_queue.put(task)
# Start worker processes
for i in range(NUMBER_OF_PROCESSES):
Process(target=worker, args=(task_queue, done_queue)).start()
# Get and print results
print('Unordered results:')
for i in range(len(TASKS1)):
print('\t', done_queue.get())
# Add more tasks using `put()`
for task in TASKS2:
task_queue.put(task)
# Get and print some more results
for i in range(len(TASKS2)):
print('\t', done_queue.get())
# Tell child processes to stop
for i in range(NUMBER_OF_PROCESSES):
task_queue.put('STOP')
if __name__ == '__main__':
freeze_support()
test()
