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

Python模塊multiprocessing & 實(shí)現(xiàn)多進(jìn)程并發(fā)方式

 更新時(shí)間:2026年05月07日 10:23:58   作者:一只勤勞的耗子  
文章介紹了Python的multiprocessing模塊,用于實(shí)現(xiàn)多進(jìn)程編程,涵蓋了Process、Queue、Pipe、Pool等主要類及其方法,以及進(jìn)程管理、進(jìn)程間通信、進(jìn)程同步等內(nèi)容,強(qiáng)調(diào)了多進(jìn)程的優(yōu)勢和應(yīng)用場景,并提供示例代碼幫助理解

簡介

multiprocessing模塊是Python標(biāo)準(zhǔn)庫中提供的一個(gè)用于實(shí)現(xiàn)多進(jìn)程編程的模塊。它基于進(jìn)程而不是線程,可以利用多核CPU的優(yōu)勢,提高程序的執(zhí)行效率,同時(shí)也可以實(shí)現(xiàn)進(jìn)程間通信和數(shù)據(jù)共享。

1. 參數(shù)說明

1.1.Process(控制進(jìn)程)

用于創(chuàng)建子進(jìn)程對象,必須指定要執(zhí)行的目標(biāo)函數(shù)。Process(target = [函數(shù)名])

語法

multiprocessing.Process(
    target   #必選參數(shù),表示進(jìn)程所執(zhí)行的目標(biāo)函數(shù)。
    group    #預(yù)留參數(shù),不需要使用。
    name     #指定子進(jìn)程的名稱,通過 [子進(jìn)程].name 獲取
    args     #指定子進(jìn)程(函數(shù))需要傳遞的參數(shù),以元組方式傳遞,例如:args=('A', 'B', 'C')
    kwargs   #指定子進(jìn)程(函數(shù))需要傳遞的參數(shù),以字典方式傳遞,例如:args={'name':'xiaowang', 'age':18}
    daemon   #將子進(jìn)程設(shè)置為守護(hù)進(jìn)程(True|False)。
)
  • 守護(hù)進(jìn)程:主進(jìn)程退出后,不論該子進(jìn)程是否結(jié)束,強(qiáng)制退出。
  • 非守護(hù)進(jìn)程:主進(jìn)程退出后,子進(jìn)程若未結(jié)束,會(huì)繼續(xù)執(zhí)行。

所有方法如下

p = multiprocessing.Process(target = xxx)

    p.run       #表示進(jìn)程在啟動(dòng)后要執(zhí)行的方法,需要應(yīng)用程序開發(fā)人員重寫。
    p.start     #啟動(dòng)進(jìn)程。
    p.join      #等待進(jìn)程終止。
    p.is_alive  #判斷進(jìn)程是否存活。
    p.ident     #獲取子進(jìn)程PID
    p.name      #獲取子進(jìn)程名稱
    p.kill      #強(qiáng)制結(jié)束子進(jìn)程。
    p.terminate #強(qiáng)制結(jié)束子進(jìn)程。
    p.close     #關(guān)閉與進(jìn)程關(guān)聯(lián)的所有管道和文件。
    authkey     #返回用于身份驗(yàn)證的進(jìn)程通信的密鑰。

1.2.Queue(進(jìn)程通信)

用于實(shí)現(xiàn)進(jìn)程間通信的隊(duì)列,支持多個(gè)進(jìn)程向同一個(gè)隊(duì)列中讀寫數(shù)據(jù)。

參數(shù)選項(xiàng)

multiprocessing.Queue([maxsize])
  • maxsize(可選)參數(shù)是一個(gè)int型整數(shù),用于指定隊(duì)列中最多可以存放的實(shí)例數(shù),超過此值會(huì)保持阻塞,直到隊(duì)列中有位置出現(xiàn)。
  • 如果maxsize是負(fù)數(shù),則表示隊(duì)列的長度為無限。如果maxsize是零,則表示隊(duì)列的長度為無限。
  • 但由于諸如pickle等協(xié)議的緣故,正在發(fā)送或接收的大型對象可能會(huì)導(dǎo)致程序死鎖,因此不建議這樣使用。

所有方法如下

'''向消息隊(duì)列發(fā)送信息'''
Queue.put(item[, block[, timeout]])
    #如果block參數(shù)是True,而且隊(duì)列已滿,那么程序就會(huì)在隊(duì)列有空間之前停滯等待。
    #如果block參數(shù)是False,并且隊(duì)列已滿,那么就會(huì)引發(fā)Full異常。
    #如果給出了可選參數(shù)timeout,它會(huì)阻塞timeout秒。
'''向消息隊(duì)列發(fā)送信息,如果隊(duì)列已滿,會(huì)引發(fā)Full異常,'''
Queue.put_nowait(item)
    #等同于Queue.put(item, False)。

'''獲取隊(duì)列中元素并刪除'''
Queue.get([block[, timeout]])
    #如果隊(duì)列為空,且block參數(shù)為True,那么程序會(huì)一直停滯等待,或者阻塞timeout秒。
    #如果block參數(shù)是False,而且隊(duì)列為空,那么就會(huì)引發(fā)Empty異常。
'''獲取隊(duì)列中元素并刪除,如果隊(duì)列為空則引發(fā)Empty異。'''
Queue.get_nowait()

'''關(guān)閉進(jìn)程間通信通道,以防止有些進(jìn)程被阻塞,從而導(dǎo)致程序死鎖或者內(nèi)存泄露'''
Queue.close()

'''判斷隊(duì)列是否為空'''
Queue.empty()    #為空返回True,不為空返回False。
'''判斷隊(duì)列是否已滿'''
Queue.full()     #隊(duì)列已滿返回True,未滿返回False。

'''獲取隊(duì)列中的元素個(gè)數(shù)'''
Queue.qsize()    #注意:此方法不可靠,因?yàn)樵?Queue.put() 和 Queue.get() 方法的過程中仍然可能發(fā)生改變。

1.3.Pipe(管道通信)

創(chuàng)建進(jìn)程間通信管道,支持兩個(gè)進(jìn)程之間的通信。

語法

