Python multiprocessing.Queue未按先进先出(FIFO)工作的问题
嘿,我来帮你捋清楚这个问题——你遇到的其实是对multiprocessing.Queue的FIFO特性的常见误解,咱们一步步拆解:
问题出在哪?
Queue本身确实是先进先出的,但这里的「进」指的是进程向队列写入数据的顺序,而不是你提交任务的顺序。因为不同帧的处理耗时可能不一样(比如帧3的内容更简单,处理更快),先完成的进程会先把结果put进队列,所以你用q.get()拿到的自然是进程完成的顺序,而不是你提交任务的1、2、3...顺序。
另外看你的代码还有个小冗余:你把同一个queue实例反复加到queueList里,最后遍历queueList其实是多次调用同一个队列的get()——虽然能拿到结果,但完全没必要,直接循环对应次数调用get()就行。
两种解决思路,任你选
思路1:给结果绑顺序标识,收集后排序
这种方法改动最小,只需要在返回结果时带上帧的编号,等所有结果收齐后再按帧号排序:
修改后的代码片段:
from multiprocessing import Process, Queue def initiateAnalysis(self): ''' Starts the SVD analysis using multiprocessing ''' queue = Queue() jobs = [] collectResponses = [] for frame in self.commonFrames: p = Process(target=self.collectResponsesFromFrame, args=(frame, queue)) p.start() print(frame, ' process started') jobs.append(p) # 收集所有结果,每个结果包含帧号和分析数据 for _ in range(len(self.commonFrames)): collectResponses.append(queue.get()) # 等待所有进程跑完 for p in jobs: p.join() # 按帧号排序,恢复提交顺序 collectResponses.sort(key=lambda x: x['frameNo']) # 验证顺序是否正确 for item in collectResponses: print(item['frameNo']) def collectResponsesFromFrame(self, frame, queue): # 这里替换成你的实际分析逻辑 analysis_result = {"frameNo": frame, "data": "你的分析结果"} queue.put(analysis_result)
思路2:用multiprocessing.Pool.map(更简洁推荐)
如果你不想手动管理队列和排序,Pool.map是更省心的选择——它会自动保证输出列表的顺序和输入的任务顺序完全一致,不管进程实际完成的先后:
from multiprocessing import Pool def initiateAnalysis(self): ''' Starts the SVD analysis using multiprocessing ''' # 创建进程池,默认用CPU核心数,也可以手动指定比如Pool(4) with Pool() as pool: # map会把self.commonFrames里的每个帧传给处理函数,返回结果按输入顺序排列 collectResponses = pool.map(self.collectResponsesFromFrame, self.commonFrames) # 直接输出就是提交顺序 for item in collectResponses: print(item['frameNo']) def collectResponsesFromFrame(self, frame): # 替换成你的实际分析逻辑,直接返回结果即可 analysis_result = {"frameNo": frame, "data": "你的分析结果"} return analysis_result
再划个重点
Queue的FIFO是队列内部的存取规则——如果进程A先put,进程B后put,那get()肯定先拿到A的结果。但你的场景里,进程完成顺序和提交顺序不一致,导致put的顺序乱了,自然get出来的结果也乱了。
内容的提问来源于stack exchange,提问作者jxw
相关产品推荐
相关产品推荐

