RabbitMQ 消息堆积把下游服务拖垮:我用死信队列 + 消费速率限流 8 分钟止损

上周二凌晨 2 点,监控群里炸了。下游订单服务的 P99 延迟从 80ms 飙到 4.2s,CPU 打到 90%,Pod 开始连环重启。我爬起来看日志,满屏都是 ConnectionTimeoutThreadPoolExecutor rejection

源头不在订单服务本身。上游 RabbitMQ 队列里堆积了 47 万条消息,消费者全速拉取,订单服务被打满了。这玩意儿比直接宕机还恶心——它不崩,只是让你生不如死。

事情是怎么发生的

我们有一个营销活动发券的异步链路:

  • 上游服务批量推送优惠券发放任务到 RabbitMQ
  • 订单服务作为消费者,每条消息要查库存、写流水、调支付网关
  • 正常流量下 TPS 也就 200 左右,撑得住

但那天晚上运营搞了个"限时秒杀"预热,活动开始前 10 分钟,一波 50 万张券同时灌进队列。消费者还是按原来的 prefetch=50、单线程处理,47 万条消息像洪水一样冲进下游。

订单服务每个 Pod 开了 20 个线程处理消息,但每条消息要调 3 个外部接口,平均 RT 200ms。算下来单个 Pod 理论吞吐也就 100 TPS,20 个 Pod 撑死 2000 TPS。消息生产速度是 5000+ TPS,队列开始堆积。

更要命的是,消费者没做限流。RabbitMQ 的 Push 模式默认就是能推多少推多少,消费者 fetch 完就立即 ack,消息源源不断地进来。订单服务的线程池被占满,新来的消息只能排队等线程,等的时间越长,堆积越多,恶性循环。

我当时的想法是:先把消费者停掉,让服务喘口气。结果发现停掉消费者也没用——消息还在堆积,只是换了个地方堆而已。治标不治本。

排查:找出真正的瓶颈

先看一下 RabbitMQ 管理后台的队列状态:

  • messages_ready: 47.3 万
  • messages_unacknowledged: 1000(prefetch 堆积在客户端)
  • message_rate: 入队 5200/s,出队 1800/s

入队比出队快将近 3 倍,不堆积才怪。

再看订单服务:

# 看线程池状态
jstack <pid> | grep "pool-" | wc -l
# 200 个线程,全忙

# 看接口超时
kubectl logs order-service-xxx | grep "Read timed out" | wc -l
# 每分钟 300+ 次

问题很明确:不是 RabbitMQ 本身扛不住,是下游消费能力跟不上了。这时候你要是去加 RabbitMQ 节点、调内存、改磁盘,全是浪费钱。

真正的解法只有两个方向:

  1. 把消息分流——堆积的 47 万条不能全挤在一个出口
  2. 把消费限速——让下游按自己的节奏吃,而不是被硬塞

方案一:死信队列做二级缓冲

死信队列(DLX)不是只用来存异常消息的。我的思路是:把队列拆成"主队列 + 死信队列",当主队列消费者忙不过来时,让消息先转去死信队列,等主队列压力降下来再手动重投。

先改 RabbitMQ 的队列声明,绑定死信交换机:

import pika

connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
channel = connection.channel()

# 死信交换机
channel.exchange_declare(exchange='coupon.dlx', exchange_type='topic')

# 主队列:设置 x-dead-letter-exchange
args = {
    'x-dead-letter-exchange': 'coupon.dlx',
    'x-dead-letter-routing-key': 'coupon.retry',
    'x-max-length': 100000,  # 主队列最大长度,超出的进死信
    'x-message-ttl': 300000  # 5分钟未消费也进死信
}

channel.queue_declare(queue='coupon.main', arguments=args)
channel.queue_bind(queue='coupon.main', exchange='coupon.direct', routing_key='coupon.issue')

# 死信队列
channel.queue_declare(queue='coupon.dlq')
channel.queue_bind(queue='coupon.dlq', exchange='coupon.dlx', routing_key='coupon.retry')

