使用Python multiprocessing.Queue仅获首个元素,求问题原因
问题分析与解决:multiprocessing.Queue只获取到第一个元素
你的代码里有两个关键问题导致只能拿到第一个元素:
消费者进程只执行了一次取队列操作
function_to_get_from_q函数里仅有一行print(Queue.get()),执行完这行代码后函数就结束了,对应的process2进程随即退出,自然不会继续处理队列里剩下的9999个元素。主进程没有等待子进程完成
主进程启动两个子进程后直接结束,可能导致子进程还没完成任务就被强制终止,即便恢复循环也可能出现元素未处理完的情况。
修复方案一:恢复循环并等待子进程完成
把被注释的while not Queue.empty()循环恢复,同时在主进程中调用join()等待子进程执行完毕:
import multiprocessing as mp import datetime as dt def function_to_get_from_q(Queue): while not Queue.empty(): print(Queue.get()) def collect(Queue): for i in range(10000): Queue.put([i, (dt.datetime.utcnow() + dt.timedelta(hours=5, minutes=30)).strftime('%H:%M:%S')]) if __name__ == "__main__": Q = mp.Queue() process1 = mp.Process(target=collect, args=(Q,)) process2 = mp.Process(target=function_to_get_from_q, args=(Q,)) process1.start() process2.start() # 等待生产者进程把所有元素放入队列 process1.join() # 等待消费者进程处理完所有元素 process2.join()
注意:Queue.empty()在多进程场景下不是绝对可靠的——如果消费者判断队列空的时候,生产者还在往队列里放元素,可能会提前退出。这种方案适合生产者先完成所有写入,再让消费者处理的场景。
修复方案二:使用终止信号(更可靠)
通过发送特定的终止标记(比如None)让消费者进程明确知道何时停止,这种方式适用于生产者和消费者异步执行的场景:
import multiprocessing as mp import datetime as dt def function_to_get_from_q(Queue): while True: item = Queue.get() if item is None: # 收到终止信号就退出循环 break print(item) def collect(Queue): for i in range(10000): Queue.put([i, (dt.datetime.utcnow() + dt.timedelta(hours=5, minutes=30)).strftime('%H:%M:%S')]) Queue.put(None) # 生产者完成后发送终止信号 if __name__ == "__main__": Q = mp.Queue() process1 = mp.Process(target=collect, args=(Q,)) process2 = mp.Process(target=function_to_get_from_q, args=(Q,)) process1.start() process2.start() process1.join() process2.join()
这种方式不会因为队列暂时为空就提前退出,必须收到明确的终止信号才停止,是多进程队列消费的推荐写法。
内容的提问来源于stack exchange,提问作者User1917931829
相关产品推荐
相关产品推荐

