Python之multiprocessing包使用及說明
簡介
multiprocessing 是 Python 標(biāo)準(zhǔn)庫中的一個(gè)包,它支持進(jìn)程間的并發(fā)執(zhí)行,提供了一個(gè)與 threading 模塊相似的API。
通過使用 multiprocessing,開發(fā)者可以在 Python 程序中創(chuàng)建新的進(jìn)程,每個(gè)進(jìn)程都擁有自己的 Python 解釋器和內(nèi)存空間,從而能夠并行執(zhí)行任務(wù)。
這在進(jìn)行CPU密集型任務(wù)時(shí)尤其有用,因?yàn)樗梢岳@過全局解釋器鎖(GIL),讓多核或多CPU的系統(tǒng)可以被充分利用.
主要功能
進(jìn)程管理
Process類:用于表示一個(gè)進(jìn)程對象。它允許你啟動(dòng)、停止、等待進(jìn)程結(jié)束等??梢酝ㄟ^傳遞目標(biāo)函數(shù)和參數(shù)來創(chuàng)建 Process 實(shí)例。
進(jìn)程間通信(IPC)
Queue:提供了一種安全的方式來在進(jìn)程之間進(jìn)行數(shù)據(jù)交換。隊(duì)列內(nèi)部使用鎖來保證數(shù)據(jù)的一致性。Pipe:提供了一種進(jìn)程間雙向通信的方法。管道可以是雙向的(雙工)或單向的(半雙工)。
同步原語
Lock:用于同步進(jìn)程間的操作,確保一次只有一個(gè)進(jìn)程可以訪問共享資源。Event:允許進(jìn)程等待某些事件的發(fā)生。進(jìn)程可以通過調(diào)用 set() 方法來觸發(fā)事件。Semaphore:用于限制對共享資源的訪問,它初始化一個(gè)計(jì)數(shù)器,該計(jì)數(shù)器表示可以同時(shí)訪問資源的進(jìn)程數(shù)量。Condition:允許一個(gè)或多個(gè)進(jìn)程等待某個(gè)條件的發(fā)生,同時(shí)提供了通知機(jī)制。
進(jìn)程池
Pool類:允許創(chuàng)建一組進(jìn)程池,可以將任務(wù)分配給池中的進(jìn)程執(zhí)行。支持同步和異步的任務(wù)執(zhí)行方式,如 map、apply、apply_async 和 map_async。
共享狀態(tài)
- 共享內(nèi)存:
Value和Array類允許在不同進(jìn)程之間共享數(shù)據(jù)。這些數(shù)據(jù)存儲(chǔ)在共享內(nèi)存中,可以由多個(gè)進(jìn)程并發(fā)訪問。 Manager:提供了一種更靈活的方式來共享數(shù)據(jù)。它可以創(chuàng)建一個(gè)服務(wù)器進(jìn)程,其他進(jìn)程通過代理的方式訪問服務(wù)器進(jìn)程中的對象,如列表、字典等。
其他組件
cpu_count:返回可用的 CPU 數(shù)量,這對于決定進(jìn)程池的大小很有用。current_process:返回當(dāng)前進(jìn)程的信息。active_children:返回當(dāng)前活躍子進(jìn)程的列表。
功能詳解
進(jìn)程管理
Process類
構(gòu)造方法
init(self, group=None, target=None, name=None, args=(), kwargs={}, *, daemon=None)
group:目前為保留參數(shù),用于未來擴(kuò)展,應(yīng)始終為 None。target:表示這個(gè)進(jìn)程實(shí)例所調(diào)用對象,你可以傳入一個(gè)可調(diào)用對象或函數(shù),這個(gè)對象會(huì)在新進(jìn)程啟動(dòng)時(shí)執(zhí)行。name:為進(jìn)程指定一個(gè)名稱,主要用于調(diào)試和跟蹤。args:是 target 調(diào)用時(shí)的位置參數(shù),傳遞給函數(shù)的參數(shù)元組。kwargs:是 target 調(diào)用時(shí)的關(guān)鍵字參數(shù),傳遞給函數(shù)的字典。daemon:如果設(shè)置為 True,則這個(gè)進(jìn)程將被標(biāo)記為守護(hù)進(jìn)程,主程序結(jié)束時(shí)會(huì)自動(dòng)殺死守護(hù)進(jìn)程。默認(rèn)值 None 表示繼承父進(jìn)程的守護(hù)標(biāo)志。
實(shí)例方法
start()
- 啟動(dòng)進(jìn)程。
- 該方法將為目標(biāo)函數(shù)創(chuàng)建一個(gè)新的進(jìn)程,并在新進(jìn)程中調(diào)用該函數(shù)。
join(timeout=None)
- 阻塞調(diào)用線程,直到進(jìn)程的執(zhí)行結(jié)束或達(dá)到指定的超時(shí)時(shí)間。
- 如果未指定 timeout,則會(huì)一直等待直到進(jìn)程結(jié)束。
is_alive()
- 返回進(jìn)程是否還活著。
- 一個(gè)活著的進(jìn)程是啟動(dòng)而且尚未終止的進(jìn)程。
terminate()
- 立即終止進(jìn)程。
- 注意,這種方法終止進(jìn)程是不安全的,可能會(huì)導(dǎo)致進(jìn)程資源未能被正確清理(如臨時(shí)文件、鎖等)。
kill()
- Python 3.7 中新增。
- 作用同 terminate(),但使用更強(qiáng)的信號(hào) SIGKILL,確保進(jìn)程被終止。
close()
- Python 3.7 中新增。關(guān)閉 Process 對象,釋放所有相關(guān)資源。
- 如果進(jìn)程尚未結(jié)束,則首先需要調(diào)用 join()
屬性
exitcode
- 進(jìn)程的退出代碼。
- 進(jìn)程還在運(yùn)行時(shí)為 None,正常退出時(shí)為 0,如果進(jìn)程因信號(hào)而終止,則為負(fù)信號(hào)號(hào)碼。
name
- 進(jìn)程的名稱。
daemon
- 獲取或設(shè)置進(jìn)程的守護(hù)標(biāo)志。
- 如果這個(gè)進(jìn)程是守護(hù)進(jìn)程,則返回 True。
pid
- 進(jìn)程的 PID。
- 如果進(jìn)程還未啟動(dòng),則為 None。
示例
from multiprocessing import Process
import time
def worker(name, sleep_time):
print(f"Started worker {name}")
time.sleep(sleep_time)
print(f"Sleep time :{sleep_time}")
if __name__ == "__main__":
for i in range(3):
process = Process(target=worker, args=(f'Bob-{i}',i), name=f'Worker-{i}')
process.start() # 啟動(dòng)進(jìn)程
print(f"進(jìn)程是否存活:{process.is_alive()}")
print(time.time())
>>>
進(jìn)程是否存活:True
1734319313.823571
進(jìn)程是否存活:True
1734319313.824584
進(jìn)程是否存活:True
1734319313.825512
Started worker Bob-1
Started worker Bob-2
Started worker Bob-0
Sleep time :0
Sleep time :1
Sleep time :2
- 這個(gè)例子展示了如何創(chuàng)建和啟動(dòng)一個(gè)進(jìn)程,以及如何等待進(jìn)程完成。使用
Process類的方法和屬性可以讓你更精確地控制進(jìn)程的生命周期和行為 - 可以觀察到程序并行運(yùn)行,而如果使用
join()方法,則各個(gè)進(jìn)程則會(huì)被阻塞
進(jìn)程間通信(IPC)
Queue類
Queue 類是 multiprocessing 模塊中提供的一個(gè)重要組件,用于實(shí)現(xiàn)進(jìn)程間通信(IPC)。通過隊(duì)列,不同的進(jìn)程可以安全地交換信息或數(shù)據(jù)。
構(gòu)造方法
Queue(maxsize=0)
- 初始化一個(gè)隊(duì)列對象。maxsize 參數(shù)指定隊(duì)列中允許的最大項(xiàng)數(shù)。
- 如果 maxsize 小于或等于零,則隊(duì)列大小為無限。
實(shí)例方法
put(item, block=True, timeout=None)
- 將 item 放入隊(duì)列中。
- 如果 block 為 True(默認(rèn)值),且 timeout 為 None,則在必要時(shí)阻塞,直到有空間可用。如果 timeout 是一個(gè)正數(shù),它將最多阻塞 timeout 秒,然后拋出 queue.Full 異常(如果在這段時(shí)間內(nèi)沒有空間可用)。
- 如果 block 為 False,但隊(duì)列已滿,則立即拋出 queue.Full 異常。
get(block=True, timeout=None)
- 從隊(duì)列中移除并返回一個(gè)項(xiàng)目。
- 如果 block 為 True(默認(rèn)值),且 timeout 為 None,則在必要時(shí)阻塞,直到有項(xiàng)目可用。如果 timeout 是一個(gè)正數(shù),它將最多阻塞 timeout 秒,然后拋出 queue.Empty 異常(如果在這段時(shí)間內(nèi)沒有項(xiàng)目可用)。
- 如果 block 為 False,但隊(duì)列為空,則立即拋出 queue.Empty 異常。
empty()
- 返回 True 如果隊(duì)列為空;否則,返回 False。
- 注意,這個(gè)結(jié)果只是一個(gè)瞬間的快照,不保證后續(xù)的 put() 或 get() 調(diào)用會(huì)成功。
full()
- 返回 True 如果隊(duì)列已滿;否則,返回 False。
- 與 empty() 方法類似,這個(gè)結(jié)果也只是一個(gè)瞬間的快照。
qsize()
- 返回隊(duì)列中大約的項(xiàng)目數(shù)。注意,“大約”,是因?yàn)槿绻瑫r(shí)有其他線程或進(jìn)程在修改隊(duì)列,那么返回的結(jié)果可能不準(zhǔn)確。
close()
- 關(guān)閉隊(duì)列。這個(gè)方法在 Python 3.7 中被引入。
join_thread()
- 阻塞,直到處理隊(duì)列的背景線程退出。
- 這通常意味著隊(duì)列已被徹底耗盡,且不再有任何未處理的項(xiàng)目。
cancel_join_thread()
- 阻止隊(duì)列的加入線程。
- 這通常用于在進(jìn)程間隊(duì)列的生產(chǎn)者進(jìn)程已經(jīng)知道消費(fèi)者進(jìn)程已經(jīng)終止時(shí),避免無限阻塞。
示例
from multiprocessing import Process, Queue
def worker(q, msg):
q.put(msg)
if __name__ == "__main__":
msg = "Hello from parent process"
q = Queue()
p_1 = Process(target=worker, args=(q, f"{msg}-1"))
p_2 = Process(target=worker, args=(q, f"{msg}-2"))
p_1.start()
p_2.start()
p_1.join()
p_2.join()
while not q.empty():
print(q.get())
>>>
Hello from parent process-1
Hello from parent process-2
注意事項(xiàng)
- 在使用 Queue 時(shí),要注意在所有項(xiàng)目都被處理并且不再需要隊(duì)列時(shí),關(guān)閉和/或加入隊(duì)列的線程,以避免資源泄漏。
- 在多進(jìn)程環(huán)境下,對于隊(duì)列的操作(put、get)可能需要考慮適當(dāng)?shù)某瑫r(shí)時(shí)間,以避免死鎖或阻塞過長時(shí)間。
Pipe函數(shù)
Pipe 函數(shù)是 multiprocessing 模塊中用于進(jìn)程間通信的另一種主要機(jī)制。它創(chuàng)建了一對連接對象,默認(rèn)情況下是雙向的(雙工),但也可以配置為單向(半雙工)。通過管道,進(jìn)程可以相互發(fā)送和接收 Python 對象。Pipe 通常用于兩個(gè)進(jìn)程間的通信。
Pipe() 函數(shù)
multiprocessing.Pipe(duplex=True)
- duplex:如果 duplex 是 True(默認(rèn)值),則創(chuàng)建的管道是雙向的。如果是 False,則創(chuàng)建的管道是單向的,即只能用于發(fā)送或接收。
Pipe() 函數(shù)返回一對 Connection 對象,代表管道的兩端。
Connection 對象的方法
管道兩端的 Connection 對象提供了一組用于通信的方法:
send(obj)
- 發(fā)送一個(gè)對象到管道的另一端。
- obj 參數(shù)可以是任何可pickle的對象。
recv()
- 從管道接收另一端發(fā)送過來的對象。
- 此方法會(huì)阻塞,直到管道的另一端有數(shù)據(jù)發(fā)送過來。
fileno()
- 返回一個(gè)代表連接的文件描述符的整數(shù)。
- 這個(gè)文件描述符可以用于底層I/O操作,例如將其整合到事件循環(huán)中。
close()
- 關(guān)閉連接。
- 關(guān)閉后,連接不再可用。
poll(timeout=None)
- 檢查管道的另一端是否發(fā)送了數(shù)據(jù)。如果 timeout 是 None(默認(rèn)值),則立即返回結(jié)果。如果 timeout 是一個(gè)數(shù)字,則最多等待 timeout 秒。
- 返回 True 表示有數(shù)據(jù)可讀取,F(xiàn)alse 表示沒有。
示例
from multiprocessing import Process, Pipe
def worker(conn):
conn.send([42, None, 'hello from child'])
conn.close()
if __name__ == '__main__':
parent_conn, child_conn = Pipe()
p = Process(target=worker, args=(child_conn,))
p.start()
print(parent_conn.recv()) # 輸出: [42, None, 'hello from child']
parent_conn.close()
p.join()
>>>
[42, None, 'hello from child']
注意事項(xiàng)
- 使用管道時(shí),務(wù)必注意關(guān)閉不再使用的連接端,以避免資源泄露或不必要的阻塞。
- 管道適用于少量進(jìn)程間的通信。對于涉及多個(gè)進(jìn)程的復(fù)雜場景,可能需要考慮其他通信機(jī)制,如
Queue。 - 在雙向管道中,如果兩個(gè)進(jìn)程同時(shí)嘗試發(fā)送數(shù)據(jù),而沒有對應(yīng)的接收操作,可能會(huì)導(dǎo)致死鎖。設(shè)計(jì)通信協(xié)議時(shí),應(yīng)確保發(fā)送和接收操作的邏輯順序能夠預(yù)防死鎖。
同步原語
Lock類
Lock 類是 multiprocessing 模塊中提供的一種簡單的同步原語。
它用于控制多個(gè)進(jìn)程對共享資源的訪問,以防止資源的并發(fā)訪問導(dǎo)致數(shù)據(jù)錯(cuò)亂或損壞。Lock 類似于線程模塊中的鎖,但它是為進(jìn)程間同步而設(shè)計(jì)的。
構(gòu)造方法
Lock 類沒有特定的構(gòu)造參數(shù),你可以直接創(chuàng)建一個(gè) Lock 實(shí)例
from multiprocessing import Lock lock = Lock()
實(shí)例方法
acquire(block=True, timeout=-1)
- 嘗試獲取鎖。
- 如果 block 為 True(默認(rèn)值),則調(diào)用將阻塞直到鎖被釋放,然后獲取鎖并返回 True。如果 block 為 False,則調(diào)用將不阻塞,如果無法立即- 獲取鎖,則返回 False。
- timeout 參數(shù)(僅當(dāng) block=True 時(shí)有效)允許指定阻塞的最大時(shí)間(秒)。timeout=-1(默認(rèn)值)表示無限等待。如果指定了 timeout 并在超時(shí)期間鎖未被釋放,則返回 False。
release()
釋放鎖。
只有在當(dāng)前進(jìn)程持有鎖時(shí)才能調(diào)用此方法;否則,將拋出 RuntimeError。釋放鎖后,其他阻塞在 acquire() 調(diào)用中等待獲取鎖的進(jìn)程將能夠獲取鎖。
locked()
- 如果鎖當(dāng)前被持有,則返回 True;否則返回 False。
- 注意,這個(gè)方法的結(jié)果僅僅是一個(gè)瞬時(shí)狀態(tài)的反映,調(diào)用它的時(shí)候鎖的狀態(tài)可能已經(jīng)改變。
示例
from multiprocessing import Process, Lock
import time
# 定義一個(gè)簡單的共享資源訪問函數(shù)
def worker_with(lock, num):
with lock:
print(f"Worker {num} has acquired the lock")
time.sleep(1)
print(f"Worker {num} is releasing the lock")
if __name__ == "__main__":
lock = Lock()
processes = [Process(target=worker_with, args=(lock, n)) for n in range(5)]
for p in processes:
p.start()
for p in processes:
p.join()
print("Processing complete.")
>>>
確保了在任何時(shí)刻只有一個(gè)進(jìn)程能夠訪問共享資源,從而避免了并發(fā)訪問的問題
注意事項(xiàng)
- 使用鎖時(shí)應(yīng)該小心,避免出現(xiàn)死鎖情況。確保每次成功獲取鎖之后,都能夠正確釋放鎖。
- 過度使用鎖可能會(huì)導(dǎo)致程序性能下降,因?yàn)樗拗屏舜a并發(fā)執(zhí)行的能力。在設(shè)計(jì)程序時(shí),應(yīng)盡量減少鎖的使用,只在絕對必要時(shí)使用鎖來保護(hù)共享資源。
- 與所有同步原語一樣,
Lock應(yīng)該謹(jǐn)慎使用,以避免復(fù)雜的同步問題。
Event類
Event 類是 Python multiprocessing 模塊中提供的一個(gè)同步原語,用于在進(jìn)程之間通信和協(xié)調(diào)。它模擬了線程模塊中的 Event 對象,允許一個(gè)進(jìn)程向其他多個(gè)進(jìn)程發(fā)出某個(gè)事件的發(fā)生,這對于在多個(gè)進(jìn)程之間同步操作非常有用
構(gòu)造方法
創(chuàng)建一個(gè) Event 對象很簡單,不需要任何參數(shù)
from multiprocessing import Event event = Event()
實(shí)例方法
set()
- 將內(nèi)部標(biāo)志設(shè)置為 True。
- 所有正在等待這個(gè)事件的進(jìn)程將被喚醒。
clear()
- 將內(nèi)部標(biāo)志設(shè)置為 False。
- 之后,調(diào)用 wait() 的進(jìn)程將會(huì)被阻塞,直到 set() 被調(diào)用。
is_set()
- 返回 Event 的內(nèi)部標(biāo)志的當(dāng)前狀態(tài)。
- 如果標(biāo)志為 True,返回 True;否則返回 False。
wait(timeout=None)
- 阻塞調(diào)用進(jìn)程,直到內(nèi)部標(biāo)志為 True。如果 timeout 為 None(默認(rèn)值),則無限期等待。timeout 是一個(gè)數(shù)值,表示等待的最大秒數(shù)。
- 如果在 wait() 調(diào)用期間內(nèi)部標(biāo)志被設(shè)置為 True,或者在調(diào)用 wait() 時(shí)內(nèi)部標(biāo)志已經(jīng)是 True,則返回 True;如果是因?yàn)槌瑫r(shí)而返回,則返回 False。
示例
from multiprocessing import Process, Event
import time
def worker(event, id):
print(f"Worker {id} waiting for event.")
event.wait() # 阻塞,直到事件被設(shè)置。
print(f"Worker {id} received event.")
if __name__ == "__main__":
event = Event()
# 創(chuàng)建并啟動(dòng)三個(gè)工作進(jìn)程
workers = [Process(target=worker, args=(event, i)) for i in range(3)]
for w in workers:
w.start()
print("Main process doing some work.")
time.sleep(2) # 模擬主進(jìn)程工作一段時(shí)間
event.set() # 設(shè)置事件,喚醒所有等待的進(jìn)程
print("Main process triggered the event.")
for w in workers:
w.join() # 等待所有工作進(jìn)程完成
print("Main process exiting.")
>>>
Main process doing some work.
Worker 0 waiting for event.
Worker 1 waiting for event.
Worker 2 waiting for event.
Main process triggered the event.
Worker 0 received event.
Worker 1 received event.
Worker 2 received event.
Main process exiting.
注意事項(xiàng)
- 使用
Event可以實(shí)現(xiàn)進(jìn)程間的簡單通信,但它不提供保護(hù)共享資源的機(jī)制。如果需要對共享資源進(jìn)行訪問控制,應(yīng)考慮使用鎖(如Lock、RLock)。 Event對象適用于當(dāng)一個(gè)進(jìn)程必須等待一個(gè)或多個(gè)進(jìn)程發(fā)出特定信號(hào)才能繼續(xù)執(zhí)行的場景。- 在使用
Event進(jìn)行同步時(shí),要注意可能的死鎖情況,特別是在復(fù)雜的進(jìn)程間通信和協(xié)調(diào)場景下。
Semaphore類
Semaphore 類是 multiprocessing 模塊中提供的一種同步原語,用于控制對共享資源的訪問數(shù)量。它可以被視為一個(gè)可用資源的計(jì)數(shù)器,是一種更為通用的同步機(jī)制,可以用來解決多個(gè)進(jìn)程訪問有限數(shù)量的資源問題。
構(gòu)造方法
Semaphore 對象在創(chuàng)建時(shí)可以接受一個(gè)可選的整數(shù)值,用于指定信號(hào)量的初始值,即同時(shí)可以訪問共享資源的進(jìn)程數(shù)量。如果不指定,則默認(rèn)值為1,此時(shí)它的行為類似于 Lock。
from multiprocessing import Semaphore sem = Semaphore(value=3)
實(shí)例方法
acquire(block=True, timeout=None)
嘗試減少信號(hào)量計(jì)數(shù)器的值(即獲取資源)。如果計(jì)數(shù)器的值大于0,則減1并立即返回 True。如果計(jì)數(shù)器的值為0,則根據(jù) block 參數(shù)的值行為會(huì)有所不同:
- 如果 block 為 True(默認(rèn)值),調(diào)用將阻塞,直到其他進(jìn)程釋放資源(即調(diào)用 release()),計(jì)數(shù)器值變?yōu)榇笥?。
- 如果 block 為 False,則調(diào)用不會(huì)阻塞,立即返回 False。
- timeout 參數(shù)(僅當(dāng) block=True 時(shí)有效)允許指定阻塞的最長時(shí)間(秒)。如果在超時(shí)期間資源未被釋放,則返回 False。
release()
- 釋放資源,即將信號(hào)量計(jì)數(shù)器的值增加1。
- 當(dāng)計(jì)數(shù)器的值從0變?yōu)?時(shí),正在等待這個(gè)信號(hào)量的其他進(jìn)程將能夠繼續(xù)執(zhí)行。
get_value() (在某些Python版本中不推薦使用)
- 返回信號(hào)量的當(dāng)前值,即當(dāng)前可用資源的數(shù)量。
- 需要注意的是,因?yàn)檫M(jìn)程間的狀態(tài)可能會(huì)迅速變化,所以返回的值可能并不準(zhǔn)確,主要用于調(diào)試目的。
示例
from multiprocessing import Process, Semaphore
import time
sem = Semaphore(2) # 最多允許2個(gè)進(jìn)程同時(shí)訪問共享資源
def worker(sem, num):
with sem:
print(f"Worker {num} is working.")
time.sleep(num) # 模擬耗時(shí)操作
print(f"Worker {num} finished.")
if __name__ == "__main__":
workers = [Process(target=worker, args=(sem, i)) for i in range(5)]
for w in workers:
w.start()
for w in workers:
w.join()
print("All workers completed.")
>>>
Worker 0 is working.
Worker 0 finished.
Worker 4 is working.
Worker 3 is working.
Worker 3 finished.
Worker 2 is working.
Worker 4 finished.
Worker 1 is working.
Worker 1 finished.
Worker 2 finished.
All workers completed.
注意事項(xiàng)
- 在使用
Semaphore時(shí)應(yīng)當(dāng)注意,acquire和release方法調(diào)用需要成對出現(xiàn),以避免信號(hào)量計(jì)數(shù)器的值被錯(cuò)誤地修改,導(dǎo)致資源訪問控制混亂。 Semaphore可用于實(shí)現(xiàn)各種同步模式,包括限制對資源的并發(fā)訪問、實(shí)現(xiàn)生產(chǎn)者消費(fèi)者問題等。- 在設(shè)計(jì)多進(jìn)程程序時(shí),合理使用
Semaphore可以有效地控制資源訪問,防止資源競爭和沖突,但也需要注意避免死鎖和其他同步問題。
Condition類
Condition 類是 multiprocessing 模塊中提供的一個(gè)同步原語,用于在進(jìn)程之間等待某些條件的滿足。它允許一個(gè)或多個(gè)進(jìn)程等待某個(gè)條件變?yōu)檎妫跅l件滿足時(shí)能夠通知一個(gè)或多個(gè)等待的進(jìn)程繼續(xù)執(zhí)行。Condition 常常與共享資源或狀態(tài)變更相關(guān)的場景一起使用,提供了一種更加靈活的進(jìn)程同步機(jī)制
構(gòu)造方法
在創(chuàng)建 Condition 對象時(shí),可以選擇傳遞一個(gè) Lock 或 RLock 對象用于內(nèi)部使用。如果不傳遞,則 Condition 對象會(huì)自動(dòng)創(chuàng)建一個(gè)新的 RLock 對象
from multiprocessing import Condition cond = Condition()
實(shí)例方法
acquire(*args, **kwargs)
- 獲得與條件關(guān)聯(lián)的鎖。這個(gè)方法直接調(diào)用內(nèi)部鎖對象的 acquire 方法,并且接受相同的參數(shù)。
- 通常在調(diào)用 wait()、notify() 或 notify_all() 之前需要先調(diào)用 acquire() 獲得鎖。
release()
- 釋放與條件關(guān)聯(lián)的鎖。這個(gè)方法直接調(diào)用內(nèi)部鎖對象的 release 方法。
- 一般在調(diào)用 wait() 后,當(dāng)條件滿足繼續(xù)執(zhí)行后,需要釋放鎖。
wait(timeout=None)
- 讓當(dāng)前進(jìn)程等待直到被通知或超時(shí)。調(diào)用這個(gè)方法時(shí),進(jìn)程必須已經(jīng)獲得了與條件關(guān)聯(lián)的鎖。
- 方法會(huì)自動(dòng)釋放鎖,然后進(jìn)程被阻塞,直到其他進(jìn)程在同一個(gè)條件變量上調(diào)用 notify() 或 notify_all()。一旦被通知,wait 方法重新獲得鎖然后返回。
- timeout 指定等待的最長時(shí)間,如果未指定,則無限等待。
notify(n=1)
- 通知等待這個(gè)條件的進(jìn)程,使其繼續(xù)執(zhí)行。
- 調(diào)用這個(gè)方法時(shí),進(jìn)程必須已經(jīng)獲得了與條件關(guān)聯(lián)的鎖。
- n 指定要喚醒的等待進(jìn)程的數(shù)量,默認(rèn)為1。
notify_all() 或 notifyAll()
- 通知等待這個(gè)條件的所有進(jìn)程,使其繼續(xù)執(zhí)行。
- 調(diào)用這個(gè)方法時(shí),進(jìn)程必須已經(jīng)獲得了與條件關(guān)聯(lián)的鎖。
示例
from multiprocessing import Process, Condition, Array
import time
def producer(cond, shared_array):
with cond:
print("Producer adding item.")
shared_array[0] += 1
print("Producer done, item added.")
cond.notify()
def consumer(cond, shared_array):
with cond:
print("Consumer waiting for item.")
cond.wait()
print("Consumer consumed item.")
shared_array[0] -= 1
if __name__ == "__main__":
condition = Condition()
shared_array = Array('i', [0])
p = Process(target=producer, args=(condition, shared_array))
c = Process(target=consumer, args=(condition, shared_array))
c.start()
time.sleep(2) # 確保消費(fèi)者先運(yùn)行
p.start()
p.join()
c.join()
print("Final shared value:", shared_array[0])
>>>
Consumer waiting for item.
Producer adding item.
Producer done, item added.
Consumer consumed item.
Final shared value: 0
注意事項(xiàng)
- 使用
Condition對象時(shí),必須確保在調(diào)用wait()、notify()和notify_all()方法前已經(jīng)獲得了與條件關(guān)聯(lián)的鎖。 wait()方法在返回前會(huì)重新獲得鎖,這意味著調(diào)用wait()后不需要再次顯式調(diào)用acquire()。notify()和notify_all()不會(huì)立即釋放鎖,而是在當(dāng)前進(jìn)程的鎖釋放操作(release())執(zhí)行后,被通知的進(jìn)程才能繼續(xù)執(zhí)行。- 正確使用
Condition可以解決復(fù)雜的同步問題,但設(shè)計(jì)不當(dāng)可能導(dǎo)致死鎖或性能問題,因此需要謹(jǐn)慎使用。
進(jìn)程池
Pool類
Pool 類是 Python 的 multiprocessing 模塊提供的一個(gè)高級(jí)接口,用于并行執(zhí)行多個(gè)進(jìn)程。Pool 類可以自動(dòng)管理進(jìn)程池中的進(jìn)程數(shù)量,讓你輕松地將任務(wù)分配給進(jìn)程池執(zhí)行,而無需手動(dòng)管理每個(gè)進(jìn)程的創(chuàng)建和終止。這對于執(zhí)行 CPU 密集型任務(wù)特別有用
構(gòu)造方法
創(chuàng)建 Pool 對象時(shí),可以指定幾個(gè)參數(shù),最常用的是 processes 參數(shù),它指定了進(jìn)程池中進(jìn)程的數(shù)量。如果不指定,默認(rèn)為機(jī)器的 CPU 核心數(shù)
from multiprocessing import Pool pool = Pool(processes=4)
實(shí)例方法
apply(func, args=(), kwds={})
- 使用阻塞模式執(zhí)行一個(gè)任務(wù)。func 是要執(zhí)行的函數(shù),args 是函數(shù)參數(shù)的元組,kwds 是函數(shù)參數(shù)的字典。
- 這個(gè)方法只有在前一個(gè)任務(wù)完成后才會(huì)返回,然后返回函數(shù)的結(jié)果。
apply_async(func, args=(), kwds={}, callback=None, error_callback=None)
- 使用非阻塞模式執(zhí)行一個(gè)任務(wù)。參數(shù)與 apply 相同,但是會(huì)立即返回一個(gè) AsyncResult 對象。
- 你可以使用這個(gè)對象查詢?nèi)蝿?wù)的狀態(tài)或獲取結(jié)果。callback 是一個(gè)可選的回調(diào)函數(shù),當(dāng)任務(wù)完成時(shí)調(diào)用。
- error_callback 是任務(wù)執(zhí)行過程中發(fā)生異常時(shí)調(diào)用的函數(shù)。
map(func, iterable, chunksize=None)
- 類似于內(nèi)置的 map() 函數(shù),它會(huì)將 iterable 中的每個(gè)元素應(yīng)用于 func 函數(shù)。這個(gè)方法會(huì)阻塞直到結(jié)果返回。
- 使用 chunksize 參數(shù)可以控制每個(gè)任務(wù)的大小,對于非常大的迭代對象,通過合適的 chunksize 可以提高性能。
map_async(func, iterable, chunksize=None, callback=None, error_callback=None)
- 非阻塞版本的 map()。它立即返回一個(gè) AsyncResult 對象,可以用來查詢?nèi)蝿?wù)狀態(tài)或獲取結(jié)果。
- callback 和 error_callback 參數(shù)的作用與 apply_async 相同。
close()
- 關(guān)閉進(jìn)程池,使其不再接受新的任務(wù)。
- 但是已經(jīng)提交的任務(wù)會(huì)繼續(xù)執(zhí)行直到完成。
join()
- 等待所有進(jìn)程池中的進(jìn)程執(zhí)行完畢。
- close() 或 terminate() 方法被調(diào)用后才能使用 join()
示例
from multiprocessing import Pool
import time
import os
def square(x):
time.sleep(1)
return x * x
if __name__ == '__main__':
print(time.time())
# 創(chuàng)建一個(gè)包含4個(gè)進(jìn)程的進(jìn)程池
with Pool(processes=os.cpu_count()) as p:
# 使用 map 函數(shù)分配任務(wù)
results = p.map(square, range(10))
print(results)
# 使用 apply_async 異步執(zhí)行任務(wù)
async_result = p.apply_async(square, (10,))
print(async_result.get()) # 使用 get 方法獲取異步任務(wù)的結(jié)果
print(time.time())
>>>
1734330507.251192
[0, 1, 4, 9, 16, 25, 36, 49, 64, 81]
100
1734330510.379658
注意事項(xiàng)
- 使用
Pool時(shí),確保只在if __name__ == '__main__': 塊中創(chuàng)建Pool對象和執(zhí)行任務(wù)。這可以防止在子進(jìn)程中無意中創(chuàng)建新的子進(jìn)rocesses,導(dǎo)致程序崩潰或行為異常。 Pool的map和apply方法會(huì)阻塞主進(jìn)程,直到所有任務(wù)完成。如果需要非阻塞操作,應(yīng)使用map_async和apply_async方法。- 在使用完
Pool后,應(yīng)調(diào)用close()方法關(guān)閉進(jìn)程池,然后調(diào)用join()方法等待所有進(jìn)程完成。這是良好的資源管理習(xí)慣,可以防止資源泄露。 terminate()方法會(huì)立即停止所有進(jìn)程,應(yīng)謹(jǐn)慎使用,以免丟失未完成的工作
共享狀態(tài)
Value和Array類
在 Python 的 multiprocessing 模塊中,Value 和 Array 是兩個(gè)用于進(jìn)程間共享數(shù)據(jù)的類。由于進(jìn)程間共享內(nèi)存并非像線程間共享全局變量那樣簡單,這兩個(gè)類提供了一種在不同進(jìn)程間安全共享數(shù)據(jù)的方法。
Value
Value 類用于在進(jìn)程之間共享一個(gè)存儲(chǔ)在共享內(nèi)存中的單一數(shù)據(jù)值。它適用于當(dāng)你只需要共享一個(gè)簡單的數(shù)據(jù)元素,如一個(gè)整數(shù)或浮點(diǎn)數(shù)時(shí)。
使用方法:創(chuàng)建 Value 對象時(shí),需要指定數(shù)據(jù)類型和初始值。數(shù)據(jù)類型是使用類似于 C 語言類型聲明的字符串來指定的,例如 ‘i’ 表示整數(shù),‘d’ 表示雙精度浮點(diǎn)數(shù)。
from multiprocessing import Value
# 創(chuàng)建一個(gè)共享的整數(shù)變量,初始值為 0
num = Value('i', 0)
訪問和修改值:可以通過 .value 屬性來訪問和修改存儲(chǔ)在 Value 中的數(shù)據(jù)。
# 讀取值 print(num.value) # 修改值 num.value = 10
Array
Array 類用于在進(jìn)程之間共享一個(gè)存儲(chǔ)在共享內(nèi)存中的數(shù)組。它適用于當(dāng)你需要共享一組數(shù)據(jù)(如一系列整數(shù)或浮點(diǎn)數(shù))時(shí)。
使用方法:創(chuàng)建 Array 對象時(shí),同樣需要指定數(shù)據(jù)類型和數(shù)組的大小。數(shù)據(jù)類型的指定方式與 Value 類似,數(shù)組的大小則通過一個(gè)整數(shù)來指定。
from multiprocessing import Array
# 創(chuàng)建一個(gè)共享的整數(shù)數(shù)組,包含 10 個(gè)元素,初始值都為 0
arr = Array('i', 10)
訪問和修改數(shù)組元素:Array 支持類似于列表的索引和切片操作來訪問和修改數(shù)組中的元素。
# 讀取第一個(gè)元素 print(arr[0]) # 修改第一個(gè)元素 arr[0] = 1 # 獲取數(shù)組的長度 print(len(arr)) # 使用切片 arr[1:3] = [2, 3]
Manager類
Manager 類在 Python multiprocessing 模塊中提供了一種方式,允許你在不同進(jìn)程間共享數(shù)據(jù)。與 Value 和 Array 直接在共享內(nèi)存中存儲(chǔ)數(shù)據(jù)不同,Manager 類創(chuàng)建的數(shù)據(jù)結(jié)構(gòu)是通過一個(gè)服務(wù)器進(jìn)程管理的。這意味著它允許更加靈活的數(shù)據(jù)共享方式,不僅限于簡單的數(shù)值或數(shù)組,也支持列表、字典、Namespace 等更復(fù)雜的數(shù)據(jù)結(jié)構(gòu)。
基本使用
使用 Manager 類時(shí),首先需要?jiǎng)?chuàng)建一個(gè) Manager 對象,然后使用這個(gè)對象來創(chuàng)建共享的數(shù)據(jù)結(jié)構(gòu)。
from multiprocessing import Manager
with Manager() as manager:
# 創(chuàng)建一個(gè)共享的列表
shared_list = manager.list([0, 1, 2])
# 創(chuàng)建一個(gè)共享的字典
shared_dict = manager.dict({'a': 1, 'b': 2})
創(chuàng)建的共享數(shù)據(jù)結(jié)構(gòu)可以像普通的數(shù)據(jù)結(jié)構(gòu)那樣被訪問和修改,但是它們實(shí)際上是在一個(gè)單獨(dú)的服務(wù)器進(jìn)程中維護(hù)的。這意味著你可以安全地在多個(gè)進(jìn)程間共享和修改這些數(shù)據(jù)結(jié)構(gòu),而無需擔(dān)心進(jìn)程安全問題。
Manager 支持的類型
Manager 類支持多種類型的共享數(shù)據(jù)結(jié)構(gòu),包括但不限于:
- list: 共享列表
- dict: 共享字典
- Value: 存儲(chǔ)單個(gè)值的共享對象,類似于 multiprocessing.Value
- Array: 存儲(chǔ)多個(gè)值的共享數(shù)組,類似于 multiprocessing.Array
- Namespace: 允許訪問屬性的簡單方式,類似于一個(gè)簡單的類或結(jié)構(gòu)體
示例代碼
以下示例展示了如何使用 Manager 類在進(jìn)程間共享列表和字典:
from multiprocessing import Process, Manager
def worker(shared_list, shared_dict):
shared_list.append('new item')
shared_dict['new_key'] = 'new value'
if __name__ == "__main__":
with Manager() as manager:
shared_list = manager.list([1, 2, 3])
shared_dict = manager.dict({'key1': 'value1', 'key2': 'value2'})
p = Process(target=worker, args=(shared_list, shared_dict))
p.start()
p.join()
print(f"Shared list: {shared_list}")
print(f"Shared dict: {shared_dict}")
其他組件
multiprocessing 模塊提供了一些用于查詢系統(tǒng)和進(jìn)程狀態(tài)的函數(shù),這些函數(shù)對于編寫并發(fā)程序時(shí)了解資源利用情況和進(jìn)程管理非常有用。以下是對 cpu_count、current_process 和 active_children 三個(gè)方法的詳細(xì)解釋:
cpu_count()
cpu_count 函數(shù)返回機(jī)器上可用的 CPU 核心數(shù)。這個(gè)信息對于決定并行程序中進(jìn)程池(Pool)的大小非常有用,因?yàn)槟憧赡芟胍鶕?jù)可用的處理器數(shù)量來優(yōu)化你的程序性能。
使用場景:在創(chuàng)建進(jìn)程池時(shí),可以根據(jù) cpu_count 的返回值來設(shè)置進(jìn)程池的大小。例如,如果你的機(jī)器有 4 個(gè) CPU 核心,那么創(chuàng)建一個(gè)包含 4 個(gè)進(jìn)程的進(jìn)程池可能是一個(gè)合理的選擇。
from multiprocessing import cpu_count
print(f"Number of CPUs: {cpu_count()}")
current_process()
current_process 函數(shù)返回當(dāng)前進(jìn)程的信息。這個(gè)函數(shù)返回一個(gè) Process 對象,其中包含有關(guān)當(dāng)前進(jìn)程的詳細(xì)信息,如進(jìn)程的名稱、PID(進(jìn)程ID)等。
使用場景:這個(gè)函數(shù)在調(diào)試并發(fā)程序時(shí)特別有用,因?yàn)槟憧梢杂盟鼇碜R(shí)別當(dāng)前正在執(zhí)行的進(jìn)程。例如,在多進(jìn)程環(huán)境中打印日志時(shí),可能會(huì)包含進(jìn)程的名稱或 PID 來區(qū)分日志消息是從哪個(gè)進(jìn)程生成的。
from multiprocessing import current_process
process = current_process()
print(f"Current process name: {process.name}")
print(f"Current process ID: {process.pid}")
active_children()
active_children 函數(shù)返回一個(gè)列表,包含當(dāng)前活躍的子進(jìn)程對象。每個(gè)子進(jìn)程對象都包含有關(guān)該進(jìn)程的信息,如名稱、PID 等。這個(gè)函數(shù)可以在父進(jìn)程中調(diào)用,以獲取當(dāng)前所有活躍的子進(jìn)程列表。
使用場景:在管理和監(jiān)控基于 multiprocessing 的并發(fā)程序時(shí),active_children 函數(shù)可以幫助你了解有多少子進(jìn)程正在運(yùn)行。這對于確保所有預(yù)期內(nèi)的子進(jìn)程都已啟動(dòng),或者在程序結(jié)束時(shí)檢查是否還有未完成的子進(jìn)程非常有用。
from multiprocessing import Process, active_children
import time
def worker():
print("Worker sleeping...")
time.sleep(2)
print("Worker done.")
if __name__ == "__main__":
for _ in range(2):
p = Process(target=worker)
p.start()
time.sleep(1) # 等待一段時(shí)間,讓子進(jìn)程啟動(dòng)
for child in active_children():
print(f"Active child: PID={child.pid}, Name={child.name}")
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
Pandas之Dropna濾除缺失數(shù)據(jù)的實(shí)現(xiàn)方法
這篇文章主要介紹了Pandas之Dropna濾除缺失數(shù)據(jù)的實(shí)現(xiàn)方法,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2019-06-06
Python屬性(Property)優(yōu)雅掌控對象數(shù)據(jù)的完全指南
這篇文章主要為大家詳細(xì)介紹了Python屬性(Property)的用法和高級(jí)應(yīng)用場景,可以優(yōu)雅掌控你的對象數(shù)據(jù),文中的示例代碼講解詳細(xì),有需要的可以了解下2026-01-01
K-近鄰算法的python實(shí)現(xiàn)代碼分享
這篇文章主要介紹了K-近鄰算法的python實(shí)現(xiàn)代碼分享,具有一定借鑒價(jià)值,需要的朋友可以參考下。2017-12-12
人工智能學(xué)習(xí)pyTorch自建數(shù)據(jù)集及可視化結(jié)果實(shí)現(xiàn)過程
這篇文章主要為大家介紹了人工智能學(xué)習(xí)pyTorch自建數(shù)據(jù)集及可視化結(jié)果的實(shí)現(xiàn)過程,有需要的朋友可以借鑒參考下,希望能夠有所幫助2021-11-11
教你用Python腳本快速為iOS10生成圖標(biāo)和截屏
這篇文章主要介紹了教你用Python快速為iOS10生成圖標(biāo)和截屏的相關(guān)資料,非常不錯(cuò),具有參考借鑒價(jià)值,需要的朋友可以參考下2016-09-09
Python數(shù)據(jù)可視化之繪制柱狀圖和條形圖
今天帶大家學(xué)習(xí)怎么利用Python繪制柱狀圖,條形圖,文中有非常詳細(xì)的代碼示例,對正在學(xué)習(xí)python的小伙伴們很有幫助,需要的朋友可以參考下2021-05-05

