Skip to content

消息队列

消息队列的价值可以概括成四个词:解耦、异步、削峰、可靠投递;代价也有三个:一致性变弱、链路变长、排障变难。本文先讲清楚什么时候该用 MQ,再对比 Kafka 与 RabbitMQ 两套架构范式,最后用实测说明投递语义、幂等、顺序性与死信这些"用了才知道"的关键细节。

一、什么时候该用,什么时候不该用

该用不该用
上下游处理能力不匹配,需要削峰强一致要求(同一事务内的操作别拆)
一次操作要触发多个下游,且允许异步(下单→发券→通知→埋点)需要立刻拿到结果(登录校验、库存强校验)
需要重试与补偿(第三方调用失败后重投)数据量极小、调用简单(直接 RPC 更简单)
需要广播(一个事件多个消费方)团队没有 MQ 运维能力(先评估代价)
需要解耦发布节奏(数据变更广播给多个系统)只是"想让接口快一点"——先查 SQL 与缓存

常见误用:把同步调用强行改成异步,却不处理失败重试与幂等——结果是"接口快了,但订单丢了还查不到原因"。引入 MQ 前先想清楚:失败了怎么办?重复了怎么办?

二、两种模型

模型语义代表
队列(点对点)一条消息只被一个消费者处理,消费后出队RabbitMQ Queue、RocketMQ、ActiveMQ
发布/订阅一条消息被所有订阅者各处理一次Kafka Consumer Group、RabbitMQ fanout

核心差异(选型时最重要的分野)

维度RabbitMQKafka
定位消息代理(Broker 负责路由、确认、重投)分布式日志(消息是持久化的 append-only 日志)
消费后是否删除确认后删除不删除,按保留期(默认 7 天)保留,可重放
消费模型Push(Broker 推给消费者)Pull(消费者按 offset 拉)
顺序保证单队列内有序分区内有序
吞吐万~十万级/秒十万~百万级/秒(顺序磁盘 IO + 零拷贝)
延迟微秒~毫秒毫秒级(批量换取吞吐)
典型场景任务队列、RPC 异步化、复杂路由日志、埋点、流处理、事件溯源

三、Kafka 架构要点

text
Topic(逻辑主题)
 └─ Partition 0  [msg0][msg1][msg2][msg3]…   ← 每个分区是一个有序、不可变的日志
 └─ Partition 1  [msg0][msg1][msg2]…
 └─ Partition 2  [msg0][msg1][msg2][msg3][msg4]…

Producer → 按 key 哈希(或轮询)决定写入哪个分区
Consumer Group → 组内每个消费者负责若干分区(一个分区只被组内一个消费者消费)
Offset → 每条消息在分区内的序号,由消费者自己维护(可回退、可重放)

关键机制:

  • 分区是并行度与顺序性的单位:分区越多吞吐越高,但只能保证分区内有序;
  • Consumer Group 实现"队列 + 广播"两种语义:组内竞争消费(队列),组间各自独立消费(广播);
  • Rebalance:消费者增减时重新分配分区,期间会短暂停止消费(避免在高峰期频繁扩缩容);
  • 副本与 ISR:每个分区有多个副本,只有 ISR(In-Sync Replicas)集合内的副本才可选为 leader;acks=all + min.insync.replicas=2 才能在不丢数据和不牺牲可用性之间取得平衡;
  • 顺序性只在分区内成立:要按用户保序,就用 userId 作为消息 key,让同一用户的消息都进同一分区(实测见 §六)。
java
// 生产者:至少一次语义的典型配置
props.put("acks", "all");                    // 等所有 ISR 副本确认
props.put("retries", 3);
props.put("enable.idempotence", true);       // 幂等生产者:避免重试产生重复消息
props.put("max.in.flight.requests.per.connection", 5);

// 消费者:手动提交 offset(处理成功后再提交)
while (true) {
  ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
  for (ConsumerRecord<String, String> r : records) {
    process(r);                              // 业务处理(必须幂等)
  }
  consumer.commitSync();                     // 处理完成后提交,避免"未处理就提交"导致丢消息
}

四、RabbitMQ 的交换机模型

生产者不直接发到队列,而是发到 Exchange,由 Exchange 按类型和绑定规则路由:

类型路由规则场景
Directrouting key 完全匹配点对点任务分发
Topic通配符匹配(order.*#.paid按事件类型分发,最灵活
Fanout忽略 key,广播到所有绑定队列广播通知
Headers按消息头属性匹配少用(复杂且慢)
python
# 生产者:发送到 topic 交换机,消息持久化
channel.exchange_declare(exchange='order.events', exchange_type='topic', durable=True)
channel.basic_publish(
    exchange='order.events',
    routing_key='order.paid',                 # 消费者可用 order.* 订阅
    body=json.dumps(payload),
    properties=pika.BasicProperties(delivery_mode=2)   # 2 = 持久化
)

# 消费者:手动 ACK,prefetch 控制并发
channel.basic_qos(prefetch_count=10)          # 一次最多推 10 条未确认消息
def callback(ch, method, properties, body):
    try:
        handle(body)
        ch.basic_ack(delivery_tag=method.delivery_tag)     # 成功后确认
    except Exception:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=True)  # 失败重回队列
