如何让multiprocessing的Queue与Pool函数配合使用并分配4个进程
用multiprocessing.Pool改造你的生产消费代码
嘿,我来帮你把这段代码改成用multiprocessing.Pool的版本!首先得明确你原来的逻辑是一个生产者持续生成数据,一个消费者处理数据,现在要换成4个消费者进程来处理,用Pool来管理这4个进程是很合适的。
先给你直接上修改后的代码,然后再拆解说明改动点:
import random from multiprocessing import Pool, Queue, Process import time from datetime import datetime def producer(q): # 原来的function1,改名让语义更清晰 while True: daydate = datetime.now() number = random.randrange(1, 215) print('Sent to consumer: ({}, {})'.format(daydate, number)) q.put((daydate, number)) time.sleep(2) def consumer(q): # 原来的function2,现在作为Pool中每个进程的执行函数 while True: try: date, number = q.get(timeout=1) # 加超时避免进程一直阻塞 print("Received values from producer: ({}, {})".format(date, number)) time.sleep(2) except: # 超时后继续循环,不影响运行 continue if __name__ == "__main__": q = Queue() # 启动独立的生产者进程,设置为守护进程,主进程退出时自动结束 producer_process = Process(target=producer, args=(q,)) producer_process.daemon = True producer_process.start() # 创建包含4个进程的Pool,每个进程都运行consumer函数 with Pool(4) as pool: # 给Pool里的每个进程分配consumer任务 [pool.apply_async(consumer, args=(q,)) for _ in range(4)] try: # 让主进程保持运行,直到你按下Ctrl+C终止 while True: time.sleep(1) except KeyboardInterrupt: print("\nStopping all processes...") pool.terminate() pool.join()
关键改动说明:
- 角色拆分更清晰:把原来的
function1和function2改名为producer(生产者)和consumer(消费者),一眼就能看懂各自的职责。 - 生产者独立运行:还是用单独的
Process来跑生产者逻辑,避免它阻塞主进程,同时设置为守护进程,保证主进程退出时它会自动停止。 - 用Pool管理消费者进程:创建4个进程的
Pool,每个进程都执行consumer函数,共享同一个队列。这样4个消费者会自动从队列里抢任务处理,比单个消费者效率更高。 - 优雅退出机制:给
q.get()加了1秒超时,避免进程一直卡在取队列的操作上;同时主进程监听键盘中断(Ctrl+C),收到中断后会终止Pool并等待所有进程结束,不会留下僵尸进程。
补充小提示:
如果你的需求不是无限运行,而是有明确的任务结束条件,可以给生产者加一个退出逻辑,比如生成N条数据后往队列里放一个哨兵值(比如None),然后消费者收到这个值就退出循环,这样Pool的进程就会自动结束啦。
内容的提问来源于stack exchange,提问作者Joemoreneau
相关产品推荐
相关产品推荐

