asyncio.PriorityQueue不是开箱即用的异步优先队列,因其线程不安全、仅限单事件循环、无取消感知与超时控制,且优先级值越小越先出队;需封装支持超时与任务取消,并在业务层引入公平性约束防止低优先级任务饿死。

为什么 asyncio.PriorityQueue 不是开箱即用的“异步优先队列”?
因为 asyncio.PriorityQueue 本身是线程不安全的,且只在协程上下文中工作——它不能直接被多个事件循环或不同线程调用。更关键的是,它的 get() 和 put() 是 awaitable,但**没有内置的取消感知或超时控制**,一旦某个高优先级任务卡住(比如 await 了一个永不 resolve 的 Future),低优先级任务就会无限阻塞。
- 它底层基于
heapq,优先级值越小越先出队(注意不是“越大越优先”) - 所有操作必须在同一个
asyncio.EventLoop中执行;跨 loop 调用会抛RuntimeError: Task attached to a different loop - 如果 put 时传入非可比较对象(如 dict、自定义类没实现
__lt__),会直接 raiseTypeError
如何让优先级队列支持任务取消和超时?
直接封装 asyncio.PriorityQueue,在 get() 上加 timeout,并对每个任务包装成带 cancel 支持的 asyncio.Task:
import asyncio
from typing import Any, Callable, Awaitable
<p>class PriorityTaskQueue:
def <strong>init</strong>(self):
self._queue = asyncio.PriorityQueue()
self._shutdown = False</p><pre class="brush:php;toolbar:false;">async def put(self, priority: int, coro: Awaitable, *args, **kwargs):
# 用 task 包装,便于后续 cancel
task = asyncio.create_task(coro(*args, **kwargs))
await self._queue.put((priority, task))
async def get(self, timeout: float = None) -> asyncio.Task:
try:
_, task = await asyncio.wait_for(self._queue.get(), timeout=timeout)
return task
except asyncio.TimeoutError:
raise TimeoutError("No task available within timeout")
async def join(self):
while not self._queue.empty():
await self._queue.join()
注意:这里把 coro 提前转成 Task,而不是存函数——否则 get 出来还要再 await,就失去“按优先级调度执行”的意义了。
怎么避免高优先级任务饿死低优先级任务?
纯 heapq 实现的优先队列容易导致“优先级倾轧”:只要不断有 priority=0 的任务进来,priority=5 的任务永远没机会执行。解决办法不是改算法,而是**在业务层引入“公平性约束”**:
快速生成专业的 Python 脚本和应用代码。一键创建完整项目结构,支持CLI、API、爬虫、Bot、Django等多种项目类型,包含完整的项目结构、配置文件、依赖管理、测试、README和文档。
- 使用滑动窗口计数:记录过去 N 秒内 priority ≤ 1 的任务数量,超过阈值则自动降权(例如临时 +2)
- 给每个优先级设置最大连续执行次数(如最多连跑 3 个 high-priority 任务,就强制 yield 一次
await asyncio.sleep(0)) - 用时间戳 + 优先级复合键:
(priority, timestamp),保证相同优先级下先进先出,防止老任务被新同级任务无限插队
示例复合键写法:await queue.put((priority, time.time(), task_id), task) —— 注意必须确保 time.time() 在 put 时计算,不能在 task 内部懒算。
生产环境要注意 asyncio.Queue 和 PriorityQueue 的混用陷阱
别在同一个队列里既用 asyncio.Queue 又用 asyncio.PriorityQueue 做“混合调度”。它们的内部锁机制不同:PriorityQueue 用的是 asyncio.Lock,而普通 Queue 默认用 asyncio.Semaphore 控制 size,混用会导致 get() 行为不一致,甚至死锁。
- 如果你需要“有限容量 + 优先级”,必须用
asyncio.PriorityQueue(maxsize=N),不要试图用Queue套一层再排序 -
PriorityQueue的qsize()返回的是当前堆长度,但empty()和full()的判断逻辑跟Queue不同——它不检查 maxsize 是否已满,除非你显式传了maxsize - 调试时打印队列内容?别直接 print(queue._queue._queue),那是内部 heap list,顺序不直观;要用
list(queue._queue._queue)并手动 heapify 检查,或者改用queue._queue.queue(Python 3.12+ 已改为_queue属性)
真正麻烦的不是实现,而是优先级语义怎么跟业务对齐:是紧急程度?资源权重?还是 SLA 等级?这些没法靠数据结构解决,得靠上游任务生成逻辑兜底。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










