You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用asyncio Queue异步执行代码异常:数据丢失问题排查与优化建议

问题排查与优化方案

核心问题分析

你的代码仅输出current_value=10的原因有两个:

  1. retrieve函数提前退出:每个锁处理块后都有break语句,导致处理完第一个item(10)后直接跳出循环,后续11-15的item从未被读取。
  2. 单协程无法并发处理:每个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')

关键优化点

  1. 替换多锁为Semaphore:用Semaphore(5)替代5个独立锁,更简洁地实现并发数控制,避免重复的锁判断逻辑。
  2. 多Worker协程:启动5个retrieve协程,充分利用5个模拟浏览器的并发能力,同时处理多个item。
  3. 移除break语句:让retrieve循环持续处理item,直到收到None终止信号。
  4. 全局并发执行:将所有main(x,y)任务放入asyncio.gather并发执行,避免串行等待每个x,y组合。
  5. 任务清理机制:通过queue.join()等待所有任务完成,发送None信号让worker正常退出,避免孤儿协程。

内容的提问来源于stack exchange,提问作者jguy

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.03 05:10:56