1. 一次“消息静默消失”的线上事故:问题到底出在哪一环
先讲一个真实发生过的案例,这个案例几乎包含了我后面要说的所有坑。
凌晨两点,线上告警群里突然热闹起来。用户反馈:支付成功扣了款,但积分没到账、短信也没收到。订单服务日志显示支付成功消息已发出,消息中间件控制台也显示消息已存储,消费端却什么都没干。我第一反应是消费代码挂了,拉日志发现消费端压根没收到这条消息。于是顺着链路排查,最后定位到三个事实:
- 生产者用的是异步发送且没有回调处理,发送失败时只打了一行 debug 日志,没人看见。
- Broker 采用默认异步刷盘策略,消息写入页缓存就返回成功,当时所在的物理机恰好断电重启,未落盘的数据全部丢失。
- 消费者开启了自动提交 offset,业务线程在处理消息时抛异常退出,offset 却已经提交了。这条消息既不重试,也不告警,直接“蒸发”。
这个案例里,生产者、Broker、消费者三端各漏了一道防线,消息就彻底找不回来了。做消息队列的人常说“消息不丢失”是个系统工程,不是调某一个参数就能解决。消息从业务系统产生开始,到最终被消费方成功处理并落库,中间要经历生产者发送、Broker 存储、消费者拉取与确认三个大阶段,每个阶段都有独立的丢失风险。
也就是说,消息不丢失 = 生产者不丢 + Broker 不丢 + 消费者不丢,三者缺一不可。这篇文章我不会只讲概念,而是把七个关键防护点一层层拆开,每层对应什么风险、什么配置、什么代码写法,以及在真实项目中踩过的坑和验证方法,一次说透。
读这篇文章之前,假定你已经了解消息队列的基本概念,比如主题、分区、消费组、offset 这些词。如果对 Kafka 和 RocketMQ 的配置不熟也没关系,我会把两套主流 MQ 的对应实现都拉出来对比,你只需要掌握思路,换到任何 MQ 体系都能用。
2. 生产者端的三道闸门:把“发后即忘”变成“确认收到”
绝大多数消息丢失事故的源头在生产者侧,因为开发者最容易在这里使用“发后即忘”的方式。发后即忘并不是说消息一定丢,而是当发送失败发生时,你没有任何感知。我见过不止一个项目,生产者发送消息的代码就是一行producer.send(record),异常全部吞掉。这在低并发、网络稳定的环境里可能谁也发现不了问题,但一旦 Broker 重启、网络抖动或者 topic 不存在,消息就会悄悄丢掉。生产者端要做三道防线,核心思想是:每一次发送都必须有明确的成功或失败结论,并且对失败要有补偿手段。
2.1 第一层:同步发送或用回调感知结果
所谓“确认收到”,是要求生产者能够拿到 Broker 的确认结果。Kafka 的 Producer 有两种常见写法:
// 不推荐:发后即忘 producer.send(new ProducerRecord<>("order-event", orderId, message)); // 推荐:带回调的异步发送 producer.send(new ProducerRecord<>("order-event", orderId, message), (metadata, exception) -> { if (exception != null) { log.error("消息发送失败,topic={}, key={}", "order-event", orderId, exception); // 进入补偿流程,见第二层 } else { log.info("消息发送成功,partition={}, offset={}", metadata.partition(), metadata.offset()); } });如果是同步发送也简单,直接判断返回值。RocketMQ 里写法更直观:
SendResult sendResult = producer.send(message); if (sendResult.getSendStatus() != SendStatus.SEND_OK) { // 处理失败 }用回调或者同步发送,核心是拿到发送结果,而不是干等或者不管。有人会问:高性能场景下同步发送不是会很慢吗?这就是典型的需求权衡问题。对于订单、支付这类必须保证不丢的消息,哪怕多付出 1 毫秒的延迟也值得;对于日志采集、行为埋点这类允许少量丢失的数据,用发后即忘可以换取吞吐,这是完全合理的。关键是你要意识到自己舍弃了什么,而不是无意识地丢掉关键数据。
2.2 第二层:重试机制与补偿表,不放过一次瞬时失败
拿到发送失败结果之后,下一步是重试。Kafka 生产者自带重试机制,核心配置是retries和retry.backoff.ms:
acks=all retries=3 retry.backoff.ms=300 max.in.flight.requests.per.connection=1max.in.flight.requests.per.connection=1很关键,它限制了在单个连接上未确认请求的最大数量。如果不设这个值,重试可能会导致消息顺序错乱——第一条消息发送失败,第二条消息却先发出去了,等第一条重试成功时顺序就颠倒了。对于强顺序要求的场景,必须设置为 1。
RocketMQ 的同步发送失败后,你也可以自己封装重试:
for (int retry = 0; retry < 3; retry++) { try { SendResult result = producer.send(message); if (result.getSendStatus() == SendStatus.SEND_OK) { break; } } catch (Exception e) { Thread.sleep(300L * (retry + 1)); } }但重试不是万能的。Broker 若真正宕机,重试十次也没用。所以我还习惯在数据库里建一张消息补偿表,结构大概是这样的:
| 字段 | 说明 |
|---|---|
| id | 主键 |
| business_key | 业务唯一键,如订单号 |
| topic | 目标主题 |
| payload | 消息内容 |
| status | 0: 待发送, 1: 已发送, 2: 已确认, 3: 发送失败 |
| retry_count | 已重试次数 |
| next_retry_time | 下次重试时间 |
发送消息前先落入本地事务,把业务操作和消息状态更新放在同一个数据库事务里,也就是经典的“本地消息表”方案。再通过一个定时任务扫描 status 为 0 或 3 且超过 next_retry_time 的记录,重新投递。这保证了:只要本地业务成功了,消息至少不会丢,最坏情况是延迟到达。
2.3 第三层:事务消息,解决“业务成功但消息发送失败”的原子性问题
本地消息表需要额外建表、写定时任务,有些人觉得麻烦,于是 RocketMQ 直接提供了事务消息机制。它的流程我简单描述一下:
- 生产者发送 half message,此时对消费者不可见。
- 执行本地业务事务。
- 如果本地事务成功,提交消息;失败则回滚消息。
- 如果步骤 3 因为宕机等原因没有执行,Broker 会回查生产者本地事务状态,根据结果决定提交或回滚。
代码大致长这样(以 RocketMQ 为例):
TransactionMQProducer producer = new TransactionMQProducer(); producer.setTransactionListener(new TransactionListener() { @Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 执行本地业务,比如写订单表 try { orderService.createOrder((OrderDO) arg); return LocalTransactionState.COMMIT_MESSAGE; } catch (Exception e) { return LocalTransactionState.ROLLBACK_MESSAGE; } } @Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // 回查本地事务状态 return orderService.isOrderCreated(msg.getKeys()) ? LocalTransactionState.COMMIT_MESSAGE : LocalTransactionState.ROLLBACK_MESSAGE; } });Kafka 也有事务 API,但用起来复杂,而且业界用 Kafka 做事务消息的比例明显低于 RocketMQ。如果你用的是 Kafka,我更推荐本地消息表方案,因为它通俗易懂,也不依赖特定版本特性。
这一层防线做的事用一句话总结:把消息发送和业务写库放进同一个“事务决策”里,杜绝“钱扣了,消息没发出去”的惨案。
3. Broker 端的两道关卡:写入可靠与存储可靠缺一不可
过了生产者这一关,消息已经送到 Broker 手里了。但 Broker 并不是保险柜,它自身也面临两个风险:第一,收到消息后只写了内存/页缓存就返回成功,机器断电就丢;第二,数据虽然写入了本地磁盘,但磁盘损坏或者机器报废,数据依然跟着没了。这两类风险对应两层防线。
3.1 第四层:同步刷盘,让消息真正落在磁盘上
先讲一个小原理。几乎所有的 MQ 写入消息时,并不是直接写磁盘文件,而是先写入操作系统的页缓存(Page Cache),然后由操作系统异步刷盘。这样做的目的是利用内存的高性能来换取吞吐量,但代价是:如果写入页缓存还没刷盘时进程崩溃或机器断电,数据就丢了。
以 RocketMQ 为例,Broker 的刷盘策略有两种:
| 刷盘策略 | 行为 | 可靠性 | 性能 |
|---|---|---|---|
| ASYNC_FLUSH | 写入页缓存即返回成功 | 低,断电可能丢 | 高 |
| SYNC_FLUSH | 写入磁盘文件后才返回成功 | 高 | 较低,但可接受 |
配置位于 Broker 的broker.conf:
flushDiskType=SYNC_FLUSHKafka 的刷盘配置不太一样,它更依赖副本机制来保证可靠性,而本机刷盘由log.flush.interval.messages和log.flush.interval.ms控制。默认值是比较宽泛的,如果希望更可靠,可以调小这些值,但代价是引入更多次磁盘写入。
我在生产环境里的经验是:订单、支付、账户等核心链路用同步刷盘,日志、行为数据用异步刷盘。不要一刀切全上同步刷盘,那样会把日志型的高吞吐场景拖垮;也不要在核心链路上贪图性能用异步刷盘,因为一次断电就能让你损失一批关键消息,业务恢复成本远比省下的那点性能高。
3.2 第五层:多副本与 ISR,坏一台机器也不丢数据
刷盘只能防断电,防不了磁盘损坏、机器报废。这时候需要副本机制。
Kafka 的副本机制核心概念是 ISR(In-Sync Replicas,同步中的副本)。Leader 分区的数据会同步到多个 Follower 副本,Producer 写入时可以通过acks参数控制需要多少个副本确认:
acks=0:不等待确认,最多丢。acks=1:Leader 写入成功就返回,Leader 宕机可能丢。acks=all:所有 ISR 副本都写入成功才返回,最安全。
配合min.insync.replicas参数,可以设定最少需要几个副本同步成功。最稳妥的组合是:
acks=all min.insync.replicas=2意思是至少要有一个 Follower 和 Leader 保持同步,写入才算成功。如果 Follower 全部挂掉,写入会失败,而不是静默接受,这逼着生产者走重试或补偿。
RocketMQ 也有对应概念,主从模式下 Broker 的配置:
brokerRole=SYNC_MASTERSYNC_MASTER表示主节点需要等待从节点复制成功后才返回。换成ASYNC_MASTER则会丢掉这条保护。
这里特别提醒一个 Kafka 的隐藏坑:unclean.leader.election.enable一定要设为 false。如果设为 true,当所有同步副本都挂了,Kafka 可以选一个不同步的副本当 Leader——这样做的好处是服务可用性提高了,但代价是消息大量丢失。对“消息不丢失”有刚需的业务,这个参数必须关掉。宁可短暂不可用,也不能让数据悄悄丢。
4. 消费者端的手动确认:自动提交是丢消息的隐形杀手
我在最开始那个案例里说消费者这端也丢了消息,很多人不理解:消费者不是只负责收消息吗,怎么会丢?问题就出在 offset 的提交时机上。
偏移量(offset)可以理解成书签。消费者读完一批消息后,要把书签记录到 Broker,下次拿着书签继续读后面的消息。如果在消息处理完之前就更新了书签,等于书签已经翻页,内容还没看懂。此时消费者进程崩溃,重启后会从书签位置继续消费——处理失败的那批消息永远不会再读到了。
Kafka 消费者默认是自动提交 offset 的,相关配置是:
enable.auto.commit=true auto.commit.interval.ms=5000每 5 秒自动提交一次。如果你的业务逻辑处理时间超过 5 秒,消息处理到一半,offset 已经提交了;或者代码在poll()之后for循环里处理消息,第 1 条处理失败抛异常,后续的几条压根没执行,但 offset 还是被自动提交了。这就是不丢消息的大忌。
正确做法是把自动提交关掉,改成手动提交,而且要在消息处理成功后再提交:
props.put("enable.auto.commit", "false"); while (running) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { try { process(record); // 业务处理 consumer.commitSync(); // 处理成功后再提交 } catch (Exception e) { log.error("消费失败,等待下次重试或进入死信流程", e); // 注意:这里 continue,不要提交 } } }RocketMQ 的消费者机制不太一样,它是通过返回状态来决定是否确认的:
consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) -> { try { businessService.process(msgs); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; } catch (Exception e) { return ConsumeConcurrentlyStatus.RECONSUME_LATER; } });返回RECONSUME_LATER的消息会在之后的重试队列中再次尝试,而不是直接被丢弃。这里最忌讳的就是不管业务是否成功,一律返回消费成功,这在 RocketMQ 中极为常见,尤其是消费逻辑里 catch 住了异常并“吞掉”,从现象上看消息全部消费成功,实际业务数据全没落库。
所以消费者这一层防线的核心是:你必须在业务处理真正成功时才确认消息。如果拿不准,宁可返回失败让它重试,也不要模棱两可地确认。重试顶多带来重复消费,而你还有最后一道防线兜底。
5. 第七层防线:重试、死信与幂等,把“最后一公里”焊死
即使你做到了生产者确认、Broker 同步刷盘、多副本、消费者手动提交,消息依然有可能重试。为什么?因为“不丢失”和“不重复”是两回事。几乎所有主流 MQ 保证的是 At Least Once(至少一次)投递,也就是说消息可能不丢,但可能重复。当消费者的业务代码处理消息后没有来得及提交 offset,进程就崩溃了,重启后会重新消费这一条——这就产生了重复。应对重复的方法不是去消灭它,而是让重复消费变得无害,这正是第七层防线的意义。
5.1 消费失败的重试策略与死信队列
消费者拿到消息后,处理失败怎么办?第一步是重试。Kafka 场景下,需要自己维护重试逻辑。最简单的方法是:把消费失败的消息写入一个本地待重试表,用定时任务扫描重发。更讲究的做法是利用 Kafka 的重试主题,设计不同延迟级别的重试 Topic,比如 1 秒、10 秒、60 秒分别投递一次。RocketMQ 则内置了重试队列,消费失败的消息会按照延迟等级自动重试,默认最多 16 次。
重试次数用完之后,消息去哪?不能直接丢掉,而应该送进死信队列(DLQ)。RocketMQ 有自动死信队列,消息重试耗尽后就进入%DLQ%消费组名,你可以在控制台查看、手动介入。Kafka 没有死信队列的概念,需要自己实现:消费者在多次重试失败后,把消息发给一个专门的dlq-topic,再由人工处理程序或者告警介入。
我在项目中给死信队列配了一套告警规则:只要死信主题有消息产生,就触发企业微信/钉钉通知,并附上最后一批消费失败的原始消息内容。因为死信队列里躺着的基本都是脏数据或者业务 bug,不及时处理就会演变成资损事故。有些团队会把死信队列的消息定期删掉,那是非常危险的,等于把日志销毁了,后面根本没法追溯。
5.2 幂等兜底:同一个消息重复消费也不怕
这里必须聊幂等。所谓幂等,就是同一个操作执行一次和执行多次结果一致。常见的做法有四种:
基于数据库唯一索引。比如消费消息后要插入一张积分流水表,给业务的唯一键(比如 orderId)建唯一索引,重复插入会被数据库拦住,直接跳过。
基于 Redis SETNX。消费前先SETNX consume_lock_{orderId} 1,拿到锁才处理,处理完再删掉。注意设置锁的过期时间,防止服务宕机锁不释放。
基于状态机。比如订单状态流转:待支付 → 已支付 → 已发货。重复消费时如果发现订单状态已经不是“待支付”了,说明已经处理过,直接跳过。
基于消息内业务键去重。在本地维护一张已处理消息表,记录 messageId,处理前先查这个表。
我做消息消费的第一原则是:不管前面几层做得有多好,消费者的落库操作都必须有幂等保护。因为重复消息是客户端无法完全避免的,哪怕你前面配置拉满,在极端情况下(比如消费者处理完、提交 offset 前的网络分区)依然会产生重复。幂等是最后一道防线,这道防线没守住,前面的努力都可能白费。
这里还要提醒一个很容易踩的坑:如果你在消费逻辑里先做了各种判断,最后发现是重复消息就 return,那你需要确保 return 之前把这条消息的商品事务、缓存、统计全部一致地跳过。我见过一个系统,数据库唯一索引防住了重复插入,但漏了 Redis 缓存的更新,导致重复消息插入被挡住但缓存状态被覆盖,反而产生了数据不一致。幂等方案必须覆盖到所有可能触发的旁路逻辑,不只是主库写操作。
6. 7 层防线的落地清单:从配置到故障演练的完整方案
讲到这里,七层防线全部出现,我先把它们汇总成一张检查表,然后在下面给出落地建议。
| 层级 | 防线名称 | 关键操作 | 常出问题的默认值 |
|---|---|---|---|
| 1 | 生产者发送确认 | 同步发送或回调 | 异步无回调 |
| 2 | 生产者重试补偿 | 设置 retries、本地消息表 | 无重试、无补偿 |
| 3 | 业务事务原子性 | 本地消息表或事务消息 | 业务成功后消息丢失 |
| 4 | Broker 持久化 | 同步刷盘 | 异步刷盘 |
| 5 | 多副本同步 | 副本数≥2、同步复制 | 单副本或异步复制 |
| 6 | 消费者手动确认 | 业务成功后提交 offset | 自动提交 |
| 7 | 重试、死信与幂等 | 重试策略+死信告警+幂等表 | 无限重试或直接丢弃 |
这张表可以当成你接手任何一个消息项目的第一份排查清单。上生产环境之前,我建议你按这个顺序逐项对照检查配置。
6.1 一套推荐的 Kafka 配置模板
如果你用的是 Kafka,并且消息属于“一定不能丢”的类型,可以直接参考这套配置:
# 生产者端 acks=all retries=5 retry.backoff.ms=300 max.in.flight.requests.per.connection=1 enable.idempotence=true # Broker 端 min.insync.replicas=2 unclean.leader.election.enable=false log.flush.interval.messages=10000 log.flush.interval.ms=1000 # 消费者端 enable.auto.commit=false auto.offset.reset=earliest这里多说一句enable.idempotence=true的作用。它让生产者具备幂等能力,即使客户端重试发送,Broker 也能识别重复消息并避免重复写入。它和acks=all配合使用,能基本消除生产者重试导致的重复消息。注意开启幂等发送后,max.in.flight.requests.per.connection会默认被设置为 5,如果你同时有顺序性要求,还是需要显式设回 1。
6.2 一套推荐的 RocketMQ 配置模板
RocketMQ 场景下,重点在 Broker 和消费者:
# Broker 配置 brokerRole=SYNC_MASTER flushDiskType=SYNC_FLUSH # 消费者:注意别吃异常,重试次数按业务设置 consumer.setConsumeTimeout(15); consumer.setMaxReconsumeTimes(16);RocketMQ 还有一点和 Kafka 不一样:虽然 Broker 用了同步刷盘、同步复制,但 producer 如果用了单向发送(sendOneWay)或者异步发送没有正确处理SendCallback,前面的努力等于白费。所以 RocketMQ 的核心链路我全部使用producer.send()同步发送,并捕获异常。
6.3 怎么验证消息真的不丢
配置完之后,怎么验证你的防线是否有效?光看配置不叫验证,要做故障演练。
我常用的验证方式分为三个步骤:
第一步:断开 Broker 网络,验证生产者补偿逻辑。挑一台测试环境上的应用,用iptables禁用其访问 MQ 的端口,观察生产者是否按预期报错、重试、进入本地补偿表。恢复网络后,检查补偿任务是否把积压的消息全部补发成功,且消息顺序没有乱。
第二步:kill -9 模拟 Broker 宕机,验证消息已刷盘。往队列里连续写入一批核心消息,在数据刚写入但可能未刷盘的时刻强制断电(测试环境可以直接断虚拟机电源),重启 Broker 后检查消息是否存在。同步刷盘模式下,写入成功的消息必须全部还在。如果你用的是异步刷盘,你会发现确实丢了最后一批——这就是为什么我反复强调核心链路用同步刷盘。
第三步:消费者 kill -9 模拟处理中崩溃,验证 offset 提交时机。开启一个大批量消费任务,在处理一批消息的中途 kill -9 消费者进程,重启后观察之前那条正在处理的消息是否被再次投递。如果配置了手动提交且没有提前提交 offset,那这条消息会被重新消费;此时你的幂等逻辑应该保证数据不重复。
这三步做完,心里才有底。我见过太多团队把生产者的acks=all配上就觉得消息一定不丢了,直到断电演练才傻眼——Broker 刷盘策略还是默认的异步模式,一把电闸就丢了几千条消息。所以验证永远不是可选项。
6.4 日志、监控与告警:让消息丢在“明处”
最后一小节,说说比技术配置更重要的东西:可观测性。消息丢没丢,你的系统自己要有能力感知。我的统一做法是:
- 生产者端,每次发送成功打一条包含 topic、partition、offset、耗时、messageId 的日志,失败打 ERROR 日志并附上完整消息内容。
- 消费者端,在消息进入和离开业务处理逻辑处各打一条日志,带上 messageId 和业务唯一键。
- Broker 层面,监控未确认消息数、死信队里堆积量、消费组滞后量(Lag)。
- 设置两条告警规则:死信队列有消息进来说明消费链路出问题了,需要立刻介入;消费组 Lag 持续增长说明消费速度跟不上,需要扩容或者排查阻塞。
日志的意义在于,当消息真的丢了,你能在第一时间定位到是在哪一环丢的。如果生产者日志显示已发送、Broker 控制台显示已写入、消费者日志显示没收到,那问题大概率出在消费者订阅或者分组配置上。只要有一环的日志缺失或不准,排查就会变成大海捞针。
我在实际项目中还养成一个习惯:对每个核心消息都会生成全局唯一的 messageId,从生成到发送、存储、消费、落库,全程透传到日志和数据库表。这样不管哪一环出问题,都可以用 messageId 把整条链路的日志串起来看。这个习惯帮我节省过太多排查时间。
最后分享一点个人习惯
文章写到这里,7 层防线和落地方法都讲完了。最后还是想补充一点非常个人的经验:接手任何一套带消息队列的系统,我都会先向团队问三个问题——消息丢了有没有感知?有没有补偿和重试兜底?重复消费有没有幂等保护?如果三问里有一问答不上来,那这个系统就没有达到“消息不丢失”的标准。七层防线不是让你全部堆满,而是让你在每一层的取舍上都有明确依据。核心业务我建议一层都不要省,非核心业务可以按风险与成本做裁剪。这样一来,既避免了“一刀切”的性能浪费,也杜绝了“裸奔”式的数据丢失。希望这篇文章能给你带来一点实实在在的参考。