multiprocessing.Pipe([duplex])
  • duplex:是否創(chuàng)建一個(gè)雙向通信的管道(默認(rèn)True),代表創(chuàng)建一個(gè)雙向管道。如果設(shè)置為False,則創(chuàng)建一個(gè)單向管道(只能從某一端寫入數(shù)據(jù),從另一端讀出數(shù)據(jù))。

所有方法如下

'''向管道中寫入消息'''
Pipe.send('消息')    #如果發(fā)送失敗會(huì)拋出 BrokenPipeError 異常。
'''從管道中讀取數(shù)據(jù),會(huì)阻塞直到有數(shù)據(jù)可讀'''
Pipe.recv():         #如果管道已關(guān)閉,則會(huì)返回一個(gè) EOFError 異常。

'''關(guān)閉管道'''
Pipe.close()

'''返回管道的文件描述符'''
Pipe.fileno()

'''判斷讀取管道是否阻塞(如果管道可讀或者關(guān)閉,返回True)'''
Pipe.poll([timeout]) #如果設(shè)置了 timeout,則會(huì)在指定時(shí)間內(nèi)返回。

'''二進(jìn)制數(shù)據(jù)-向管道中寫入消息'''
Pipe.send_bytes(buf[, offset[, size]]) #與send一樣功能,但是接受的參數(shù)為二進(jìn)制字符串。
'''二進(jìn)制數(shù)據(jù)-從管道中讀取數(shù)據(jù)'''
Pipe.recv_bytes([maxlength])           #與recv類似,但是返回的是二進(jìn)制字符串。

1.4.Pool(進(jìn)程池)

語法

multiprocessing.Pool(
    processes         #可選參數(shù),指定進(jìn)程池中的線程數(shù)(默認(rèn)CPU最大數(shù))
    initializer       #可選參數(shù),指定每個(gè)工作進(jìn)程啟動(dòng)時(shí)要調(diào)用的函數(shù)。默認(rèn)值為 None。
    initargs          #可選參數(shù),指定傳遞給初始化器函數(shù)的參數(shù)元組。默認(rèn)值為 ()。
    maxtasksperchild  #可選參數(shù),限制每個(gè)工作進(jìn)程可以執(zhí)行的任務(wù)數(shù)量,然后被終止,以避免內(nèi)存泄漏問題。默認(rèn)值為 None,表示進(jìn)程將一直存在。
)

所有方法如下

Pool.[方法]

'''進(jìn)程池中同步執(zhí)行函數(shù), 類似于 func(*args, **kwds) 表達(dá)式'''
Pool.apply(func, args=(), kwds={})    #如果進(jìn)程池中的一個(gè)進(jìn)程崩潰,那么就將函數(shù)重新提交到池中。

'''在進(jìn)程池中異步執(zhí)行函數(shù)'''
apply_async(func, args=(), callback=None)  #支持回調(diào)函數(shù)。

'''進(jìn)程池中同步執(zhí)行函數(shù),將執(zhí)行結(jié)果存儲(chǔ)于列表中返回, 類似于 map(func, iterable) 表達(dá)式'''
Pool.map(func, iterable, chunksize=None)  #返回一個(gè)包含結(jié)果的列表。如果進(jìn)程池中的一個(gè)進(jìn)程崩潰,那么就將函數(shù)重新提交到池中。

'''進(jìn)程池中異步執(zhí)行函數(shù),并返回生成器對象, 可以逐個(gè)獲取函數(shù)執(zhí)行結(jié)果'''
Pool.imap(func, iterable, chunksize=1)  #這可以在內(nèi)存有限的情況下處理大型數(shù)據(jù)集。如果進(jìn)程池中的一個(gè)進(jìn)程崩潰,那么就將函數(shù)重新提交到池中。

'''類似于 imap() 函數(shù),但是結(jié)果不保證按照迭代器的順序產(chǎn)生'''
Pool.imap_unordered(func, iterable, chunksize=1)  #如果進(jìn)程池中的一個(gè)進(jìn)程崩潰,那么就將函數(shù)重新提交到池中。

'''關(guān)閉進(jìn)程池'''
Pool.close()

'''等待所有進(jìn)程完成'''
Pool.join()

'''強(qiáng)制關(guān)閉所有進(jìn)程'''
Pool.terminate()

1.5.Lock(進(jìn)程鎖)

語法

multiprocessing.Lock(locking = True)
    locking=True:使用底層操作系統(tǒng)提供的鎖機(jī)制,例如 POSIX 信號量或 Windows 臨界區(qū)。這是默認(rèn)值。
    locking=False:會(huì)更快地獲取和釋放鎖,但是在一些平臺(tái)上可能會(huì)出現(xiàn)未定義的行為。

所有方法如下

'''獲取鎖'''
Lock.acquire(block=True, timeout=None)  #如果 block=True,則會(huì)阻塞直到獲取到鎖,否則立即返回。如果 timeout 不是 None,則在超時(shí)之前沒有獲取到鎖就會(huì)拋出 TimeoutError 異常。如果鎖已經(jīng)被當(dāng)前進(jìn)程持有,則會(huì)拋出 AssertionError 異常。
'''釋放鎖'''
Lock.release()  #如果鎖沒有被當(dāng)前進(jìn)程持有,則會(huì)拋出 AssertionError 異常。

Lock.locked()            #返回當(dāng)前鎖是否被持有的布爾值。
Lock.notify(n=1)         #喚醒 wait() 方法中等待鎖的 n 個(gè)線程。
Lock.notify_all()        #喚醒 wait() 方法中等待鎖的所有線程。
Lock.wait(timeout=None)  #等待鎖。如果鎖已經(jīng)被釋放,則會(huì)立即返回。如果鎖仍然被持有,則會(huì)釋放鎖,并阻塞直到其他線程喚醒該線程或者超時(shí),則會(huì)拋出 TimeoutError 異常。

2. 進(jìn)程管理

2.1. 判斷子進(jìn)程狀態(tài)

通過is_alive 判斷某個(gè)子進(jìn)程是否存活

import multiprocessing
from time import sleep

# 定義一個(gè)函數(shù)(作為子進(jìn)程)
def proc():
    print('=========== 子進(jìn)程開始運(yùn)行 ===========')
    sleep(3)
    print('=========== 子進(jìn)程運(yùn)行結(jié)束===========')

