如何让asyncio Worker任务创建后直接启动,无需调用asyncio.sleep?
你遇到的这个问题本质上是asyncio事件循环的调度时机问题:asyncio.create_task()只是把Worker任务注册到事件循环中,但当前的main协程还在持续运行,事件循环没有机会切换去执行新创建的Worker任务——直到遇到一个会让出控制权的await操作。
原来的代码里用await asyncio.sleep(0.00001)就是为了强制触发一次调度,但其实完全可以去掉这个冗余操作,因为后续的await queue.join()本身就会让出控制权,事件循环会自动调度已经注册的Worker任务启动执行。
修改后的代码
import asyncio import random import time async def worker(name, queue): print(f"starting worker {name}") while True: sleep_for = await queue.get() await asyncio.sleep(sleep_for) queue.task_done() print(f'{name} has slept for {sleep_for:0.2f} seconds') async def main(n): queue = asyncio.Queue() total_sleep_time = 0 for _ in range(20): sleep_for = random.uniform(0.05, 1.0) total_sleep_time += sleep_for await queue.put(sleep_for) tasks = [] for i in range(n): task = asyncio.create_task(worker(f'worker-{i}', queue)) tasks.append(task) started_at = time.monotonic() # 移除多余的sleep操作,直接等待队列任务完成 await queue.join() total_slept_for = time.monotonic() - started_at for task in tasks: task.cancel() # 等待所有Worker任务被取消 await asyncio.gather(*tasks, return_exceptions=True) print('====') print(f'{n} workers slept in parallel for {total_slept_for:.2f} seconds') print(f'total expected sleep time: {total_sleep_time:.2f} seconds') if __name__ == '__main__': import sys n = 3 if len(sys.argv) == 1 else int(sys.argv[1]) # 修正原代码的小问题:将命令行参数转为整数 asyncio.run(main(n))
补充说明
为什么
await queue.join()能触发调度?
当main协程执行到await queue.join()时,它会进入挂起状态,等待队列中所有任务都被标记为task_done。此时事件循环会去查找所有可运行的任务(也就是我们刚创建的Worker任务),并开始执行它们——Worker会立即打印"starting worker",然后从队列中获取任务开始处理。如果队列初始为空怎么办?
哪怕队列一开始没有元素,Worker任务在执行到await queue.get()时会挂起,直到有新元素被放入队列,事件循环会自动唤醒它,完全不需要额外的sleep操作。极端场景下的手动调度
如果在某些场景下,你需要在创建任务后立即强制调度(比如后续没有其他await操作),可以用await asyncio.sleep(0)代替原来的长sleep——这是一个轻量的让出操作,只会让事件循环调度一次其他任务,不会产生实际的延迟。
内容的提问来源于stack exchange,提问作者ZeroCool

