RabbitMQ 消费限流与 QoS

默认情况下 broker 会把消息尽可能快地推给消费者,消费者一旦处理不过来就会积压。QoS(服务质量)与 prefetch 参数用于限制"同一时刻一个消费者手上有多少未确认消息",是消费端最重要的限流手段。

prefetch 是什么

prefetch 定义信道(channel)级别未确认消息的数量上限:消费者手上未 ack 的消息数达到该值时,broker 不再推送新消息,直到消费者确认掉一部分;值为 0 表示不限制。

ch.basic_qos(prefetch_count=1)   # 每次最多 1 条未确认
ch.basic_qos(prefetch_count=10)  # 每次最多 10 条未确认

轮询与公平分发

不设 prefetch 时,broker 按轮询把消息轮流发给各个消费者,不分处理快慢:快消费者干完活等新消息,慢消费者却越积越多,整体吞吐被最慢者拖累。

设 prefetch=1 后变为公平分发:一次只给消费者一条,处理完(ack)才给下一条,快消费者自然多拿、慢消费者少拿,实现"能者多劳"。

配置分发方式效果
不设 QoS轮询快慢不分,易倾斜
prefetch=1公平分发逐条处理,天然限流
prefetch=N批量流水吞吐与内存的折中

配合手动 ack 的限流示例

限流要生效必须用手动确认,否则消息发出去即删除,prefetch 就失去意义:

import pika
conn = pika.BlockingConnection(pika.ConnectionParameters("127.0.0.1"))
ch = conn.channel()
ch.basic_qos(prefetch_count=1)     # 信道级未确认上限=1
def on_msg(ch, method, props, body):
    ok = handle(body)              # 模拟处理(如调第三方接口)
    if ok:
        ch.basic_ack(delivery_tag=method.delivery_tag)   # 处理完才放行下一条
    else:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)
ch.basic_consume("q.task", on_message_callback=on_msg, auto_ack=False)

处理耗时长、下游有并发压力的任务(发短信、调外部 API)非常适合 prefetch=1:天然保证同一时刻只有一条在途,避免把下游打挂。

处理很快的轻量任务可适当放大 prefetch(如 30~100),摊薄每条消息的网络往返,提升吞吐。

prefetch 调优建议

  • prefetch 过小(始终 =1)而消息很轻时,每次都要等 ack 往返,吞吐上不去,跨机房高延迟场景尤其明显;
  • prefetch 过大(=0 不限制)时,broker 瞬间把大量消息塞进消费者内存,慢消费者处理不过来,unacked 积压,还可能因处理超时导致消息被重复投递;
  • 经验起点:轻量任务 50200,中等任务 1030,重量任务 1~5,再按队列积压与消费延迟实测调整。

同一信道上的多个消费者共享该信道的 prefetch 额度;多线程消费建议每线程一个信道、各自设 prefetch。

小结

QoS 用 prefetch 给每个信道设"未确认消息上限",prefetch=1 实现逐条公平分发,配合手动 ack 就是消费端限流的标准做法;prefetch 太小伤吞吐、太大会倾斜,需按任务轻重实测调优。

笔记加载中…