可直接用 asyncio.PriorityQueue 搭建轻量可控的优先级任务引擎:将任务封装为 (priority, coro) 元组入队,由永驻调度协程循环取高优任务 await 执行,并通过 task_done()、异常捕获、超时控制及取消令牌实现健壮调度。
可以直接用 asyncio.PriorityQueue 搭建一个轻量、可控的优先级任务引擎,核心是把“任务”封装成带优先级的协程,并由一个长期运行的调度协程按序取出执行。
Python 的 asyncio.PriorityQueue 是线程安全、协程友好的优先队列,数值越小优先级越高。它天然适合作为任务入口:所有待执行的协程都以 (priority, coro) 元组形式入队。
queue = asyncio.PriorityQueue()
await queue.put((1, fetch_user()))(高优)、await queue.put((5, write_log()))(低优)put 协程对象,必须 await 它——所以实际应传入已 awaitable 的协程调用,如 fetch_user() 而非 fetch_user
这个协程持续监听队列,每次取最高优先级任务并 await 执行,执行完标记完成。它不退出,靠 queue.join() 控制生命周期。
while True 循环 + await queue.get() 阻塞等待新任务(priority, coro) 后,先 try/finally 包裹,确保无论成功失败都调用 queue.task_done()
except Exception 记录错误但不中断循环任务不是自动运行的,需要显式创建调度协程,并用 create_task 启动;外部通过 put 注入任务,用 join 等待全部结束。
main() 中调用 asyncio.create_task(dispatcher(queue)) 启动调度器asyncio.gather() 或 asyncio.wait_for() 管理多个调度器实例(如按业务域隔离)asyncio.Event,await event.wait() 实现条件阻塞真实场景中,高优任务可能需要抢占或限时执行。可以在调度协程中加入超时控制,或对任务本身做包装。
asyncio.wait_for(coro, timeout=2.0) 限制最大耗时asyncio.CancelledError 捕获逻辑,或用 asyncio.shield() 保护关键步骤不被意外取消heapq + asyncio.Lock 手动维护