队列里积压了两千多个任务,消费者也没挂,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()
优先级数字越小,越早出队。消费者处理完手上的任务后,会先取告警,再取库存同步,最后继续导出历史数据。
不过这里别理解错了。优先级队列不能中断一个已经执行中的任务。消费者如果正在跑一个耗时两分钟的导出任务,告警进来后仍然要等它执行完。
线上任务耗时差异很大时,我通常会再拆一层:告警任务单独一组消费者,普通任务走另一组。优先级队列负责同一类任务里的先后顺序,不要拿它硬扛所有隔离需求。
还有一种坑是低优先级任务长期饿死。高优先级任务持续进入,历史任务可能一直抢不到执行机会。可以根据等待时间动态提高优先级,或者限制连续处理高优先级任务的数量。
但别一上来就把调度规则写成半个操作系统。先盯住三个值:队列长度、任务等待时间、任务执行时间。哪一个开始持续上涨,再决定是加消费者、拆队列,还是调整优先级。
队列消费最怕的不是慢,而是看起来一直在工作,真正该处理的任务却卡在后面。