multiprocessing --- 基于進(jìn)程的并行?

源代碼 Lib/multiprocessing/


概述?

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)程。 Processthreading.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上使用 spawnforkserver 啟動方法也將啟動一個 信號量跟蹤器 進(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)建的鎖不能傳遞給使用 spawnforkserver 啟動方法啟動的進(jìn)程。

想要使用特定啟動方法的庫應(yīng)該使用 get_context() 以避免干擾庫用戶的選擇。

警告

'spawn''forkserver' 啟動方法當(dāng)前不能在Unix上和“凍結(jié)的”可執(zhí)行內(nèi)容一同使用(例如,有類似 PyInstallercx_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)存

可以使用 ValueArray 將數(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)建 numarr 時使用的 '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 、 dictNamespace 、 Lock 、 RLockSemaphore 、 BoundedSemaphore 、 ConditionEvent 、 BarrierQueue 、 ValueArray 。例如

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è)置為 TrueFalse 。如果是 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ù)(如果有),分別從 argskwargs 參數(shù)中獲取順序和關(guān)鍵字參數(shù)。

start()?

啟動進(jìn)程活動。

每個進(jìn)程對象最多只能調(diào)用一次。它安排對象的 run() 方法在一個單獨的進(jìn)程中調(diào)用。

join([timeout])?

