因为queue.queue的priorityqueue在优先级相同时会比较任务对象,而自定义对象若未实现__lt__方法将抛typeerror。

为什么不能直接用 queue.Queue 做优先级任务调度?
queue.Queue 确实支持优先级(靠 PriorityQueue 子类),但它默认用元组比较:如果优先级相同,会尝试比较任务本身——而自定义对象通常没实现 __lt__,直接抛 TypeError: '。更麻烦的是,它不支持动态调整任务优先级,也没法去重或取消待处理任务。
用 heapq 手写可取消的优先队列
核心思路是用 heapq 维护堆,配合一个 dict 记录任务状态,实现“逻辑删除”。这样能避免重复执行、支持取消、还能在插入时覆盖同 ID 任务:
import heapq
import threading
<p>class PriorityQueue:
def <strong>init</strong>(self):
self._heap = []
self._entry_map = {} # task_id → [priority, task_id, task]
self._counter = 0
self._lock = threading.Lock()</p><pre class="brush:python;toolbar:false;">def put(self, task_id, priority, task):
with self._lock:
# 已存在则标记为失效,再推新条目
if task_id in self._entry_map:
self._entry_map[task_id][-1] = None # 标记旧任务失效
entry = [priority, self._counter, task_id, task]
self._entry_map[task_id] = entry
heapq.heappush(self._heap, entry)
self._counter += 1
def get(self):
while self._heap:
priority, count, task_id, task = heapq.heappop(self._heap)
if task is not None: # 跳过被取消的
del self._entry_map[task_id]
return task_id, task
raise IndexError("pop from empty priority queue")
关键点:
-
heapq只保证堆顶最小,不保证整体有序;依赖[priority, count, ...]结构避免比较 task 对象 - 取消不是真删,而是置
task = None,get()时跳过——这是性能和线程安全的折中 -
threading.Lock必须包裹所有涉及_heap和_entry_map的操作,否则多线程下堆结构可能损坏
如何让任务支持延迟执行(比如定时调度)?
单纯优先级不够,很多场景需要“未来某时刻才可执行”。这时得把触发时间也纳入排序依据:
把堆元素改成 [scheduled_time, priority, count, task_id, task],其中 scheduled_time 是绝对时间戳(如 time.time() + delay)。消费者线程需持续检查堆顶:if scheduled_time ,否则 <code>time.sleep(min(0.1, scheduled_time - time.time())) 避免忙等。
注意:
- 不要用
datetime对象做堆元素——它不可比较(除非都带 tzinfo);统一用 float 秒级时间戳 - 如果任务允许“过期即丢弃”,在
get()中加判断:if scheduled_time > time.time(): continue - 修改堆内元素(如提前触发)会导致堆损坏,只能取消后重插
生产环境别绕开 asyncio 或专用库
手写队列适合学习或极简嵌入场景。真实服务中,优先考虑:
- 异步场景用
asyncio.PriorityQueue(注意它也不支持取消,需自己 wrap) - 需要持久化、重试、监控,直接上
celery+redis或rq - 轻量但需可靠,用
arq(基于 redis 的 async 任务队列)
手写代码最难缠的不是逻辑,是边界:时钟跳变导致延迟不准、高并发下锁竞争、OOM 前没及时清理失效条目——这些在成熟库里都有对应防护。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