if __name__ == '__main__':
    # 將函數(shù)proc定義為子進(jìn)程
    p = multiprocessing.Process(target=proc)

    # 啟動(dòng)子進(jìn)程
    p.start()

    # 判斷子進(jìn)程是否結(jié)束
    while p.is_alive():
        print('[CHECK] 子進(jìn)程運(yùn)行中')
        sleep(1)
    else:
        print('子進(jìn)程已死亡!')

結(jié)果

[CHECK] 子進(jìn)程運(yùn)行中
=========== 子進(jìn)程開始運(yùn)行 ===========
[CHECK] 子進(jìn)程運(yùn)行中
[CHECK] 子進(jìn)程運(yùn)行中
[CHECK] 子進(jìn)程運(yùn)行中
=========== 子進(jìn)程運(yùn)行結(jié)束===========
子進(jìn)程已死亡!

2.2. 獲取子進(jìn)程信息

獲取子進(jìn)程名稱和PID

import multiprocessing

def proc1():
    '''模仿一個(gè)子進(jìn)程'''
    pass

def proc2():
    '''模仿一個(gè)子進(jìn)程'''
    pass

if __name__ == '__main__':
    # 創(chuàng)建2個(gè)子進(jìn)程
    p1 = multiprocessing.Process(target=proc1, name='【自定義的p1】')
    p2 = multiprocessing.Process(target=proc2)

    # 啟動(dòng)進(jìn)程1
    p1.start()
    # 啟動(dòng)進(jìn)程2
    p2.start()

    # 獲取子進(jìn)程的名稱和PID
    print(f'進(jìn)程1的名稱為:{p1.name}, PID為:{p1.ident}')
    print(f'進(jìn)程2的名稱為:{p2.name}, PID為:{p2.ident}')
  • 可以通過參數(shù) name 設(shè)置子進(jìn)程名稱,未設(shè)置默認(rèn)Process-[n]

結(jié)果

進(jìn)程1的名稱為:【自定義的p1】, PID為:20136
進(jìn)程2的名稱為:Process-2, PID為:15816

2.3. 子進(jìn)程執(zhí)行多任務(wù)

實(shí)現(xiàn)簡單的多任務(wù)執(zhí)行

import multiprocessing
import time

def progress1():
    '''模仿一個(gè)子進(jìn)程'''
    for i in range(1,6):
        print(f'[{i}/5] 這是進(jìn)程1,每隔1s輸出一次...')
        time.sleep(1)

def progress2():
    '''模仿一個(gè)子進(jìn)程'''
    for i in range(1,5):
        print(f'[{i}/4] 這是進(jìn)程2,每隔2s輸出一次...')
        time.sleep(2)

if __name__ == '__main__':
    # 創(chuàng)建2個(gè)子進(jìn)程
    p1 = multiprocessing.Process(target=progress1)
    p2 = multiprocessing.Process(target=progress2)

    # 啟動(dòng)進(jìn)程1
    p1.start()
    # 啟動(dòng)進(jìn)程2
    p2.start()

    # 程序執(zhí)行其他事情
    time.sleep(1)
    print('=============== 結(jié)束 =============== ')

輸出結(jié)果(主進(jìn)程并沒有去等待子進(jìn)程結(jié)束,直接做其他事)

[1/5] 這是進(jìn)程1,每隔1s輸出一次...
[1/4] 這是進(jìn)程2,每隔2s輸出一次...
=============== 結(jié)束 =============== 
[2/5] 這是進(jìn)程1,每隔1s輸出一次...
[3/5] 這是進(jìn)程1,每隔1s輸出一次...
[2/4] 這是進(jìn)程2,每隔2s輸出一次...
[4/5] 這是進(jìn)程1,每隔1s輸出一次...
[3/4] 這是進(jìn)程2,每隔2s輸出一次...
[5/5] 這是進(jìn)程1,每隔1s輸出一次...
[4/4] 這是進(jìn)程2,每隔2s輸出一次...

使用 join 等待某個(gè)子進(jìn)程結(jié)束再運(yùn)行主進(jìn)程

import multiprocessing
import time

def progress1():
    '''模仿一個(gè)子進(jìn)程'''
    for i in range(1,6):
        print(f'[{i}/5] 這是進(jìn)程1,每隔1s輸出一次...')
        time.sleep(1)

def progress2():
    '''模仿一個(gè)子進(jìn)程'''
    for i in range(1,5):
        print(f'[{i}/4] 這是進(jìn)程2,每隔2s輸出一次...')
        time.sleep(2)

if __name__ == '__main__':
    # 創(chuàng)建2個(gè)子進(jìn)程
    p1 = multiprocessing.Process(target=progress1)
    p2 = multiprocessing.Process(target=progress2)

    # 啟動(dòng)進(jìn)程1
    p1.start()
    # 啟動(dòng)進(jìn)程2
    p2.start()

    # 等待子進(jìn)程結(jié)束
    p1.join()
    p2.join()

    # 程序執(zhí)行其他事情
    print('=============== 結(jié)束 =============== ')

結(jié)果

[1/5] 這是進(jìn)程1,每隔1s輸出一次...
[1/4] 這是進(jìn)程2,每隔2s輸出一次...
[2/5] 這是進(jìn)程1,每隔1s輸出一次...
[3/5] 這是進(jìn)程1,每隔1s輸出一次...
[2/4] 這是進(jìn)程2,每隔2s輸出一次...
[4/5] 這是進(jìn)程1,每隔1s輸出一次...
[5/5] 這是進(jìn)程1,每隔1s輸出一次...
[3/4] 這是進(jìn)程2,每隔2s輸出一次...
[4/4] 這是進(jìn)程2,每隔2s輸出一次...
=============== 結(jié)束 =============== 

將 進(jìn)程1(p1) 設(shè)置為守護(hù)進(jìn)程(隨主進(jìn)程退出而退出)

import multiprocessing
import time

def progress1():
    '''模仿一個(gè)子進(jìn)程'''
    for i in range(1,6):
        print(f'[{i}/5] 這是進(jìn)程1,每隔1s輸出一次...')
        time.sleep(1)

