当前位置:首页>python>Python实现消费者优先级队列,别让补数据任务堵住线上告警

Python实现消费者优先级队列,别让补数据任务堵住线上告警

  • 2026-09-25 15:37:52
Python实现消费者优先级队列,别让补数据任务堵住线上告警
队列里积压了两千多个任务,消费者也没挂,CPU、内存看着都正常,但一条线上告警等了四分钟还没执行。

日志顺序更离谱:

14:20:11 consume task=history_export priority=8
14:20:13 consume task=history_export priority=8
14:20:15 publish task=payment_alarm priority=1
14:20:17 consume task=history_export priority=8

告警任务已经进队列,消费者还在慢吞吞地导历史数据。

这种问题我一般不先加消费者。加线程只能提高吞吐量,不能解决执行顺序。队列先进先出,前面压着一批低价值任务,后面的紧急任务照样得排队。

这里需要的是优先级队列。

Python 自带的 queue.PriorityQueue 已经把线程安全处理好了,消费者不需要自己围着 heapq 再套一层锁。不过直接往里面塞 (priority, task),我不太建议。

代码很容易写成这样:

task_queue.put((1, {"type": "payment_alarm"}))

问题藏在优先级相同的时候。两个元组的第一个元素相等,Python 会继续比较第二个元素,而字典之间不能比较大小,运行到线上就会看到:

TypeError: '<' not supported between instances of 'dict' and 'dict'

我更习惯额外放一个递增序号。优先级相同时,谁先进队列谁先执行,顺序也不会乱。

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

log = logging.getLogger("task-consumer")
logging.basicConfig(
    level=logging.INFO,
    format="%(asctime)s %(threadName)s %(levelname)s %(message)s",
)

sequence = itertools.count()
task_queue: queue.PriorityQueue["QueueItem"] = queue.PriorityQueue(
    maxsize=5000
)


@dataclass(order=True)
classQueueItem:
    priority: int
    sequence_id: int
    payload: dict[str, Any] = field(compare=False)


defpublish(task_type: str, priority: int, **data: Any) -> None:
    item = QueueItem(
        priority=priority,
        sequence_id=next(sequence),
        payload={
"task_type": task_type,
"data": data,
"created_at": time.time(),
        },
    )

try:
        task_queue.put(item, timeout=1)
except queue.Full:
        log.error(
"queue_full task=%s priority=%s size=%s",
            task_type,
            priority,
            task_queue.qsize(),
        )
raise

这里有两个地方不能省。

一个是 payload 上的 compare=False,否则同优先级任务还是可能比较业务数据。另一个是 maxsize。无界队列写起来省事,出问题时也省事,内存会替你把进程停掉。

消费者拿到任务后,先记录等待时间。我排查队列问题时,这个值比“执行成功”有用得多。

defhandle_task(payload: dict[str, Any]) -> None:
    task_type = payload["task_type"]
    data = payload["data"]

if task_type == "payment_alarm":
        log.warning("send_alarm order_id=%s", data["order_id"])
        time.sleep(0.2)
return

if task_type == "sync_inventory":
        log.info("sync_inventory sku=%s", data["sku"])
        time.sleep(0.8)
return

if task_type == "history_export":
        log.info("export_history batch=%s", data["batch"])
        time.sleep(2)
return

raise ValueError(f"unknown task type: {task_type}")


defconsume() -> None:
whileTrue:
        item = task_queue.get()

try:
            payload = item.payload
            wait_ms = int(
                (time.time() - payload["created_at"]) * 1000
            )

            log.info(
"consume task=%s priority=%s wait_ms=%s queue_size=%s",
                payload["task_type"],
                item.priority,
                wait_ms,
                task_queue.qsize(),
            )
            handle_task(payload)

except Exception:
            log.exception(
"consume_failed task=%s priority=%s",
                item.payload.get("task_type"),
                item.priority,
            )
finally:
            task_queue.task_done()

启动两个消费者,再塞几类任务进去:

for index in range(2):
    threading.Thread(
        target=consume,
        name=f"consumer-{index}",
        daemon=True,
    ).start()

for batch in range(6):
    publish("history_export", priority=8, batch=batch)

publish("sync_inventory", priority=4, sku="SKU-9082")
publish("payment_alarm", priority=1, order_id="ORD-7315")

task_queue.join()

优先级数字越小,越早出队。消费者处理完手上的任务后,会先取告警,再取库存同步,最后继续导出历史数据。

不过这里别理解错了。优先级队列不能中断一个已经执行中的任务。消费者如果正在跑一个耗时两分钟的导出任务,告警进来后仍然要等它执行完。

线上任务耗时差异很大时,我通常会再拆一层:告警任务单独一组消费者,普通任务走另一组。优先级队列负责同一类任务里的先后顺序,不要拿它硬扛所有隔离需求。

还有一种坑是低优先级任务长期饿死。高优先级任务持续进入,历史任务可能一直抢不到执行机会。可以根据等待时间动态提高优先级,或者限制连续处理高优先级任务的数量。

但别一上来就把调度规则写成半个操作系统。先盯住三个值:队列长度、任务等待时间、任务执行时间。哪一个开始持续上涨,再决定是加消费者、拆队列,还是调整优先级。

队列消费最怕的不是慢,而是看起来一直在工作,真正该处理的任务却卡在后面。

最新文章

随机文章