Skip to content

不丢、不重、不积压:MQ 生产手册

本文是消息队列系统学习系列的 L3 实战篇。前置:[31. MQ 选型与容量规划]。 学完可以配合面试题食用:04-kafka-message-loss15-message-backlog-kafka-rabbitmq-rocketmq24-dead-letter-queue-message-replay

什么才算"不丢"——三段论

消息从生产端出发,经过 Broker 中转,最终被消费端处理。这中间每一段都可能丢,而且丢的原因各不相同。生产环境里最怕的不是丢,是"你不知道丢了"。

三段论把消息生命周期切成三个独立区域,各自防御:

  • 生产端 — Broker:生产者发出消息,Broker 确认接收。如果确认没回到客户端,消息可能根本没到。
  • Broker 存储:消息落盘后被调度到副本。单点宕机、磁盘坏道、刷盘策略不当都会丢。
  • Broker — 消费端:消费端拿到消息,处理一半挂了,offset 已经提交 — 消息"已消费"但实际没处理。

下面逐段拆解,每段给一句"保底原则"。

生产端确认:acks 与重试

生产者发消息,Broker 返回 ack 才算成功。Kafka 的参数 acks 控制确认力度:

  • acks=0:发出去就不管,不确认。磕一下网线就丢。
  • acks=1:Leader 写完本地日志就返回。Leader 挂了但数据还没同步到副本,返回 ack 之后 Leader 宕机,消息就丢了。
  • acks=all(或 -1):Leader 等所有 ISR 副本都写完了才返回 ack。这是最安全的。

重试不等同于可靠。retriesretry.backoff.ms 控制重试行为,但重试针对的是"可恢复错误"(如 Leader 选举中的 NOT_LEADER_FOR_PARTITION),不是"Broker 挂了"——后者要等超时抛异常,业务方自己兜底。

保底原则:生产端不能只靠 MQ 的 ack,还要给自己的数据库加一条"待确认"记录。落数据库、发 MQ、等回调、更新状态——这就是知名"本地消息表"。

Broker 存储:副本数与刷盘

Kafka 关闭了同步刷盘(flush.messages 默认无限大,flush.ms 默认也无限大),靠副本冗余来抗丢失。replication.factor 至少 3,min.insync.replicas 至少 2,才能保证一个副本挂了还能正常服务。

RocketMQ 默认 flushDiskType=ASYNC_FLUSH,核心场景改成 SYNC_FLUSH 会掉 50% 写入性能,但无需等副本同步就确认了安全性。权衡:SYNC_FLUSH + 异步副本 vs ASYNC_FLUSH + 同步副本,大多数业务选前者。

RabbitMQ 的 publisher-confirm + delivery-mode=2(持久化队列)才是稳妥组合。

保底原则:副本数 3 起步,核心队列打开持久化。不要相信"重启就好了"。

消费端:先处理,后 ack

消费端丢消息最常见的原因:先提交了 offset,业务处理还没跑完,进程挂了。

Kafka Consumer 的 enable.auto.commit=true 配合 auto.commit.interval.ms=5000,每 5 秒自动提交一批。如果消费处理耗时 > 5 秒,消息可能在处理到一半时就被标记为"已消费"。

java
// 正确做法:手动提交
Properties props = new Properties();
props.put("enable.auto.commit", "false");
// ...
while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
    for (ConsumerRecord<String, String> record : records) {
        process(record);  // 先处理业务逻辑
    }
    consumer.commitSync();  // 全部处理完再提交
}

RocketMQ 的 ConsumeOrderly 模式默认是"先处理再 ack",但 ConsumeConcurrently 模式如果 ConsumeConcurrentlyStatus.CONSUME_SUCCESS 返回太快,同样有丢失窗口。

保底原则:手动提交 offset,业务处理放在 commit 之前。消费失败时重试,重试耗尽进死信,不要默默跳过。

幂等消费:不重的四种武器

消息系统的"至多一次"和"至少一次"两个语义只保证不丢,不保证不重。Kafka 的 enable.idempotence=true 解决的是生产者到 Broker 的重复(重试导致同一条消息被 Broker 接受两次),不是消费端重复。

消费端重复的原因:消费端处理成功,offset 提交超时,Broker 认为消费失败,重新推送。所以消费端必须自己扛幂等。

四种落地方式,按推荐程度排序:

1. 唯一键去重表

每来一条消息,提取业务唯一键(如订单号),INSERT 到一张幂等表中。INSERT 成功才处理,失败(唯一键冲突)说明已处理过,直接 ack 并跳过。

sql
CREATE TABLE idempotent_record (
    idempotent_key VARCHAR(128) PRIMARY KEY,
    create_time DATETIME DEFAULT CURRENT_TIMESTAMP
);
java
public boolean tryProcess(String orderId, ConsumerRecord<String, String> record) {
    // 幂等检查
    try {
        jdbcTemplate.update("INSERT INTO idempotent_record (idempotent_key) VALUES (?)", orderId);
    } catch (DuplicateKeyException e) {
        log.info("重复消息,跳过: {}", orderId);
        consumer.commitSync();
        return true;
    }
    // 执行业务逻辑...
    return true;
}

