asyncio.priorityqueue 是 asyncio 生态中唯一原生支持优先级的队列,基于 heapq 实现,要求元素为可比较元组(如 (priority, item)),priority 越小优先级越高。

asyncio.PriorityQueue 是唯一原生支持优先级的队列
Python 标准库中,asyncio.PriorityQueue 是 asyncio 生态里**唯一开箱即用、线程安全且协程友好的优先级队列**。它底层基于 heapq,要求每个入队元素是可比较的(通常是元组,如 (priority, item))。别试图给 asyncio.Queue 手动加排序逻辑——它不维护顺序,put_nowait() 和 get() 都不保证优先级行为。
常见错误现象:RuntimeError: await was called on a coroutine object created by a non-async function,往往是因为误把同步的 queue.PriorityQueue 丢进 await queue.get();或者用 asyncio.Queue 后手动 sorted(queue._queue) ——这不仅非法访问私有属性,还破坏了 asyncio 的事件循环调度契约。
- 必须用
(priority, item)元组形式入队,priority 通常为 int 或 float;数值越小,优先级越高(符合 heapq 行为) - item 本身不能是不可比较对象(如 dict),否则
get()时抛TypeError: ' - 不支持重复修改已入队元素的优先级;如需动态调整,得先取出、再按新 priority 重新
put()
如何让任务对象自带优先级并支持取消
直接塞裸数据不够灵活,实际业务中你往往需要带状态的任务(比如含超时、重试次数、回调函数)。这时推荐封装一个 PriorityTask 类,让它实现 __lt__ 协议,避免依赖元组位置和类型限制。
使用场景:爬虫任务按响应时间预期排序、告警消息按严重等级分级处理、后台作业按 SLA 分级调度。
- 在
__lt__(self, other)中只比较 priority 字段,不要引入时间戳或随机数,否则 heapq 可能崩溃 - 把
asyncio.Task对象存为属性,并在__del__或显式cancel()时调用task.cancel(),防止任务泄漏 - 示例结构:
class PriorityTask: def __init__(self, priority: int, coro): self.priority = priority self.coro = coro self.task = None def __lt__(self, other): return self.priority
并发消费时如何避免高优任务被低优任务阻塞
即使用了 asyncio.PriorityQueue,如果消费者协程里执行的是 CPU 密集型或未加 await 的同步阻塞操作(比如 time.sleep(1)),高优任务仍会卡住。asyncio 不会自动抢占,优先级只体现在「谁先从队列里被 get() 出来」,不等于「谁先被执行完」。
性能影响:一个耗时 5 秒的低优任务若没做异步切分,会让后续所有高优任务至少等待 5 秒才开始执行——队列空了也没用,因为消费者被占着。
- 消费逻辑中禁止出现
time.sleep()、requests.get()、json.loads(大字符串)等同步阻塞调用 - 必须用
await asyncio.sleep()、aiohttp.ClientSession、asyncio.to_thread()(Python 3.9+)包裹 CPU 密集操作 - 考虑设置消费者数量上限(如
asyncio.create_task()控制并发数),否则大量低优任务并发启动反而挤占事件循环资源
生产环境要注意的边界问题
真实系统里,任务可能因异常退出、消费者崩溃、网络分区而堆积。这时候 PriorityQueue 本身不会自动清理或降级,全靠你设计兜底策略。
容易被忽略的地方:队列满时 put() 默认会挂起,但如果你设置了 maxsize 却没配超时或拒绝策略,整个生产者链路就卡死了;另外,PriorityQueue 不提供「按优先级范围批量取」或「查询某优先级是否存在任务」的接口,这些都得自己 wrap。
- 始终为
put()加timeout(如await queue.put(item, timeout=3)),超时后走降级逻辑(日志 + 告警 + 丢弃或转存) - 定期用
queue.qsize()监控堆积量,但注意它不是原子操作,仅作参考;真要精确统计,得用len(queue._queue)(接受私有属性风险)或自定义计数器 - 不要依赖
queue.empty()判断是否还有任务——它返回 False 不代表马上能get(),中间可能有其他协程插入
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