def progress2():
    '''模仿一個(gè)子進(jìn)程'''
    for i in range(1,5):
        print(f'[{i}/4] 這是進(jìn)程2,每隔2s輸出一次...')
        time.sleep(2)

if __name__ == '__main__':
    # 創(chuàng)建2個(gè)子進(jìn)程
    p1 = multiprocessing.Process(target=progress1, name='【自定義的p1】', daemon=True)
    p2 = multiprocessing.Process(target=progress2)

    # 啟動(dòng)進(jìn)程1
    p1.start()
    # 啟動(dòng)進(jìn)程2
    p2.start()

    # 程序執(zhí)行其他事情
    time.sleep(2)
    print('====================== 結(jié)束 ======================')
    print(f'進(jìn)程1的名稱為:{p1.name}, PID為:{p1.ident}')
    print(f'進(jìn)程2的名稱為:{p2.name}, PID為:{p2.ident}')
    print('====================== 結(jié)束 ======================')

結(jié)果

[1/5] 這是進(jìn)程1,每隔1s輸出一次...
[1/4] 這是進(jìn)程2,每隔2s輸出一次...
[2/5] 這是進(jìn)程1,每隔1s輸出一次...
====================== 結(jié)束 ======================
進(jìn)程1的名稱為:【自定義的p1】, PID為:21252
進(jìn)程2的名稱為:Process-2, PID為:21253
====================== 結(jié)束 ======================
[2/4] 這是進(jìn)程2,每隔2s輸出一次...
[3/4] 這是進(jìn)程2,每隔2s輸出一次...
[4/4] 這是進(jìn)程2,每隔2s輸出一次...

2.4. 運(yùn)行多個(gè)并發(fā)

通過循環(huán)的方式去構(gòu)造多個(gè)并發(fā)

import multiprocessing
from time import sleep

class MyClass(object):
    def __init__(self, thread_num=1, sleep_proc=None):
        # thread_num表示進(jìn)程數(shù),sleep_proc表示是否等待退出
        self.thread_num = thread_num
        self.sleep_proc = sleep_proc

    def proc(self):
        # 定義一個(gè)并發(fā)的子進(jìn)程
        for i in range(1,4):
            print(f'[{i}/3] 我是一個(gè)子進(jìn)程')
            sleep(1)

    def call_proc(self):
        # 調(diào)用子進(jìn)程
        processes = []
        # 循環(huán)調(diào)用多個(gè)子進(jìn)程
        for num in range(self.thread_num):
            # 定義子進(jìn)程屬性
            p = multiprocessing.Process(target=MyClass().proc)
            # 啟動(dòng)子進(jìn)程
            p.start()
            # 將循環(huán)的子進(jìn)程放入列表,用于后面的等待退出
            processes.append(p)

        # 判斷是否等待子進(jìn)程結(jié)束
        if self.sleep_proc is True:
            for p in processes:
                p.join()

if __name__ == '__main__':
    # 指定2個(gè)并發(fā),并等待子進(jìn)程結(jié)束
    mc = MyClass(2,True)
    mc.call_proc()

    print('=============== 結(jié)束 ===============')

結(jié)果

[1/3] 我是一個(gè)子進(jìn)程
[1/3] 我是一個(gè)子進(jìn)程
[2/3] 我是一個(gè)子進(jìn)程
[2/3] 我是一個(gè)子進(jìn)程
[3/3] 我是一個(gè)子進(jìn)程
[3/3] 我是一個(gè)子進(jìn)程
=============== 結(jié)束 ===============

2.5. 強(qiáng)制殺死子進(jìn)程

通過 terminate 或者 kill 強(qiáng)制殺死子進(jìn)程

import multiprocessing
from time import sleep

# 定義2個(gè)函數(shù)(作為子進(jìn)程)
def proc1():
    print('=========== 子進(jìn)程1開始運(yùn)行 ===========')
    sleep(10)
    print('=========== 子進(jìn)程1運(yùn)行結(jié)束===========')

def proc2():
    print('=========== 子進(jìn)程2開始運(yùn)行 ===========')
    sleep(10)
    print('=========== 子進(jìn)程2運(yùn)行結(jié)束===========')

if __name__ == '__main__':
    # 將函數(shù)定義為子進(jìn)程
    p1 = multiprocessing.Process(target=proc1)
    p2 = multiprocessing.Process(target=proc2)

    # 啟動(dòng)子進(jìn)程
    p1.start()
    p2.start()

    # 3秒后手動(dòng)殺死子進(jìn)程
    sleep(3)
    p1.kill()       #使用kill殺死子進(jìn)程
    p2.terminate()  #使用terminate殺死子進(jìn)程
    print('子進(jìn)程p1、p2已被強(qiáng)制殺死!')

3. 進(jìn)程間的通信

進(jìn)程間通信(IPC,Interprocess communication)是一組編程接口,讓程序員能夠協(xié)調(diào)不同的進(jìn)程,使之能在一個(gè)操作系統(tǒng)里同時(shí)運(yùn)行,并相互傳遞、交換信息。這使得一個(gè)程序能夠在同一時(shí)間里處理許多用戶的要求。因?yàn)榧词怪挥幸粋€(gè)用戶發(fā)出要求,也可能導(dǎo)致一個(gè)操作系統(tǒng)中多個(gè)進(jìn)程的運(yùn)行,進(jìn)程之間必須互相通話。IPC接口就提供了這種可能性。

python 的 multiprocessing 模塊中,可以使用 Queue、Pipe、Manager 等數(shù)據(jù)結(jié)構(gòu)實(shí)現(xiàn)進(jìn)程間通信(IPC),也就是進(jìn)程之間交換數(shù)據(jù)。進(jìn)程間通信是多進(jìn)程編程中非常重要和常見的一部分,通常用于在多個(gè)進(jìn)程之間共享并傳遞信息、數(shù)據(jù)或任務(wù)結(jié)果。

3.1. Queue 進(jìn)程通信

from multiprocessing import Process, Queue
from time import sleep


def proc1(q):
    '''這是一個(gè)發(fā)送消息的函數(shù)'''
    msgs = ["香蕉", "蘋果", "水蜜桃"]
    for msg in msgs:
        # 將迭代對象放入通信隊(duì)列
        q.put(msg)
        # 打印當(dāng)前迭代的內(nèi)容
        print(f"[進(jìn)程1] 發(fā)送信息({msg})")
        sleep(1)
    # 迭代完成后關(guān)閉通信
    q.close()