如果可選參數(shù) timeoutNone (默認(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 個孩子。

is_alive()?

返回進(jìn)程是否還活著。

粗略地說,從 start() 方法返回到子進(jìn)程終止之前,進(jìn)程對象仍處于活動狀態(tài)。

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.Thread API ,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)句柄,可以與 WaitForSingleObjectWaitForMultipleObjects 系列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.Emptyqueue.Full 異常去表示超時。 你需要從 queue 中導(dǎo)入它們,因為它們并不在 multiprocessing 的命名空間中。

注解

當(dāng)一個對象被放入一個隊列中時,這個對象首先會被一個后臺線程用 pickle 序列化,并將序列化后的數(shù)據(jù)通過一個底層管道的管道傳遞到隊列中。 這種做法會有點讓人驚訝,但一般不會出現(xiàn)什么問題。 如果它們確實妨礙了你,你可以使用一個由管理器 manager 創(chuàng)建的隊列替換它。

  1. 將一個對象放入一個空隊列后,可能需要極小的延遲,隊列的方法 empty()? 才會返回 False 。而 get_nowait() 可以不拋出 queue.Empty 直接返回。

  2. 如果有多個進(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.Emptyqueue.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ù) blockTrue (默認(rèn)值) 而且 timeoutNone (默認(rèn)值), 將會阻塞當(dāng)前進(jìn)程,直到有空的緩沖槽。如果 timeout 是正數(shù),將會在阻塞了最多 timeout 秒之后還是沒有可用的緩沖槽時拋出 queue.Full? 異常。反之 (blockFalse 時),僅當(dāng)有可用緩沖槽時才放入對象,否則拋出 queue.Full 異常 (在這種情形下 timeout 參數(shù)會被忽略)。

put_nowait(obj)?

相當(dāng)于 put(obj, False)。

get([block[, timeout]])?

從隊列中取出并返回對象。如果可選參數(shù) blockTrue (默認(rèn)值) 而且 timeoutNone (默認(rèn)值), 將會阻塞當(dāng)前進(jìn)程,直到隊列中出現(xiàn)可用的對象。如果 timeout 是正數(shù),將會在阻塞了最多 timeout 秒之后還是沒有可用的對象時拋出 queue.Empty 異常。反之 (blockFalse 時),僅當(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.current_process()?

返回與當(dāng)前進(jìn)程相對應(yīng)的 Process 對象。

threading.current_thread() 相同。

multiprocessing.freeze_support()?

為使用了 multiprocessing? 的程序,提供凍結(jié)以產(chǎn)生 Windows 可執(zhí)行文件的支持。(在 py2exe, PyInstallercx_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 新版功能.

連接對象(Connection)?

Connection 對象允許收發(fā)可以序列化的對象或字符串。它們可以看作面向消息的連接套接字。

通常使用 Pipe 創(chuàng)建 Connection 對象。詳見 : 監(jiān)聽者及客戶端.

class multiprocessing.connection.Connection?
send(obj)?

將一個對象發(fā)送到連接的另一端,可以用 recv() 讀取。

發(fā)送的對象必須是可以序列化的,過大的對象 ( 接近 32MiB+ ,這個值取決于操作系統(tǒng) ) 有可能引發(fā) ValueError? 異常。

recv()?

返回一個由另一端使用 send() 發(fā)送的對象。該方法會一直阻塞直到接收到對象。 如果對端關(guān)閉了連接或者沒有東西可接收,將拋出 EOFError? 異常。

fileno()?

返回由連接對象使用的描述符或者句柄。

close()?

關(guān)閉連接對象。

當(dāng)連接對象被垃圾回收時會自動調(diào)用。

poll([timeout])?

返回連接對象中是否有可以讀取的數(shù)據(jù)。

如果未指定 timeout ,此方法會馬上返回。如果 timeout 是一個數(shù)字,則指定了最大阻塞的秒數(shù)。如果 timeoutNone? ,那么將一直等待,不會超時。

注意通過使用 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 并且該連接對象將不再可讀。

在 3.3 版更改: 曾經(jīng)該函數(shù)拋出 IOError ,現(xiàn)在這是 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? 異常,并且完整的消息將會存放在異常實例 ee.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? 對象。

Lock supports the context manager protocol and thus may be used in with statements.

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) timeoutNone (默認(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)過了同步包裝器包裝過的。可以通過 Valuevalue 屬性訪問這個對象本身。

typecode_or_type 指明了返回的對象類型: 它可能是一個 ctypes 類型或者 array? 模塊中每個類型對應(yīng)的單字符長度的字符串。 *args 會透傳給這個類的構(gòu)造函數(shù)。

如果 lock 參數(shù)是 True (默認(rèn)值), 將會新建一個遞歸鎖用于同步對于此值的訪問操作。 如果 lockLock 或者 RLock 對象,那么這個傳入的鎖將會用于同步對這個值的訪問操作,如果 lockFalse , 那么對這個對象的訪問將沒有鎖保護(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ù)組的長度。

如果 lockTrue (默認(rèn)值) 則將創(chuàng)建一個新的鎖對象用于同步對值的訪問。 如果 lock 為一個 LockRLock 對象則該對象將被用于同步對值的訪問。 如果 lockFalse 則對返回對象的訪問將不會自動得到鎖的保護(hù),也就是說它不是“進(jìn)程安全的”。

請注意 lock 是一個僅限關(guān)鍵字參數(shù)。

請注意 ctypes.c_char 的數(shù)組具有 valueraw 屬性,允許被用來保存和提取字符串。

multiprocessing.sharedctypes 模塊?

multiprocessing.sharedctypes 模塊提供了一些函數(shù),用于分配來自共享內(nèi)存的、可被子進(jìn)程繼承的 ctypes 對象。

注解

雖然可以將指針存儲在共享內(nèi)存中,但請記住它所引用的是特定進(jìn)程地址空間中的位置。 而且,指針很可能在第二個進(jìn)程的上下文中無效,嘗試從第二個進(jìn)程對指針進(jìn)行解引用可能會導(dǎo)致崩潰。

multiprocessing.sharedctypes.RawArray(typecode_or_type, size_or_initializer)?

從共享內(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() 來借助其中的鎖保證操作的原子性。

multiprocessing.sharedctypes.RawValue(typecode_or_type, *args)?

從共享內(nèi)存中申請并返回一個 ctypes 對象。

typecode_or_type 指明了返回的對象類型: 它可能是一個 ctypes 類型或者 array? 模塊中每個類型對應(yīng)的單字符長度的字符串。 *args 會透傳給這個類的構(gòu)造函數(shù)。

注意對 value 的訪問、賦值操作可能是非原子操作 - 使用 Value() 來借助其中的鎖保證操作的原子性。

請注意 ctypes.c_char 的數(shù)組具有 valueraw 屬性,允許被用來保存和提取字符串 - 請查看 ctypes 文檔。

multiprocessing.sharedctypes.Array(typecode_or_type, size_or_initializer, *, lock=True)?

返回一個純 ctypes 數(shù)組, 或者在此之上經(jīng)過同步器包裝過的對象,這取決于 lock 參數(shù)的值,除此之外,和 RawArray() 一樣。

如果 lockTrue (默認(rèn)值) 則將創(chuàng)建一個新的鎖對象用于同步對值的訪問。 如果 lock 為一個 LockRLock 對象則該對象將被用于同步對值的訪問。 如果 lockFalse 則對所返回對象的訪問將不會自動得到鎖的保護(hù),也就是說它將不是“進(jìn)程安全的”。

注意 locl 只能是命名參數(shù)。

multiprocessing.sharedctypes.Value(typecode_or_type, *args, lock=True)?

返回一個純 ctypes 數(shù)組, 或者在此之上經(jīng)過同步器包裝過的進(jìn)程安全的對象,這取決于 lock 參數(shù)的值,除此之外,和 RawArray() 一樣。

如果 lockTrue (默認(rèn)值) 則將創(chuàng)建一個新的鎖對象用于同步對值的訪問。 如果 lock 為一個 LockRLock 對象則該對象將被用于同步對值的訪問。 如果 lockFalse 則對所返回對象的訪問將不會自動得到鎖的保護(hù),也就是說它將不是“進(jìn)程安全的”。

注意 locl 只能是命名參數(shù)。

multiprocessing.sharedctypes.copy(obj)?

從共享內(nèi)存中申請一片空間將 ctypes 對象 obj 過來,然后返回一個新的 ctypes 對象。

multiprocessing.sharedctypes.synchronized(obj[, lock])?

將一個 ctypes 對象包裝為進(jìn)程安全的對象并返回,使用 lock 同步對于它的操作。如果 lockNone (默認(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對象的語法。(表格中的 MyStructctypes.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)程可以通過代理訪問這些共享對象。

multiprocessing.Manager()?

返回一個已啟動的 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)聽的地址。如果 addressNone ,則允許和任意主機(jī)的請求建立連接。

authkey 是認(rèn)證標(biāo)識,用于檢查連接服務(wù)進(jìn)程的請求合法性。如果 authkeyNone, 則會使用 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()
shutdown()?

停止管理器的進(jìn)程。這個方法只能用于已經(jīng)使用 start() 啟動的服務(wù)進(jìn)程。

它可以被多次調(diào)用。

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。

proxytypeBaseProxy? 的子類,可以根據(jù) typeid 為共享對象創(chuàng)建一個代理,如果是 None , 則會自動創(chuàng)建一個代理類。

exposed 是一個函數(shù)名組成的序列,用來指明只有這些方法可以使用 BaseProxy._callmethod() 代理。(如果 exposedNone, 則會在 proxytype._exposed_ 存在的情況下轉(zhuǎn)而使用它) 當(dāng)暴露的方法列表沒有指定的時候,共享對象的所有 “公共方法” 都會被代理。(這里的“公共方法”是指所有擁有 __call__() 方法并且不是以 '_' 開頭的屬性)

method_to_typeid 是一個映射,用來指定那些應(yīng)該返回代理對象的暴露方法所返回的類型。(如果 method_to_typeidNone, 則 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.Lockthreading.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 屬性的對象并返回它的代理。

dict()?
dict(mapping)
dict(sequence)

創(chuàng)建一個共享的 dict 對象并返回它的代理。

list()?
list(sequence)

創(chuàng)建一個共享的 list 對象并返回它的代理。

在 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}

如果指涉對象包含了普通 listdict 對象,對這些內(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ù)目。如果 processesNone,則使用 os.cpu_count() 返回的值。

如果 initializer 不為 None,則每個工作進(jìn)程將會在啟動時調(diào)用 initializer(*initargs)。

maxtasksperchild 是一個工作進(jìn)程在它退出或被一個新的工作進(jìn)程代替之前能完成的任務(wù)數(shù)量,為了釋放未使用的資源。默認(rèn)的 maxtasksperchildNone,意味著工作進(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)程。 Poolmaxtasksperchild 參數(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í)行速度。

同樣,如果 chunksize1 , 那么 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 命名管道)。其中只有第一個保證各平臺可用。如果 familyNone ,那么 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 是一個浮點型,該方法會最多阻塞這么多秒。如果 timeoutNone ,則會允許阻塞的事件沒有限制。timeout為負(fù)數(shù)的情況下和為0的情況相同。

對于 Unix 和 Windows ,下列對象都可以出現(xiàn)在 object_list

當(dāng)一個連接或者套接字對象擁有有效的數(shù)據(jù)可被讀取的時候,或者另一端關(guān)閉后,這個對象就處于就緒狀態(tài)。

Unix: wait(object_list, timeout)select.select(object_list, [], [], timeout) 幾乎相同。差別在于,如果 select.select() 被信號中斷,它會拋出一個附帶錯誤號為 EINTROSError 異常,而 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)險的。所以 ListenerClient() 之間使用 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ù)目。 如果 processesNone,則使用 os.cpu_count() 返回的值。

如果 initializer 不為 None,則每個工作進(jìn)程將會在啟動時調(diào)用 initializer(*initargs)。

不同于 Poolmaxtasksperchildcontext 不可被提供。

注解

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

需要更多信息,請查看 bpo-5155, bpo-5313 以及 bpo-5331

spawnforkserver 啟動方式?

相對于 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)程)

例如,使用 spawnforkserver 啟動方式執(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()