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 积压,还可能因处理超时导致消息被重复投递;
- 经验起点:轻量任务 50
200,中等任务 1030,重量任务 1~5,再按队列积压与消费延迟实测调整。
同一信道上的多个消费者共享该信道的 prefetch 额度;多线程消费建议每线程一个信道、各自设 prefetch。
小结
QoS 用 prefetch 给每个信道设"未确认消息上限",prefetch=1 实现逐条公平分发,配合手动 ack 就是消费端限流的标准做法;prefetch 太小伤吞吐、太大会倾斜,需按任务轻重实测调优。