def proc2(q):
    '''這是一個(gè)接收消息的函數(shù)'''
    while True:
        try:
            # 刪除通信隊(duì)列中的一個(gè)元素
            msg = q.get(block=False)
            print(f"[進(jìn)程2] 收到信息({msg}), 并刪除該信息")
        except:
            q.close()   # 通信完成后關(guān)閉
            break
        sleep(1)


if __name__ == '__main__':
    # 使用Queue方法通信
    q = Queue()
    # 定義子進(jìn)程屬性,將 Queue 方法傳入子進(jìn)程
    p1 = Process(target=proc1, args=(q,))
    p2 = Process(target=proc2, args=(q,))

    # 啟動(dòng)2個(gè)子進(jìn)程
    p1.start()
    p2.start()

    # 等待2個(gè)子進(jìn)程結(jié)束
    p1.join()
    p2.join()

    print("=============== 結(jié)束 ===============")
  • 進(jìn)程1發(fā)送信息,進(jìn)程2讀取信息并刪除隊(duì)列(生產(chǎn)者-消費(fèi)者)

邏輯視圖

3.2.Queue控制多個(gè)子進(jìn)程

  • 設(shè)置多個(gè)子進(jìn)程,一個(gè)用于控制,其他執(zhí)行對應(yīng)程序
  • 控制的方法:運(yùn)行子進(jìn)程通過消息隊(duì)列判斷是否繼續(xù)運(yùn)行。

若消息隊(duì)列為空,繼續(xù)運(yùn)行子進(jìn)程;若消息隊(duì)列不為空,停止子進(jìn)程。由控制子進(jìn)程決定

from multiprocessing import Process,Queue
from time import sleep


class MyClass(object):
    def __init__(self, time):
        self.time = time    #指定n秒后退出子進(jìn)程
        self.q = Queue()    #將Queue賦值給公共方法

    def sed_msgs(self):
        '''該方法決定子進(jìn)程是否退出'''
        sleep(self.time)    #按指定時(shí)間休眠
        self.q.put('over')  #休眠后向消息隊(duì)列發(fā)送一條消息
        self.q.close()      #關(guān)閉通信

    def proc1(self):
        '''這是一個(gè)子進(jìn)程,當(dāng)消息隊(duì)列不為空則停止運(yùn)行'''
        while self.q.empty():
            print('[進(jìn)程1] 執(zhí)行當(dāng)前任務(wù)...')
            sleep(1)
        self.q.close()

    def proc2(self):
        '''這是一個(gè)子進(jìn)程,當(dāng)消息隊(duì)列不為空則停止運(yùn)行'''
        while self.q.empty():
            print('[進(jìn)程2] 執(zhí)行當(dāng)前任務(wù)...')
            sleep(1)
        self.q.close()

if __name__ == '__main__':
    mc = MyClass(3)     #給定休眠時(shí)間3s

    # 定義子進(jìn)程
    s1 = Process(target=mc.sed_msgs)
    p1 = Process(target=mc.proc1)
    p2 = Process(target=mc.proc2)
    communication = [s1, p1, p2]

    # 啟動(dòng)所有子進(jìn)程
    for s in communication:
        s.start()

    # 等待所有子進(jìn)程結(jié)束
    for s in communication:
        s.join()

    print('====================== 結(jié)束 ======================')

邏輯視圖

3.3. Pipe 管道通信

Pipe()支持多個(gè)進(jìn)程間的管道通信。管道可以被多個(gè)進(jìn)程訪問,但是一次只能有一個(gè)進(jìn)程對管道進(jìn)行操作。

Pipe()賦值給兩個(gè)對象:p1,p2 =Pipe(),p1是發(fā)送消息的方法,p2是接收消息的方法。

定義2個(gè)進(jìn)程,分別向管道中發(fā)送消息和接收消息

from multiprocessing import Pipe, Process

class MyClass(object):
    '''定義2個(gè)管道和2個(gè)進(jìn)程,相互發(fā)送和接收消息'''
    def __init__(self):
        # parent_conn表示發(fā)送消息,child_conn表示接收消息
        self.parent_conn1, self.child_conn1 = Pipe()
        self.parent_conn2, self.child_conn2 = Pipe()

    def proc1(self):
        # 向管道1中發(fā)送消息
        self.parent_conn1.send('蘋果')
        # 接收管道2中的消息
        received_message = self.child_conn2.recv()
        print(f'[進(jìn)程1] 收到的消息是:{received_message}')
        # 關(guān)閉管道1
        self.parent_conn1.close()

    def proc2(self):
        # 接收管道1中的消息
        received_message = self.child_conn1.recv()
        print(f'[進(jìn)程2] 收到的消息是:{received_message}')
        # 向管道2中發(fā)送消息
        self.parent_conn2.send('香蕉')
        # 關(guān)閉管道2
        self.parent_conn2.close()

if __name__ == '__main__':
    mc = MyClass()

    # 創(chuàng)建子進(jìn)程
    p1 = Process(target=mc.proc1)
    p2 = Process(target=mc.proc2)

    # 啟動(dòng)子進(jìn)程
    p1.start()
    p2.start()

    # 等待子進(jìn)程結(jié)束
    p1.join()
    p2.join()

    print('=================== 結(jié)束 ===================')

說明:

進(jìn)程1向管道1中發(fā)送一條消息(蘋果),發(fā)送完成之后接收管道2消息,如果沒有則等待。

進(jìn)程2接收管道1的消息并輸出,再向管道2中發(fā)送一條消息(香蕉)。

4. 進(jìn)程池

4.1. 同步與異步并發(fā)區(qū)別

  • 同步(阻塞):同步并發(fā)類似于串行。例如A、B兩個(gè)進(jìn)程,必須A執(zhí)行完成后才能執(zhí)行B;若A出現(xiàn)意外(一直阻塞),那么B將無法執(zhí)行。
  • 異步(非阻塞):異步并發(fā)在執(zhí)行多任務(wù)時(shí)不會(huì)出現(xiàn)相互阻塞的情況(進(jìn)程池限制除外)。如有A、B兩個(gè)進(jìn)程,他們倆可以同時(shí)進(jìn)行,不會(huì)出現(xiàn)相互等待的情況。

