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

Asyncio队列间通信未按预期工作问题排查

问题描述

我正在实现多队列协作的流水线处理模式,用哨兵值在队列间传递停止信号,但运行结果存在两个困惑点:

  • write_task()读取write_q时,第一个传入的值是哨兵None,而非request_task()按顺序放入的任务结果,预期应该按任务创建顺序接收处理
  • write_task()中发现哨兵后打印qsize()显示队列有0个元素,但main函数中打印write_q的qsize()时却显示有2个元素。另外程序会出现无限阻塞的情况,我知道aiofiles使用了run_in_executor(),这可能和队列处理分歧有关
复现代码
import asyncio
import aiofiles
import aiocsv
import json

async def fetch(t: float) -> dict:
    print(f"INFO: Sleeping for {t}s")
    await asyncio.sleep(t)
    return t

async def task(l: list,  request_q: asyncio.Queue) -> None:

    # Read tasks from source of data
    for i in l:
        await request_q.put(
            asyncio.create_task(fetch(i))
        )

    # Sentinel value to signal we are done receiving from source
    await request_q.put(None)

async def request_task(request_q: asyncio.Queue, write_q: asyncio.Queue) -> None:
    while True:
        req = await request_q.get()

        # If we received sentinel for tasks, pass message to next queue
        if not req:
            print("INFO: received sentinel from request_q")
            request_q.task_done()
            await request_q.put(None) # put back into the queue to signal to other consumers we are done
            break
        
        # Make the request
        resp = await req
        await write_q.put(resp)
        request_q.task_done()

async def write_task(write_q: asyncio.Queue) -> None:

    headers: bool = True
    async with aiofiles.open("file.csv", mode="w+", newline='') as f:
        w = aiocsv.AsyncWriter(f)
        while True:
            # Get data out of the queue to write it
            data = await write_q.get()
            print(data)

            # if not data:
            #     print(f"INFO: Found sentinel in write_task, queue size was: {write_q.qsize()}")
            #     write_q.task_done()
            #     await f.flush()
            #     break

            if headers:
                await w.writerow([
                    "status",
                    "data",
                ])
                headers = False

            # Write the data from the response
            await w.writerow([
                "200",
                json.dumps(data)
            ])
            await f.flush()
            write_q.task_done()

async def main() -> None:

    # Create fake data to POST
    items: list[str] = [.2, .5, 1] 

    # Queues for orchestrating 
    request_q = asyncio.Queue()
    write_q = asyncio.Queue()

    # one producer
    producer = asyncio.create_task(
        task(items, request_q)
    )

    # 5 request consumers
    request_consumers = [
        asyncio.create_task(
            request_task(request_q, write_q)
        )
        for _ in range(2)
    ]

    # 5 write consumers
    write_consumer = asyncio.create_task(
        write_task(write_q)
    )

    errors = await asyncio.gather(producer, return_exceptions=True)
    print(f"INFO: Producer has completed! exceptions: {errors}")

    await request_q.join()
    for c in request_consumers:
        c.cancel()
    print("INFO: request consumer has completed! ")
    print(f"INFO: write_q in main qsize: {write_q.qsize()}")
    
    await write_q.join()
    print("INFO: write queue has completed! ")
    # await write_consumer
    write_consumer.cancel()
    print("INFO: Complete!")

if __name__ == "__main__":
    # loop = asyncio.new_event_loop()
    # loop.run_until_complete(main())
    asyncio.run(main())

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:01:13