一聚教程网:一个值得你收藏的教程网站

最新下载

热门教程

如何手动实现一个具备"优先级调度"能力的简易版异步任务执行引擎

时间:2026-07-27 16:12:02 编辑:袖梨 来源:一聚教程网

可直接用 asyncio.PriorityQueue 搭建轻量可控的优先级任务引擎:将任务封装为 (priority, coro) 元组入队,由永驻调度协程循环取高优任务 await 执行,并通过 task_done()、异常捕获、超时控制及取消令牌实现健壮调度。

可以直接用 asyncio.PriorityQueue 搭建一个轻量、可控的优先级任务引擎,核心是把“任务”封装成带优先级的协程,并由一个长期运行的调度协程按序取出执行。

用 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.Eventawait event.wait() 实现条件阻塞

扩展支持取消和超时

真实场景中,高优任务可能需要抢占或限时执行。可以在调度协程中加入超时控制,或对任务本身做包装。

  • 执行前用 asyncio.wait_for(coro, timeout=2.0) 限制最大耗时
  • 为任务附加取消令牌:传入 asyncio.CancelledError 捕获逻辑,或用 asyncio.shield() 保护关键步骤不被意外取消
  • 若需支持“取消排队中任务”,PriorityQueue 本身不提供 remove 接口,可改用自定义堆 + 标记删除,或换用 heapq + asyncio.Lock 手动维护

热门栏目