同步的簡單代碼如下

from multiprocessing import Pool
from time import sleep

def proc1():
    '''子進(jìn)程1'''
    for i in range(2):
        print(f'[進(jìn)程1] 執(zhí)行{i}...')
        sleep(1)

def proc2():
    '''子進(jìn)程2'''
    for i in range(2):
        print(f'[進(jìn)程2] 執(zhí)行{i}...')
        sleep(1)

if __name__ == '__main__':
    # 定義進(jìn)程池
    pool = Pool()
    # 同步執(zhí)行2個(gè)子進(jìn)程
    pool.apply(proc1)
    pool.apply(proc2)
    # 關(guān)閉和等待子進(jìn)程結(jié)束
    pool.close()
    pool.join()

使用pool.apply ([函數(shù)名]) 定義同步執(zhí)行:先執(zhí)行 proc1,再執(zhí)行 proc2

使用異步 pool.apply_async ([函數(shù)名])(將上述代碼 apply 修改為 apply_async 即可)

    # 異步執(zhí)行2個(gè)子進(jìn)程
    pool.apply_async (proc1)
    pool.apply_async (proc2)

其結(jié)果為:proc1 和 proc2 同時(shí)執(zhí)行

4.2. 異步調(diào)用多個(gè)子進(jìn)程

進(jìn)程池異步并發(fā)的基本使用方法

'''定義子進(jìn)程'''
def start_func():
    pass
def proc1():
    pass
def proc2():
    pass

if __name__ == '__main__':
    # 定義進(jìn)程池,設(shè)置屬性
    pool = Pool(processes=1, maxtasksperchild=2, initializer=start_func)
    # 異步啟動(dòng)子進(jìn)程(子進(jìn)程為指定的某個(gè)函數(shù))
    pool.apply_async(proc1)
    pool.apply_async(proc2)
    # 關(guān)閉和等待子進(jìn)程結(jié)束
    pool.close()
    pool.join()
  • processes:并發(fā)數(shù)(默認(rèn)為系統(tǒng)最大CPU數(shù))。如果將并發(fā)數(shù)設(shè)置為1,即使進(jìn)程池中有2個(gè)進(jìn)程,也會(huì)按同步的方式執(zhí)行。一般設(shè)置不超過CPU數(shù),最好等于子進(jìn)程數(shù)。
  • maxtasksperchild:可以限制在處理完一定數(shù)量的任務(wù)后重建進(jìn)程池中的工作進(jìn)程,以避免內(nèi)存泄漏和資源泄漏等問題。使用maxtasksperchild參數(shù)時(shí)需要權(quán)衡性能和資源利用率之間的平衡。如果你將maxtasksperchild設(shè)置得太小,進(jìn)程重建的成本可能會(huì)顯著影響性能。如果你將maxtasksperchild設(shè)置得太大,那么可能會(huì)導(dǎo)致內(nèi)存泄漏或資源泄漏。需要根據(jù)具體情況進(jìn)行調(diào)整。
  • initializer:指定每個(gè)工作進(jìn)程啟動(dòng)時(shí)要調(diào)用的函數(shù)(func)。如果子進(jìn)程只有2個(gè),而processes 設(shè)置為3,那么將會(huì)調(diào)用3次 func。

進(jìn)程池中有4個(gè)子進(jìn)程,僅使用2個(gè)并發(fā)(A、B、C、D,先執(zhí)行AB,再執(zhí)行CD)

from multiprocessing import Pool
from time import sleep

def start_func():
    print('啟動(dòng)前此函數(shù)')

def proc1():
    '''子進(jìn)程'''
    for i in range(2):
        print(f'[進(jìn)程1] 執(zhí)行{i}...')
        sleep(1)

def proc2():
    '''子進(jìn)程'''
    for i in range(2):
        print(f'[進(jìn)程2] 執(zhí)行{i}...')
        sleep(1)

def proc3():
    '''子進(jìn)程'''
    for i in range(2):
        print(f'[進(jìn)程3] 執(zhí)行{i}...')
        sleep(1)

def proc4():
    '''子進(jìn)程'''
    for i in range(2):
        print(f'[進(jìn)程4] 執(zhí)行{i}...')
        sleep(1)

if __name__ == '__main__':
    # 定義進(jìn)程池,設(shè)置屬性
    pool = Pool(processes=2, initializer=start_func)
    func_list = [proc1, proc2, proc3, proc4]
    for i in func_list:
        pool.apply_async(i)
    # 關(guān)閉和等待子進(jìn)程結(jié)束
    pool.close()
    pool.join()

結(jié)果如下:

如果將processes 設(shè)置為 4,那么4個(gè)進(jìn)程將會(huì)同時(shí)進(jìn)行

4.3. 進(jìn)程池的高并發(fā)

  • 使用 map 或 imap 迭代來高效完成工作

使用 map(同步) 執(zhí)行高并發(fā)。會(huì)自動(dòng)等待子進(jìn)程完成后,才會(huì)進(jìn)行執(zhí)行主進(jìn)程。

from multiprocessing import Pool

def proc(num):
    print(f'我是子進(jìn)程')

if __name__ == '__main__':
    # 定義線程池,設(shè)置并發(fā)數(shù)為5
    pool = Pool(5)
    # 使用map向函數(shù)傳入多個(gè)參數(shù)(每傳入一個(gè)參數(shù)則會(huì)調(diào)用一次函數(shù),同時(shí)調(diào)用的數(shù)量由進(jìn)程數(shù)決定)
    pool.map(proc, range(5))
    # 關(guān)閉進(jìn)程池
    pool.close()
    print('========== 結(jié)束 ==========')

使用 imap(異步) 執(zhí)行高并發(fā)。調(diào)度子進(jìn)程運(yùn)行后不會(huì)等待,繼續(xù)執(zhí)行主進(jìn)程任務(wù)。若主進(jìn)程運(yùn)行結(jié)束,則子進(jìn)程不論是否運(yùn)行完成都將強(qiáng)制結(jié)束。

