使用asyncio Queue异步执行代码异常:数据丢失问题排查与优化建议
问题排查与优化方案
核心问题分析
你的代码仅输出current_value=10的原因有两个:
retrieve函数提前退出:每个锁处理块后都有break语句,导致处理完第一个item(10)后直接跳出循环,后续11-15的item从未被读取。- 单协程无法并发处理:每个
main仅启动一个retrieve协程,即使去掉break,也只能串行处理item,无法利用5个模拟浏览器的并发能力。
优化后的代码
import asyncio import random import time async def retrieve(queue, semaphore): while True: item = await queue.get() if item is None: queue.task_done() break # 收到终止信号后退出 async with semaphore: # 模拟浏览器请求耗时 await asyncio.sleep(random.uniform(3, 5)) print(item) queue.task_done() # 标记item处理完成 async def process(queue, x, y): start_value = 10 stop_value = 15 current_value = start_value # 批量添加任务到队列 while current_value <= stop_value: await queue.put([x, y, current_value]) current_value += 1 # 给每个worker发送终止信号 for _ in range(5): await queue.put(None) async def main(x, y): semaphore = asyncio.Semaphore(5) # 限制最多5个并发"浏览器" queue = asyncio.Queue() # 启动5个worker协程 workers = [asyncio.create_task(retrieve(queue, semaphore)) for _ in range(5)] # 填充任务队列 await process(queue, x, y) # 等待所有任务处理完成 await queue.join() # 清理worker协程 for worker in workers: worker.cancel() await asyncio.gather(*workers, return_exceptions=True) async def run_all(): # 创建所有x,y组合的任务 tasks = [] for x in range(2): for y in range(5): tasks.append(main(x, y)) # 并发执行所有任务 await asyncio.gather(*tasks) if __name__ == "__main__": start = time.perf_counter() asyncio.run(run_all()) end = time.perf_counter() print(f'Time elapsed:{end - start:.2f}s')
关键优化点
- 替换多锁为Semaphore:用
Semaphore(5)替代5个独立锁,更简洁地实现并发数控制,避免重复的锁判断逻辑。 - 多Worker协程:启动5个
retrieve协程,充分利用5个模拟浏览器的并发能力,同时处理多个item。 - 移除break语句:让
retrieve循环持续处理item,直到收到None终止信号。 - 全局并发执行:将所有
main(x,y)任务放入asyncio.gather并发执行,避免串行等待每个x,y组合。 - 任务清理机制:通过
queue.join()等待所有任务完成,发送None信号让worker正常退出,避免孤儿协程。
内容的提问来源于stack exchange,提问作者jguy
相关产品推荐
相关产品推荐

