把 Kafka 当队列用,丢了 0.3% 的消息:Kafka 与 RocketMQ 在可靠性、顺序、事务上的 4 笔真实账
title: 把 Kafka 当队列用,丢了 0.3% 的消息:Kafka 与 RocketMQ 在可靠性、顺序、事务上的 4 笔真实账
date: 2026-08-22
category: 消息队列
tags: [Java, 后端, 消息队列, Kafka, RocketMQ, 微服务]
我们交易系统的异步链路一开始全用 Kafka。理由很朴素:社区火、吞吐高、大家都用。直到一次对账发现,有 0.3% 的支付成功消息下游没收到,导致订单状态卡在「支付中」。排查下来不是 Kafka 的 bug,而是我们用错了它的可靠性模型——Kafka 默认是「不丢 broker 端」,但「不保证消费端一定处理成功」。
这篇不聊谁更牛,只聊我们实打实在这俩队列上花的 4 笔账:可靠性、顺序、事务、堆积。选型翻车,往往是因为你需要的那一项,正好是它弱的那一项。
第一笔账:可靠性,Kafka 默认 at-least-once
Kafka 的生产者默认acks=1(leader 收到就返回),broker 宕机可能丢消息。要「不丢」得acks=all+ 副本min.insync.replicas>1:
// Kafka 生产者:要 broker 端不丢,得这么配 Properties p = new Properties(); p.put("acks", "all"); // 等所有 ISR 副本确认 p.put("min.insync.replicas", "2"); // 至少 2 个副本落盘 p.put("enable.idempotence", "true"); // 开启幂等,避免重试导致重复 p.put("retries", "5"); Producer<String, String> producer = new KafkaProducer<>(p);但即便 broker 不丢,消费端也可能丢:Kafka 是「拉取 + 手动提交 offset」模型,如果你处理完业务逻辑、还没提交 offset 就崩了,这条消息会被重新投递——这是 at-least-once。我们那次丢消息的根因是反过来的:为了快,我们先提交 offset 再处理,结果处理逻辑抛异常,消息算「已消费」永远丢了。
// 错误示范:先 commit 再处理,处理失败消息就没了 consumer.commitSync(); // offset 已提交 process(record); // 这里抛异常,消息丢失 // 正确做法:先处理再提交,用幂等兜底重复 process(record); consumer.commitSync(); // 处理成功才提交RocketMQ 默认是at-least-once+消费重试队列,消费失败会进重试队列(默认 16 次),还不行进死信队列,消息不会凭空消失。这对「消息不能丢」的业务更友好。
第二笔账:顺序消息,Kafka 靠分区、RocketMQ 有原生支持
我们需要「同一个订单的支付、发货、完成消息严格有序」。Kafka 只能保证单个分区内有序,所以得把同一个订单 ID 哈希到同一个 partition:
// Kafka:相同 orderId 进同一 partition,靠分区器保证局部顺序 producer.send(new ProducerRecord<>( "order-topic", orderId.hashCode() % partitionCount, // key 决定分区 orderId, payload));但代价是:一旦某个分区消费慢,会阻塞整个分区的后续消息(Kafka 单分区是串行消费的)。RocketMQ 则提供「顺序消息」语义,broker 端用分段锁保证同一 queue 的消息顺序,消费者也顺序拉取,语义更省心。我们订单状态机用 RocketMQ 顺序消息后,省掉了自己维护「分区拥堵监控」的那套告警。
第三笔账:事务消息,这是两者分水岭
我们要「本地订单落库」和「发支付消息」要么都成、要么都败。Kafka 有事务 API,但用起来很重:
// Kafka 事务消息:要 encloser id + initTransactions producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord<>("order-topic", payload)); // 本地事务(这里要自己保证,Kafka 事务只管消息) orderDb.insert(order); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); }Kafka 的事务本质是「把多条消息原子提交」,它不保证和你的本地 DB 事务绑定——中间还是可能「DB 成了、消息回滚」或反过来(要用最大努力通知或额外对账兜底)。RocketMQ 的事务消息则专门为此设计:半消息 + 回查机制,broker 会主动问你「本地事务到底成了没」:
// RocketMQ 事务消息:实现 TransactionListener 回查本地事务状态 TransactionMQProducer producer = new TransactionMQProducer("order_group"); producer.setTransactionListener(new TransactionListener() { @Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { // 执行本地 DB 事务 return orderDb.insert(order) ? COMMIT_MESSAGE : ROLLBACK_MESSAGE; } @Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { // broker 回查:去 DB 查这笔订单到底在不在 return orderDb.exists(msg.getKeys()) ? COMMIT_MESSAGE : UNKNOW; } }); producer.sendMessageInTransaction(new Message("order-topic", payload.getBytes()), null);这笔账我们交了学费:用 Kafka 硬做事务,最后还是靠每小时对账补单,多写了 200 行对账代码。后来涉及「DB + 消息」强一致的链路迁到 RocketMQ 事务消息,对账代码删了一大半。
第四笔账:延迟消息,RocketMQ 开箱即用
还有一个我们真实用到的场景:订单创建后 30 分钟未支付自动关单。RocketMQ 原生支持延迟消息:
// RocketMQ:指定延迟级别,18 级覆盖 1s~2h Message msg = new Message("order-close-topic", payload.getBytes()); msg.setDelayTimeLevel(16); // 16 级 = 30 分钟 producer.send(msg);Kafka 没有原生延迟消息,我们得自己搞「时间轮 + 外部存储」或者「建多个延迟 topic 定时搬运」,多了一套组件要维护。如果你业务里关单、重试、定时通知这类场景多,RocketMQ 的延迟消息能省不少事。
第五笔账:堆积与吞吐,Kafka 仍是天花板
反过来,纯吞吐和堆积量,Kafka 确实更强。我们埋点日志(单条小、允许少量重复、不需要顺序)用 Kafka,单集群轻松扛百万 TPS,堆积上亿条也没压力。RocketMQ 在「功能丰富」(延迟消息、事务、顺序、重试)上更全,但极端吞吐略逊。
| 维度 | Kafka | RocketMQ |
|---|---|---|
| 默认可靠性 | at-least-once(需配 acks=all) | at-least-once + 重试/死信 |
| 顺序消息 | 单分区有序,靠 key 哈希 | 原生顺序消息语义 |
| 事务消息 | API 重,不绑 DB 事务 | 半消息 + 回查,绑 DB 友好 |
| 延迟消息 | 无原生,需自建 | 18 级原生支持 |
| 堆积/吞吐 | 极高 | 高,功能更全 |
消费端幂等:两个队列都需要的兜底
不管选 Kafka 还是 RocketMQ,at-least-once 语义都意味着「消息可能重复投送」。所以消费端必须自己做幂等,这是绕不开的。我们用「消息唯一键 + Redis 去重」兜底:
// 消费幂等:同一 messageId 只处理一次 public void consume(String messageId, String payload) { // SETNX:key 不存在才设成功,返回 1 表示首次处理 Long r = redis.opsForValue() .setIfAbsent("msg:" + messageId, "1", 24, TimeUnit.HOURS); if (r == null || r == 0) { return; // 重复消息,直接丢弃 } businessProcess(payload); // 真正处理 }这个幂等层对两个队列都通用,也正是前面说的「用 Kafka 做关键队列要自己补的兜底」之一。RocketMQ 虽然自带重试/死信,但重试过来的消息 messageId 会变,去重得用业务自己的唯一键(比如订单 ID),不能依赖消息系统给的 ID。这点踩过:我们一度用 RocketMQ 的 msgId 去重,结果重试消息 msgId 不同,去重失效,还是重复处理了。
我的选型判断
别再问「Kafka 和 RocketMQ 谁更好」——这问题本身错。我们现在的架构是混用:
- 日志、埋点、流式计算 → Kafka,要的是吞吐和生态(Flink 集成好)。
- 交易、订单、履约这类「消息不能丢、要顺序、要事务」的核心链路 → RocketMQ,它把可靠性相关的坑都帮你填了。
踩完 0.3% 丢消息这个坑后我形成了一个观点:用 Kafka 做关键业务队列,你得自己补幂等、补对账、补重试,等于自己造半个 RocketMQ。如果团队人力紧张、业务又关键,直接用 RocketMQ 省下的排障时间,远比那点吞吐差距值钱。只有当你真的需要 Kafka 的生态(流处理、海量日志)时,才值得为它额外写那些兜底代码。
复盘数字
那次丢消息持续约 9 小时(夜间批处理 + 白天高峰叠加),最终对账补单 1.2 万笔,客诉 37 起。把交易链路迁到 RocketMQ 事务消息 + 消费幂等后,3 个月同类「消息丢失」告警为零;而日志链路继续用 Kafka,单日吞吐稳定在 800 万条以上,堆积峰值 2 亿条时消费无阻塞。混用之后,核心链路补单率从 0.3% 降到万分之一以下(剩余的是业务侧自身重试)。
思考题
RocketMQ 的事务消息用「半消息 + 回查」解决 DB 与消息的一致性,但回查本身有次数上限,如果本地事务卡死、回查一直返回 UNKNOW,这条半消息会怎样?你会在业务侧怎么兜底?