这里用了两个条件:

  • x-max-length: 主队列最多存 10 万条,超过的直接进死信
  • x-message-ttl: 消息 5 分钟没消费完也进死信

这样做的好处是:当流量突增时,主队列不会无限堆积,超出的消息被"暂存"到死信队列。等主队列消费完了,我们再批量把死信队列的消息重投回主队列。

但有个坑要注意:x-max-length 默认是按消息数量算的,如果你的消息体很大,内存还是可能炸。建议改成按字节数限制:

args = {
    'x-max-length': 100000,
    'x-max-length-bytes': 1073741824,  # 1GB
    # ... 其他参数
}

两个限制同时生效,哪个先触发都行。

方案二:消费速率限流,QoS + 动态阈值

死信队列是"泄洪",但真正治本的是让消费者按需拉取,而不是被 Broker 硬塞。

RabbitMQ 的 basic.qos 就是这个作用。默认消费者会尽可能多地预取消息,但你可以通过 prefetch_count 限制每个消费者同时处理的消息数。

channel.basic_qos(prefetch_count=10)

但这只是静态限制。活动期间流量是波动的,平时 200 TPS 够用,活动时可能 5000 TPS。固定 prefetch=10 在高峰期还是不够用,低谷期又浪费。

我写了一个动态限流的脚本,根据下游服务的负载自动调整 prefetch:

import requests
import pika
import time

# 监控下游服务的 CPU 和线程池使用率
def get_service_load():
    try:
        resp = requests.get('http://order-service:8080/actuator/metrics', timeout=2)
        data = resp.json()
        cpu = data.get('cpu.usage', 0)
        threads = data.get('executor.active', 0)
        return cpu, threads
    except:
        return 0, 0

# 动态计算 prefetch
def calculate_prefetch(cpu, active_threads, max_threads=100):
    if cpu > 80:
        return 1  # 服务快炸了,只给一条
    elif cpu > 60:
        return 5
    elif active_threads > max_threads * 0.8:
        return 3
    else:
        return 20  # 正常状态,放开吃

# 消费者循环
def consume_with_dynamic_qos():
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()
    
    # 每 30 秒调整一次 QoS
    last_qos_update = 0
    current_prefetch = 20
    
    def callback(ch, method, properties, body):
        try:
            process_message(body)
            ch.basic_ack(delivery_tag=method.delivery_tag)
        except Exception as e:
            # 处理失败,进入死信队列
            ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
    
    while True:
        now = time.time()
        if now - last_qos_update > 30:
            cpu, threads = get_service_load()
            new_prefetch = calculate_prefetch(cpu, threads)
            if new_prefetch != current_prefetch:
                channel.basic_qos(prefetch_count=new_prefetch)
                current_prefetch = new_prefetch
                print(f"[QoS] CPU={cpu:.1f}%, Threads={threads}, prefetch={new_prefetch}")
            last_qos_update = now
        
        channel.basic_consume(queue='coupon.main', on_message_callback=callback)
        channel.start_consuming()

这个脚本的核心逻辑:

  • CPU > 80%: 只给 1 条消息,让服务先把眼前的处理完
  • CPU 60-80%: 给 5 条,保守消费
  • 线程池占用 > 80%: 给 3 条,防止线程池被打满
  • 正常状态: 给 20 条,正常吞吐

每 30 秒调整一次,不会频繁抖动。实际压测下来,这个策略能把高峰期的消息堆积时间从原来的"无限堆积"压到 3-5 分钟以内。

死信重投:别让人工来搬消息

消息进死信队列后,不能让它烂在那里。我写了一个定时任务,监控主队列长度,低于阈值时自动把死信消息搬回主队列:

import pika
import json
from datetime import datetime