channel.basic_consume(queue='pay.queue', on_message_callback=callback)

prefetch_count 很重要:不设的话 Broker 会把队列里所有消息一次性推给消费者,导致内存爆掉与负载不均。

五、可靠性与投递语义

一条消息要经过三个环节,每个环节都可能丢:

环节风险保障手段
生产 → Broker网络失败、Broker 未落盘发送确认(acks=all / publisher confirm)+ 失败重试 + 本地消息表
Broker 存储宕机丢失持久化(Kafka 落盘 + 多副本;RabbitMQ 持久化队列 + 消息 delivery_mode=2
Broker → 消费消费失败但已确认手动 ACK:业务处理成功后再提交(Kafka offset / RabbitMQ ack)

三种投递语义:

语义含义实现
最多一次 At-most-once可能丢,绝不重复收到就确认,失败不管(日志、埋点可接受)
至少一次 At-least-once绝不丢,可能重复处理成功再确认,失败重投 —— 主流选择,必须配套幂等
恰好一次 Exactly-once不丢不重Kafka 事务 + 幂等生产者(限 Kafka 内部流转;跨数据库仍需幂等兜底)

六、幂等消费:至少一次的前提

实测(100 条消息,其中 30 条被重复投递,模拟 at-least-once):

text
无幂等:投递 130 次,扣款 1300 元,余额 8700        ← 多扣了 300 元
有幂等:投递 130 次,实际处理 100 次(拦截重复 30 条),扣款 1000 元,余额 9000
js
// 幂等键去重:用消息里的业务唯一 ID(订单号 / 消息 ID)
const processed = new Set()

async function handle(msg) {
  if (processed.has(msg.msgId)) return 'skipped'      // 已处理过,直接跳过
  processed.add(msg.msgId)
  await deductBalance(msg.orderId, msg.amount)
  return 'done'
}

幂等落地的三种做法:

做法实现适用
唯一键去重表用业务唯一号建唯一索引,插入成功才执行最通用,推荐
数据库唯一约束INSERT ... ON DUPLICATE KEY UPDATE / upsert写入类操作
状态机WHERE status = 'INIT' 才更新为 PAID有明确业务状态流转时最优雅

幂等键必须是业务唯一的(订单号、流水号),不要用随机生成的消息 ID——消息重投时 ID 不变才叫重投,变了就是新消息。去重记录要有过期清理,否则表无限膨胀。

七、顺序性:只在分区内成立

实测(12 条事件按 userId 哈希分到 3 个分区,各分区并行消费):

text
分区 0 内顺序:op3 → op6 → op9 → op12
分区 1 内顺序:op1 → op4 → op7 → op10
分区 2 内顺序:op2 → op5 → op8 → op11

实际消费(跨分区交错):P2:u2-op2 , P0:u0-op3 , P1:u1-op1 , P0:u0-op6 , P0:u0-op9 , P2:u2-op5 ...

同一用户 u0 的 4 条事件落在分区 {0} → 分区内保序成立

要点:

  • 全局顺序 = 只有 1 个分区(牺牲并行度),几乎所有场景都不需要
  • 需要保序的维度(用户、订单、商品)用同一个 key,保证进同一分区;
  • 消费者内不要并发处理同一分区的消息(多线程池会打乱顺序),要么单线程按序处理,要么按 key 分组并发(同 key 串行、不同 key 并行)。

八、重试、退避与死信队列

实测(20 条消息,第 7 条永久失败,最多重试 3 次并指数退避):

text
20 条消息:总投递 22 次(含重试)
最终失败进入 DLQ:1 条 → msg#7
其余 19 条成功;msg#7 重试 3 次后进入 DLQ,未阻塞其他消息
js
const MAX_RETRY = 3

async function consumeWithRetry(msg) {
  for (let attempt = 1; attempt <= MAX_RETRY; attempt++) {
    try {
      await handle(msg)
      return ack(msg)
    } catch (err) {
      if (attempt === MAX_RETRY) return sendToDLQ(msg, err)   // 重试耗尽 → 死信队列
      await sleep(BASE_DELAY * 2 ** (attempt - 1))            // 指数退避
    }
  }
}

要点:

  • 必须限制重试次数:永久失败的消息如果无限重试,会堵住整个队列("毒丸消息"问题);
  • 必须有退避:立即重试会把故障放大(见 微服务 的熔断与重试预算);
  • 必须有死信队列(DLQ):失败消息转存到 DLQ,保留上下文供排查与人工重放;
  • 区分可重试与不可重试错误:参数错误(不可重试)应直接进 DLQ,网络超时(可重试)才重试。

九、削峰:用延迟换成功率

实测(突发 1000 个请求,下游处理能力固定 50/s):

text
无队列:成功 108 条,失败(被拒绝)892 条,耗时 1175ms
有队列:成功 1000 条,失败 0 条,耗时 20145ms(队列峰值缓冲 1000 条)
差异:无队列丢 892 条请求;有队列零丢失,代价是耗时从 1175ms 拉长到 20145ms

这就是削峰的本质:用时间换成功率。代价是延迟上升,因此不适合"必须立刻返回结果"的链路。队列长度要有监控——队列持续堆积说明消费者能力不足,需要扩容而不是继续等

十、选型对比

产品吞吐延迟顺序事务消息典型场景
Kafka极高毫秒分区内支持日志、埋点、流处理、事件总线
RabbitMQ低(微秒~毫秒)队列内插件支持任务队列、业务异步、复杂路由
RocketMQ队列内原生支持电商交易、订单链路(阿里系)
Pulsar分区内支持计算存储分离、多租户
Redis Stream极低分区内不支持轻量场景,已有 Redis 时

简单决策:业务消息(需要确认、重试、路由)选 RabbitMQ/RocketMQ;数据管道(高吞吐、可重放、流处理)选 Kafka。已有 Redis 且量不大时,Redis Stream 可以省掉一个组件,但要接受它持久化与运维能力的局限(见 缓存)。

十一、检查清单

  1. 引入 MQ 前已想清楚:失败重试策略、幂等键、堆积告警阈值;
  2. 生产者开启了发送确认,失败有重投或落本地消息表;
  3. Broker 端开启持久化与多副本(acks=all / 持久化队列);
  4. 消费者手动确认,且确认发生在业务处理成功之后;
  5. 消费逻辑幂等,幂等键是业务唯一号;
  6. 需要保序的维度用同一个 key 分区,消费者不对同 key 并发处理;
  7. 重试有次数上限、指数退避与死信队列;
  8. 队列长度与消费延迟有监控告警,堆积时有扩容预案;
  9. 消费者设置了并发上限(Kafka 分区数 / RabbitMQ prefetch);
  10. 消息有 schema 约定与版本兼容策略(加字段不改语义);
  11. 有消息轨迹/TraceID 便于排障(见 微服务 的可观测性);
  12. 有重放能力(Kafka 可按 offset 回退,RabbitMQ 可从 DLQ 重投)。

十二、常见坑速查

  1. 自动确认(auto-ack):消息刚推过来就确认,业务失败即丢失——必须手动确认;
  2. 没有幂等:at-least-once 必然重复投递,不幂等就重复扣款(实测多扣 300 元);
  3. 幂等键用随机 ID:重投时 ID 变了,去重失效——必须用业务唯一号;
  4. 无限重试:一条毒丸消息堵死整个队列——限制次数 + 退避 + DLQ;
  5. 不区分错误类型:参数错误也重试,白费资源——不可重试错误直接进 DLQ;
  6. 假设全局有序:只在分区/队列内有序,跨分区必然交错(实测见 §七);
  7. 分区数不够或过多:不够则并行度上不去,过多则 Rebalance 变慢、文件句柄暴涨;
  8. 消费者无并发上限:RabbitMQ 不设 prefetch 会把消息全推过来,内存爆掉;
  9. 消息体过大:大消息拖慢整个 Broker 并放大网络压力——只传 ID 或摘要,大字段走对象存储;
  10. 用 MQ 做同步 RPC:需要立刻拿到结果的链路不该异步化;
  11. 队列堆积不处理:只在监控里看着涨,不扩容也不降级——堆积是故障前兆;
  12. 消息无版本/schema 管理:上游加字段导致下游反序列化失败,全量回滚;
  13. 事务消息被误用:以为发了事务消息就"既不丢也不重",其实仍需消费端幂等兜底。

状态与参考

下一步

  • [ ] 用 Docker 起一个 Kafka(或 Redpanda)单节点,实测分区数与吞吐的关系、Consumer Group Rebalance 行为
  • [ ] 用 RabbitMQ 跑通 topic 交换机 + 手动 ACK + DLQ 的完整链路,补充真实管理界面截图
  • [ ] 写一篇"本地消息表 vs RocketMQ 事务消息"的分布式事务对比实验

写作规范请参阅领域概览

基于 VitePress 构建 · 内容以知识共享方式沉淀