from multiprocessing import Pool

def proc(num):
    print(f'我是子進(jìn)程')

if __name__ == '__main__':
    # 定義線程池,設(shè)置并發(fā)數(shù)為5
    pool = Pool(5)
    # 使用map向函數(shù)傳入多個(gè)參數(shù)(每傳入一個(gè)參數(shù)則會(huì)調(diào)用一次函數(shù),同時(shí)調(diào)用的數(shù)量由進(jìn)程數(shù)決定)
    pool.imap(proc, range(5))
    # 關(guān)閉進(jìn)程池
    pool.close()
    print('========== 結(jié)束 ==========')

這樣的好處是不會(huì)影響主進(jìn)程其他任務(wù)的執(zhí)行,如果需要等待則使用 join

from multiprocessing import Pool

def proc(num):
    print(f'我是子進(jìn)程')

if __name__ == '__main__':
    # 定義線程池,設(shè)置并發(fā)數(shù)為5
    pool = Pool(5)
    # 使用map向函數(shù)傳入多個(gè)參數(shù)(每傳入一個(gè)參數(shù)則會(huì)調(diào)用一次函數(shù),同時(shí)調(diào)用的數(shù)量由進(jìn)程數(shù)決定)
    pool.imap(proc, range(5))
    # 關(guān)閉進(jìn)程池
    pool.close()
    # 等待子進(jìn)程結(jié)束
    pool.join()
    print('========== 結(jié)束 ==========')

使用高并發(fā)需要注意以下幾點(diǎn):

  • 注意控制進(jìn)程池?cái)?shù)量,避免過多的進(jìn)程導(dǎo)致性能下降??梢酝ㄟ^調(diào)用multiprocessing.cpu_count()獲取當(dāng)前系統(tǒng)的 CPU 數(shù)量,并使用該值來指定進(jìn)程池?cái)?shù)量。
  • 使用maxtasksperchild參數(shù)來限制每個(gè)工作進(jìn)程可以處理的任務(wù)數(shù)量。這有助于避免長時(shí)間運(yùn)行的進(jìn)程引起的資源泄漏或內(nèi)存泄漏。
  • 在提交任務(wù)之前,盡可能地減小要處理數(shù)據(jù)的大小。例如,可以使用迭代器來代替列表,或者使用生成器來延遲計(jì)算。
  • 對于長時(shí)間運(yùn)行的任務(wù),可以使用消息隊(duì)列來緩解性能問題。例如,可以使用 RabbitMQ 或 ZeroMQ 等消息隊(duì)列來異步處理任務(wù)。

5. 進(jìn)程同步

進(jìn)程同步是指多個(gè)進(jìn)程在共享資源時(shí)的協(xié)調(diào)與同步。單個(gè)進(jìn)程運(yùn)行時(shí) ,程序的執(zhí)行是順序的;多個(gè)進(jìn)程并發(fā)執(zhí)行時(shí)可能會(huì)造成資源競爭、死鎖等問題。因此,進(jìn)程同步是保證每個(gè)進(jìn)程在使用共享資源時(shí)的順序和正確性。

5.1. 進(jìn)程加鎖的方式

兩個(gè)進(jìn)程間實(shí)現(xiàn)同步有2種方式,1、手動(dòng)加鎖(釋放鎖);2、with自動(dòng)加鎖(釋放鎖)

手動(dòng)加鎖方式

from multiprocessing import Lock

def func():
    # 手動(dòng)加鎖
    Lock().acquire()

    '''執(zhí)行程序'''
    pass

    # 手動(dòng)釋放鎖
    Lock().release()

with 自動(dòng)加鎖(釋放鎖)

from multiprocessing import Lock

def func():
    # with自動(dòng)加鎖(釋放鎖)
    with Lock():
        '''執(zhí)行程序'''
        pass

5.2. 實(shí)現(xiàn)進(jìn)程同步

  • 多進(jìn)程之間同步全局變量的修改,則需要使用共享內(nèi)存。需要用到 multiprocessing 模塊中的 Value 共享變量類型來實(shí)現(xiàn)。

實(shí)現(xiàn)邏輯

代碼如下

from multiprocessing import Process, Lock, Value
from time import sleep

def func1(lock, shared_var):
    for i in range(3):
        with lock:
            shared_var.value -= 1
            print(f'[func1] 共享變量var當(dāng)前值為:{shared_var.value}')
            sleep(1)

def func2(lock, shared_var):
    for i in range(3):
        with lock:
            shared_var.value += 2
            print(f'[func2] 共享變量var當(dāng)前值為:{shared_var.value}')
            sleep(1)

if __name__ == '__main__':
    lock = Lock()
    shared_var = Value('i', 10)

    # 定義子進(jìn)程
    p1 = Process(target=func1, args=(lock, shared_var))
    p2 = Process(target=func2, args=(lock, shared_var))

    # 啟動(dòng)子進(jìn)程
    p1.start()
    p2.start()

    # 等待子進(jìn)程結(jié)束
    p1.join()
    p2.join()

    print(f'最終共享變量var為:{shared_var.value}')

錯(cuò)誤示例(使用了多個(gè) Lock() )

from multiprocessing import Process, Lock, Value
from time import sleep

def func1(shared_var):
    '''定義子進(jìn)程,讓共享變量-1,執(zhí)行3次'''
    for i in range(3):
        with Lock():    #錯(cuò)誤語句
            shared_var.value -= 1
            print(f'[func1] 共享變量var當(dāng)前值為:{shared_var.value}')
            sleep(1)

def func2(shared_var):
    '''定義子進(jìn)程,讓共享變量+2,執(zhí)行3次'''
    for i in range(3):
        with Lock():    #錯(cuò)誤語句
            shared_var.value += 2
            print(f'[func2] 共享變量var當(dāng)前值為:{shared_var.value}')
            sleep(1)

if __name__ == '__main__':
    shared_var = Value('i', 10)

    # 定義子進(jìn)程
    p1 = Process(target=func1, args=(shared_var,))
    p2 = Process(target=func2, args=(shared_var,))

    # 啟動(dòng)子進(jìn)程
    p1.start()
    p2.start()

    # 等待子進(jìn)程結(jié)束
    p1.join()
    p2.join()

    print(f'最終共享變量var為:{shared_var.value}')

