page contents

Python 实现消费者优先级队列,别直接把字典塞进去

Python 的 PriorityQueue 底层用的是小顶堆。它先比较元组第一个元素,也就是优先级;优先级相同时,还会继续比较第二个元素。两个字典没法比较大小,异常就出来了。

attachments-2026-07-82nZUl5S6a6aac07d5900.pngTypeError: '<' not supported between instances of 'dict' and 'dict'

这个报错一出来,我基本不用往业务代码里翻,先看优先级队列里放了什么。

十有八九,代码写成了这样:

task_queue.put((1, {"order_id": "A1024"}))
task_queue.put((1, {"order_id": "A1025"}))

第一条任务没事,第二条任务一进去就可能报错。

Python 的 PriorityQueue 底层用的是小顶堆。它先比较元组第一个元素,也就是优先级;优先级相同时,还会继续比较第二个元素。两个字典没法比较大小,异常就出来了。

这类问题开发环境不一定撞得到。任务少、优先级刚好不同,看着一直正常。等线上同一批高优任务连续进来,消费者还没开始处理,生产者先挂了。

我一般不会让业务对象直接参与队列排序,而是额外放一个递增序号。

import itertools
import logging
import queue
import threading
import time
from dataclasses import dataclass, field
from typing import Any

logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(threadName)s %(levelname)s %(message)s",
)

_sequence = itertools.count()


@dataclass(order=True)
class ConsumeTask:
    priority: int
    sequence: int = field(init=False)
    topic: str = field(compare=False)
    payload: Any = field(compare=False)

    def __post_init__(self) -> None:
        self.sequence = next(_sequence)

这里有两个细节。

priority 越小,任务越先出队。紧急任务可以用 0,普通任务用 10,补偿任务用 20。

sequence 用来保证同优先级任务按照进入队列的顺序执行。topic 和 payload 都设置成 compare=False,业务数据是什么类型,不影响堆排序。

消费者代码不要写成无限阻塞后什么都不管。任务执行失败、队列退出、task_done(),这些地方少一个,后面排查都很烦。

class PriorityConsumer:
    def __init__(self, worker_count: int = 3, capacity: int = 500):
        self.tasks: queue.PriorityQueue[ConsumeTask] = queue.PriorityQueue(
            maxsize=capacity
        )
        self.worker_count = worker_count
        self.closed = threading.Event()
        self.workers: list[threading.Thread] = []

    def start(self) -> None:
        for index in range(self.worker_count):
            worker = threading.Thread(
                target=self._consume,
                name=f"priority-consumer-{index}",
                daemon=True,
            )
            worker.start()
            self.workers.append(worker)

    def submit(self, priority: int, topic: str, payload: Any) -> bool:
        task = ConsumeTask(
            priority=priority,
            topic=topic,
            payload=payload,
        )

        try:
            self.tasks.put(task, timeout=0.2)
            logging.info(
                "task queued topic=%s priority=%s backlog=%s",
                topic,
                priority,
                self.tasks.qsize(),
            )
            return True
        except queue.Full:
            logging.warning(
                "task rejected topic=%s priority=%s backlog=%s",
                topic,
                priority,
                self.tasks.qsize(),
            )
            return False

    def _consume(self) -> None:
        while not self.closed.is_set() or not self.tasks.empty():
            try:
                task = self.tasks.get(timeout=0.5)
            except queue.Empty:
                continue

            try:
                self._handle(task)
            except Exception:
                logging.exception(
                    "task failed topic=%s priority=%s",
                    task.topic,
                    task.priority,
                )
            finally:
                self.tasks.task_done()

    @staticmethod
    def _handle(task: ConsumeTask) -> None:
        logging.info(
            "task consuming topic=%s priority=%s payload=%s",
            task.topic,
            task.priority,
            task.payload,
        )
        time.sleep(0.2)

    def stop(self) -> None:
        self.tasks.join()
        self.closed.set()

        for worker in self.workers:
            worker.join(timeout=1)

跑一组任务看看顺序:

consumer = PriorityConsumer(worker_count=1)
consumer.start()

consumer.submit(10, "sync_profile", {"user_id": 81})
consumer.submit(0, "freeze_account", {"user_id": 19})
consumer.submit(5, "send_notice", {"user_id": 33})
consumer.submit(0, "cancel_payment", {"order_id": "P908"})

consumer.stop()

单消费者情况下,两条优先级为 0 的任务会先执行,并且保持提交顺序。消费者数量一旦大于一个,只能保证“出队顺序”,不能保证“完成顺序”。前一个任务调用外部接口卡了两秒,后一个任务只查了一次缓存,后提交的照样可能先完成。

这个区别经常被忽略。

还有一个地方我会盯着看:队列长度。代码里用了 maxsize,不是为了省那点内存,而是给生产速度加一道闸。队列无上限时,下游一慢,任务就会不断堆在进程内存里。最后看到的未必是业务超时,可能直接变成进程被系统杀掉。

日志至少要留下这几个字段:

topic=freeze_account priority=0 backlog=37

没有 backlog,消费者变慢时只能猜;没有 priority,普通任务为什么十分钟没处理,也不好判断。

优先级队列还有一个天然问题:低优先级任务可能一直饿着。只要高优先级任务持续进入,后面的补偿任务就永远轮不到。业务上真有这种流量,我更愿意拆成“紧急队列”和“普通队列”,消费者按比例拉取,比如处理五个紧急任务后,强制取一个普通任务。

别指望一个整数优先级把所有调度问题都解决。

PriorityQueue 适合单进程里的任务调度、批量文件处理、接口补偿和日志分级消费。跨进程、需要持久化、要求消息确认的场景,就别继续给这段代码打补丁了,该上消息队列还是上消息队列。

进程内队列最怕的不是代码长,而是看着简单,边界一个没处理。

更多相关技术内容咨询欢迎前往并持续关注好学星城论坛了解详情。

想高效系统的学习Python编程语言,推荐大家关注一个微信公众号:Python编程学习圈。每天分享行业资讯、技术干货供大家阅读,关注即可免费领取整套Python入门到进阶的学习资料以及教程,感兴趣的小伙伴赶紧行动起来吧。

attachments-2022-05-rLS4AIF8628ee5f3b7e12.jpg

 

  • 发表于 2026-07-30 09:42
  • 阅读 ( 39 )
  • 分类:Python开发

你可能感兴趣的文章

相关问题

0 条评论

请先 登录 后评论
Pack
Pack

2315 篇文章

作家榜 »

  1. 轩辕小不懂 2403 文章
  2. Pack 2315 文章
  3. 小柒 2228 文章
  4. Nen 576 文章
  5. 王昭君 209 文章
  6. 文双 71 文章
  7. 小威 64 文章
  8. Cara 36 文章