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
相关产品推荐
相关产品推荐

