asyncio Queue多消费者场景下如何保证处理结果按生产顺序输出
解决方案
要保留多消费者并行处理的性能优势,同时保证最终输出顺序和生产顺序严格一致,核心思路是处理并行化、输出序列化,具体实现如下:
实现思路
- 复用生产者生成的天然递增序号
i作为全局顺序标识 - 新增三类共享状态,并用
asyncio.Lock做并发保护:- 结果缓存字典:存储已处理完成、但还未轮到输出的任务结果,key为任务序号
- 下一个待输出序号:初始值为0,标识当前应该输出的最小序号任务
- 互斥锁:保护共享状态读写、以及结果打印操作,避免竞态问题
- 消费者任务处理环节完全并行,仅在结果写入缓存和输出环节加锁:处理完成后先将结果存入缓存,再循环检查缓存中是否存在下一个待输出序号的结果,存在则连续输出直到出现空缺,再释放锁
可运行修改后代码
import asyncio import random import time # 模拟耗时不固定的业务处理逻辑 def do_sometask_and_return_the_result(n): # 模拟随机耗时,若为IO密集型异步任务可改为async实现,调用时加await即可 time.sleep(random.random() * 2) return n * 2 async def produce(Q, n_jobs): for i in range(n_jobs): print(f"Producing :{i}") await Q.put(i) async def consume(Q, result_cache, next_output_idx, lock): while True: n = await Q.get() print(f"Consumed :{n}") # 任务处理环节无锁,完全并行执行,保证性能 x = do_sometask_and_return_the_result(n) # 结果存储和输出环节加锁,保证顺序 async with lock: result_cache[n] = x # 连续输出所有已就绪的、顺序靠前的结果 while next_output_idx[0] in result_cache: current_n = next_output_idx[0] res = result_cache.pop(current_n) print(f"Finished :{current_n} and Result: {res}") next_output_idx[0] += 1 Q.task_done() async def main(): Q = asyncio.Queue(maxsize=3) # 共享状态,用list包装int是为了方便协程内修改全局计数 result_cache = {} next_output_idx = [0] lock = asyncio.Lock() # 启动生产者任务 produce_task = asyncio.create_task(produce(Q, 10)) # 启动3个消费者任务 consume_tasks = [asyncio.create_task(consume(Q, result_cache, next_output_idx, lock)) for _ in range(3)] # 等待生产者完成所有任务生产 await produce_task # 等待队列中所有任务处理完成 await Q.join() # 终止消费者死循环 for t in consume_tasks: t.cancel() print("Done") if __name__ == "__main__": asyncio.run(main())
其他可选实现
你也可以单独启动一个输出协程,新增一个结果队列存储所有消费者返回的(序号, 结果)元组,由输出协程负责缓存结果并按顺序输出,逻辑和上述方案一致,适合结果输出逻辑更复杂的场景。
内容的提问来源于stack exchange,提问作者Madhu Soodhan
相关产品推荐
相关产品推荐

