如何在运行时向asyncio添加任务以实现单消息单协程异步处理
问题根因
你当前的消费逻辑是串行执行的:仅启动了1个output协程死循环消费队列,拿到每条数据后会同步等待5秒再打印,处理完第一条才会消费第二条,所以两次等待时间会累加,导致第二条输出比第一条晚5秒。
解决方案
调整消费逻辑,每从队列拿到一条数据就单独创建子协程处理延迟和打印,不要阻塞主消费循环,这样多条数据的等待过程会并行执行。
修改后完整代码
import asyncio # 单独封装单条数据的处理逻辑 async def process_single_data(data): await asyncio.sleep(5) print(f"<<< {data}") async def output(queue): while True: data = await queue.get() if data: # 直接创建异步任务,不等待执行完成,继续消费下一条队列数据 asyncio.create_task(process_single_data(data)) async def main(): queue = asyncio.Queue() asyncio.create_task(output(queue)) loop = asyncio.get_event_loop() while True: data = await loop.run_in_executor(None, input) await queue.put(data) if __name__ == "__main__": loop = asyncio.get_event_loop() loop.create_task(main()) loop.run_forever()
效果验证
调整后输入1后快速输入2,两条数据的5秒等待会同时计时,5秒后会先后打印<<< 1和<<< 2,间隔只有协程调度的毫秒级延迟,符合预期效果。
内容的提问来源于stack exchange,提问作者vladkhard
相关产品推荐
相关产品推荐

