最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

Python之multiprocessing包使用及說明

 更新時(shí)間:2026年05月07日 10:15:59   作者:BBluster  
multiprocessing是 Python標(biāo)準(zhǔn)庫中的的一個(gè)包,支持進(jìn)程間的并發(fā)執(zhí)行,它提供了一個(gè)類似 threading 的 API,允許創(chuàng)建新的進(jìn)程,每個(gè)進(jìn)程有自己的 Python 解釋器器和內(nèi)存空間,通過使用multiprocessing,開發(fā)者可以利用多核或多 CPU 的系統(tǒng),進(jìn)行 CPU 密集型任務(wù)

簡介

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)存:ValueArray 類允許在不同進(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)注意,acquirerelease 方法調(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)致程序崩潰或行為異常。
  • Poolmapapply 方法會(huì)阻塞主進(jìn)程,直到所有任務(wù)完成。如果需要非阻塞操作,應(yīng)使用 map_asyncapply_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 模塊中,ValueArray 是兩個(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)文章

最新評論

北流市| 榆中县| 普洱| 永康市| 潮州市| 井研县| 临清市| 军事| 大关县| 和平区| 海林市| 东至县| 万州区| 平武县| 东台市| 徐闻县| 合水县| 西林县| 武山县| 宜兰县| 大田县| 澄迈县| 铁岭县| 潼关县| 南溪县| 易门县| 灵川县| 湄潭县| 错那县| 且末县| 上杭县| 太谷县| 靖远县| 漳浦县| 瓮安县| 玉田县| 内江市| 盐山县| 阳江市| 东莞市| 辛集市|