2. 状态机幂等

业务本身有状态流转(订单:待支付→已支付→已发货),消费消息做"如果当前状态 ≤ 目标状态才推进"的判断。重复消息到达时状态已在目标状态之后,直接跳过。

3. Redis setnx

短窗口幂等。用 SET order:processed:12345 value NX EX 3600 做防重标记。适合时间窗口内(如 1 小时)不重复即可的场景。注意 Redis 挂了幂等就失效了。

4. 分布式锁

全局锁 + 业务处理,锁释放前别的重复消息拿不到。粒度粗、性能差,适合对账补偿这种低频场景。

选型建议:数据库唯一键最稳,状态机最适合有流转的业务,Redis 适合性能敏感窗口可控的场景。

积压治理:从发现到止血

积压的根因只有两种:消费者跑慢了,或者生产者突然塞多了。

发现:lag 监控

最简单的监控是 kafka-consumer-groups 命令:

bash
# 查看所有 group 的 lag
kafka-consumer-groups --bootstrap-server localhost:9092 --all-groups --describe

# 输出示例
# GROUP           TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
# order-group     order-topic     0          15230           15680           450
# order-group     order-topic     1          13100           13100           0

单分区 lag 持续增长超过阈值(如 10000)就告警。生产环境可以集成到 Prometheus + AlertManager,用 Kafka Exporter 暴露 lag 指标。

止血三板斧

积压已经发生了,第一步不是修复,是先让系统活过来。

  1. 扩消费者:如果分区数 > 消费者数,直接加消费者实例,并行度提升。如果分区数已经等于消费者数,需要先增加分区再扩消费者(Kafka 分区数可动态增加,RocketMQ 类似)。
  2. 临时跳过:非核心消息直接抛弃,不处理。把积压量降下来,保全核心消息。
  3. 临时队列:把积压消息从原队列转移到备用的"快速队列",备用的消费者只做简单落库,不做复杂业务逻辑,处理速度翻几倍。

根治方向

止血之后,针对根因修复:

  • 批量消费max.poll.records 从默认 500 调到 2000,配合 fetch.max.bytes 调大,减少网络 RTT 占比。
  • 异步落库:消费端拿到的消息先写本地队列,再批量刷入数据库。把一次一条的 INSERT 变成一次 100 条的 batch INSERT。
  • 拆分 topic:核心消息和日志消息在同一 topic 抢资源,拆分后核心消息优先处理。

乱序治理:分区有序就够了

消息乱序的根因几乎都是重试。生产者发了一条消息到分区 P0,ack 超时(但实际 Broker 写成功了),重试发了另一条到 P0,后发的先到了。

Kafka 保证分区内有序:同一条消息多次重试,如果 max.in.flight.requests.per.connection=1(或开启幂等时自动设为 5 但幂等保证顺序),消息按序写入。但跨分区从来不保证顺序。

全局有序是个伪需求:把全局 100 万订单发到同一个分区,吞吐量就被压到那个分区了。真实场景里,按订单 ID 哈希到分区,同一个订单的消息一定在同一个分区,足够。

如果业务确实需要全局有序,方案是:单分区,关闭重试,消费端自己做补偿。代价是吞吐量降到单分区上限(约 5-10 MB/s)。这通常不是架构问题,而是需求没想清楚。

事故复盘模板

生产环境出问题不可怕,可怕的是复盘成了"反思会"而不是"改进会"。下面是一个三段式复盘骨架,每次 MQ 事故都按这个模板填:

## 事故标题:[时间] [Topic] MQ 问题

### 时间线
- T1: 告警触发(lag 飙到 XX)
- T2: 定位根因
- T3: 止血操作
- T4: 恢复
- T5: 复盘

### 根因
- 直接原因:一句话
- 触发条件:什么条件下才会复现
- 为什么没被预防:现有监控/限流/防护为什么没拦住

### 改进项
- [ ] 监控:补充 xxx 指标告警
- [ ] 代码:修复 xxx 幂等问题
- [ ] 流程:灰度发布增加 xxx 检查

小结

不丢、不重、不积压,这三件事在 MQ 日常运维里占了 80% 的精力。生产端落本地消息表兜底,消费端手动提交 + 幂等去重,lag 监控随手搭好,80% 的坑就已经填上了。剩下的 20% 靠事故复盘积累,每炸一次就把一个缺失的防护补上。

下一篇将进入分布式系统模块,从分布式理论开始,走完共识、选举、CAP 这些基础。

参考

参考:Kafka 官方文档(Reliability 章节)、RocketMQ 最佳实践、Google SRE 手册(监控与告警章节)

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。
粤ICP备2026104257号-1