def requeue_dead_letters():
    connection = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
    channel = connection.channel()
    
    # 检查主队列长度
    queue_info = channel.queue_declare(queue='coupon.main', passive=True)
    main_count = queue_info.method.message_count
    
    if main_count > 50000:
        print(f"[Requeue] Main queue still heavy ({main_count}), skip requeue")
        return 0
    
    # 批量搬回死信消息
    requeued = 0
    batch_size = 1000
    
    for _ in range(batch_size):
        method, properties, body = channel.basic_get(queue='coupon.dlq', auto_ack=False)
        if not method:
            break
        
        # 去掉死信标记,重新发回主队列
        channel.basic_publish(
            exchange='coupon.direct',
            routing_key='coupon.issue',
            body=body,
            properties=pika.BasicProperties(
                delivery_mode=2,  # persistent
                headers={'x-requeued-at': datetime.now().isoformat()}
            )
        )
        channel.basic_ack(delivery_tag=method.delivery_tag)
        requeued += 1
    
    print(f"[Requeue] Moved {requeued} messages from DLQ to main queue")
    return requeued

if __name__ == '__main__':
    requeue_dead_letters()

这个脚本放在 cron 里每 2 分钟跑一次,或者塞进 K8s CronJob。关键逻辑:

  • 主队列消息数 > 5 万时不搬,避免二次冲击
  • 每次最多搬 1000 条,细水长流
  • 搬的时候给消息加 x-requeued-at 头,方便后续排查哪些消息被重投过

效果:8 分钟止血

把这套方案上线后,效果是这样的:

  • 02:14 发现堆积,停止原来的消费者,改成动态 QoS 消费者
  • 02:16 主队列长度从 47 万降到 10 万(超出的 37 万进了死信队列)
  • 02:18 订单服务 CPU 从 90% 降到 45%,线程池恢复正常
  • 02:22 死信队列开始自动重投,主队列维持在 2-5 万条
  • 02:30 所有消息消费完毕,服务恢复

总耗时 8 分钟。如果没有死信队列做二级缓冲,单靠消费者限速,47 万条消息按 2000 TPS 的消费能力要消化将近 4 分钟,但期间新消息还在不断进来,实际时间会更长。而且下游服务一直高负载,随时可能真的 OOM 或者连接池耗尽。

死信队列的作用就是"把洪峰切平"——让主队列只处理"当下能处理得动"的量,剩下的先存起来,等压力降下来再慢慢消化。

监控告警:别等 47 万条才报警

事后我补了三个监控指标:

# 1. 队列堆积告警(5 分钟超过 10 万条)
rabbitmq_queue_messages{queue="coupon.main"} > 100000

# 2. 消息消费速率告警(入队 > 出队 2 倍持续 3 分钟)
rate(rabbitmq_queue_messages_published_total[5m]) > 2 * rate(rabbitmq_queue_messages_delivered_total[5m])

# 3. 死信队列堆积告警(死信队列超过 5 万条)
rabbitmq_queue_messages{queue="coupon.dlq"} > 50000

第一个告警让我能在堆积到 10 万条时就收到通知,而不是等 47 万条了才发现。第二个告警直接告诉我"生产比消费快",也就是消费能力不足的早期信号。第三个告警防止死信队列本身变成新的炸弹。

写在最后

RabbitMQ 的坑,很多时候不是 Broker 本身扛不住,而是消费者和 Broker 之间没有"缓冲带"。

Kafka 有消费者组的背压机制,RocketMQ 有消费速率调整,但 RabbitMQ 的 AMQP 协议天然是 Push 模式,消费者如果不主动限速,就只能被动接受。死信队列在这里不是"异常消息的垃圾桶",而是"流量过载时的临时仓库"。

几个建议:

  1. 所有队列都该绑定死信交换机。不只是为了存异常消息,更是为了泄洪。
  2. QoS 不要写死。根据下游负载动态调整,让服务自己决定能吃多少。
  3. 监控入队和出队的比值。这个指标比单纯的队列长度更敏感,能提前 3-5 分钟发现问题。

如果你也在用 RabbitMQ,建议现在就检查一下自己的消费者有没有做 prefetch 限制。没有这个,等于开车不系安全带——出事只是时间问题。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