按同步邏輯執(zhí)行的過程應(yīng)該是:共享變量 = 10(初始) - 1 + 2 - 1 + 2 - 1 + 2 = 13,但實(shí)際的值卻是11。下圖所示,其中一個(gè)步驟出錯(cuò),并沒有加鎖。

仔細(xì)翻看代碼,發(fā)現(xiàn)在每個(gè)函數(shù)中都使用了 Lock() 方法,也就是說2個(gè)函數(shù)使用的是2個(gè)鎖,自然不會(huì)加鎖,出現(xiàn)了臟讀現(xiàn)象。

  • Lock():用于實(shí)現(xiàn)在進(jìn)程間同步訪問共享資源,保證同一時(shí)間只有一個(gè)進(jìn)程能夠訪問。
  • RLock():與Lock類似,但支持遞歸鎖,同一個(gè)進(jìn)程中多次獲取仍然有效。
  • Semaphore():用于控制進(jìn)程間的并發(fā)數(shù)量。
  • Event():用于實(shí)現(xiàn)進(jìn)程間的事件通知。
  • Condition():用于實(shí)現(xiàn)復(fù)雜的進(jìn)程間同步,比如等待某個(gè)條件變?yōu)檎鏁r(shí)再繼續(xù)執(zhí)行。
  • Barrier():用于實(shí)現(xiàn)多個(gè)進(jìn)程間的協(xié)調(diào)操作,比如等待所有進(jìn)程都到達(dá)某個(gè)狀態(tài)再繼續(xù)執(zhí)行。

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • Python寫的英文字符大小寫轉(zhuǎn)換代碼示例

    Python寫的英文字符大小寫轉(zhuǎn)換代碼示例

    這篇文章主要介紹了Python寫的英文字符大小寫轉(zhuǎn)換代碼示例,本文例子相對簡單,本文直接給出代碼實(shí)例,需要的朋友可以參考下
    2015-03-03
  • keras Lambda自定義層實(shí)現(xiàn)數(shù)據(jù)的切片方式,Lambda傳參數(shù)

    keras Lambda自定義層實(shí)現(xiàn)數(shù)據(jù)的切片方式,Lambda傳參數(shù)

    這篇文章主要介紹了keras Lambda自定義層實(shí)現(xiàn)數(shù)據(jù)的切片方式,Lambda傳參數(shù),具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2020-06-06
  • pyspark操作MongoDB的方法步驟

    pyspark操作MongoDB的方法步驟

    這篇文章主要介紹了pyspark操作MongoDB的方法步驟,小編覺得挺不錯(cuò)的,現(xiàn)在分享給大家,也給大家做個(gè)參考。一起跟隨小編過來看看吧
    2019-01-01
  • Python判斷MySQL表是否存在的兩種方法

    Python判斷MySQL表是否存在的兩種方法

    在數(shù)據(jù)庫開發(fā)中,經(jīng)常需要檢查某個(gè)表是否存在,如果不存在則創(chuàng)建它,下面我們就來看看如何使用Python連接MySQL數(shù)據(jù)庫并實(shí)現(xiàn)表存在性檢查與創(chuàng)建的功能吧
    2026-02-02
  • python實(shí)現(xiàn)C4.5決策樹算法

    python實(shí)現(xiàn)C4.5決策樹算法

    這篇文章主要為大家詳細(xì)介紹了python實(shí)現(xiàn)C4.5決策樹算法,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2018-08-08
  • Python實(shí)現(xiàn)把json格式轉(zhuǎn)換成文本或sql文件

    Python實(shí)現(xiàn)把json格式轉(zhuǎn)換成文本或sql文件

    這篇文章主要介紹了Python實(shí)現(xiàn)把json格式轉(zhuǎn)換成文本或sql文件,本文直接給出代碼實(shí)例,需要的朋友可以參考下
    2015-07-07
  • Pytorch之view及view_as使用詳解

    Pytorch之view及view_as使用詳解

    今天小編就為大家分享一篇Pytorch之view及view_as使用詳解,具有很好的參考價(jià)值,希望對大家有所幫助。一起跟隨小編過來看看吧
    2019-12-12
  • pandas中使用數(shù)據(jù)透視表的示例代碼

    pandas中使用數(shù)據(jù)透視表的示例代碼

    本文主要介紹了pandas中使用數(shù)據(jù)透視表的示例代碼,主要包含pivot_table函數(shù)的使用,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2024-12-12
  • 玩轉(zhuǎn)串口通信:利用pyserial庫,Python打開無限可能

    玩轉(zhuǎn)串口通信:利用pyserial庫,Python打開無限可能

    想要學(xué)習(xí)如何使用pyserial庫實(shí)現(xiàn)串口通信嗎?這篇指南將帶你一步步了解Python中的串口通信,無論是控制硬件設(shè)備還是與外部設(shè)備進(jìn)行數(shù)據(jù)交換,pyserial庫都能為你提供便捷的解決方案,快來跟著我們的指南,輕松掌握串口通信的技巧吧!
    2023-11-11
  • Pytorch中torch.cat()函數(shù)舉例解析

    Pytorch中torch.cat()函數(shù)舉例解析

    一般torch.cat()是為了把多個(gè)tensor進(jìn)行拼接而存在的,下面這篇文章主要給大家介紹了關(guān)于Pytorch中torch.cat()函數(shù)舉例解析的相關(guān)資料,文中通過示例代碼介紹的非常詳細(xì),需要的朋友可以參考下
    2022-12-12

最新評論

措美县| 独山县| 门源| 曲沃县| 潮州市| 临城县| 鄂托克前旗| 正蓝旗| 武山县| 新巴尔虎左旗| 武威市| 富顺县| 商丘市| 广宁县| 绥江县| 宕昌县| 淳化县| 永城市| 手游| 威信县| 黔江区| 乌拉特中旗| 济南市| 柘荣县| 绩溪县| 健康| 临夏市| 沈阳市| 邹平县| 疏勒县| 马边| 郧西县| 比如县| 江都市| 牟定县| 东明县| 云安县| 东丰县| 永福县| 独山县| 拉萨市|