Kafka 和 RocketMQ 到底谁扛得住事务消息:我们切了 3 次队,第 2 次差点背锅

发布时间:2026/7/29 15:37:39
Kafka 和 RocketMQ 到底谁扛得住事务消息:我们切了 3 次队,第 2 次差点背锅 我们交易系统的异步化做过三次消息队列选型一路从 Kafka 换到 RocketMQ 又想换回来中间有一次因为用错消息队列做分布式事务导致一笔订单扣了库存却没发券客诉直接砸到我们组头上。那次事故逼我把两家的底层差异啃了一遍。这篇不堆参数从我们真实踩的坑出发对比 Kafka 和 RocketMQ 在事务消息、顺序消息、消费模型上的不同并附上代码。先说结论性的差异出身决定能力Kafka 为高吞吐日志流而生设计目标是百万级 TPS、持久化到磁盘但追求顺序写RocketMQ 出身阿里电商为业务消息可靠投递设计强在事务消息、定时消息、顺序消息、丰富的消费重试。所以选谁先看你要流还是业务通知。// Kafka 生产端追求吞吐acks 配置决定可靠性 Properties p new Properties(); p.put(bootstrap.servers, k1:9092,k2:9092); p.put(acks, all); // 1. 等所有 ISR 副本确认最可靠 p.put(enable.idempotence, true); // 2. 开启幂等避免重试导致重复 p.put(retries, 3); ProducerString, String producer new KafkaProducer(p); producer.send(new ProducerRecord(order-event, orderId, json)); // 3. 发完不等结果异步逐行第 1 行acksall让生产者在所有同步副本都写入后才算成功牺牲一点延迟换不丢第 2 行幂等生产者给每条消息编序号broker 去重避免网络重试产生重复第 3 行send默认异步吞吐极高但发出去不等于消费到了。Kafka 的可靠靠副本 幂等 异步批量但它在事务消息和精准一次投递上比 RocketMQ 繁琐。事务消息RocketMQ 的 Half Message 更顺手我们那次背锅正是因为用 Kafka 硬做扣库存 发券的事务一致性。Kafka 有事务 API但它是生产端原子发送多条消息的语义不是本地事务和发消息绑定的语义用它做数据库扣库存成功了消息才发很容易写错。RocketMQ 的事务消息是专门为此设计的 Half Message 机制// RocketMQ 事务消息本地事务和发消息绑定 TransactionMQProducer producer new TransactionMQProducer(order_group); producer.setTransactionListener(new TransactionListener() { Override public LocalTransactionState executeLocalTransaction(Message msg, Object arg) { try { deductStock((Order) arg); // 1. 先扣库存本地事务 return LocalTransactionState.COMMIT; // 2. 成功 - 半消息转可消费 } catch (Exception e) { return LocalTransactionState.ROLLBACK; // 3. 失败 - 半消息丢弃 } } Override public LocalTransactionState checkLocalTransaction(MessageExt msg) { return stockDeducted(msg) ? COMMIT : ROLLBACK; // 4. 回查Broker 未收到确认时调用 } }); Message m new Message(order-topic, , orderId, json.getBytes()); producer.sendMessageInTransaction(m, order); // 5. 先发 Half Message再执行本地事务逐行第 5 行先发一条半消息Half Message对消费者不可见第 1 行执行本地扣库存第 2/3 行根据结果决定 COMMIT半消息变可见还是 ROLLBACK半消息删除。关键在于第 4 行checkLocalTransaction——如果本地事务执行完、Broker 没收到确认比如 producer 宕机Broker 会主动回查本地状态决定提交还是回滚保证本地事务成功则消息必发、失败则必不发。Kafka 没有这种回查机制要自己额外做事务表 补偿任务复杂得多。这就是我们那次事故的根用 Kafka 做分布式事务补偿逻辑我们自己写得有漏洞导致库存扣了、券没发。消费模型拉 vs 推影响重试体验Kafka 是拉pull模型consumer 主动去 partition 拉RocketMQ 是推push模型底层也是拉但 SDK 封装成推并内置消费失败重试、死信队列。// Kafka 消费手动提交 offset控制精确一次 while (true) { ConsumerRecordsString, String recs consumer.poll(Duration.ofMillis(100)); // 1. 拉取 for (ConsumerRecordString, String r : recs) { process(r); // 2. 业务处理 } consumer.commitSync(); // 3. 处理完再同步提交 offset }逐行第 1 行 consumer 主动poll第 3 行commitSync在处理完后提交位移如果处理到一半宕机重启会重拉这批可能重复处理——所以 Kafka 的不丢要靠业务幂等兜底。RocketMQ 的推模型在此基础上消费失败自动重试默认 16 次、间隔递增并进入死信队列对业务更友好// RocketMQ 消费失败自动重试 死信 consumer.registerMessageListener((MessageListenerConcurrently) (msgs, context) - { try { for (MessageExt m : msgs) handle(m); return ConsumeConcurrentlyStatus.CONSUME_SUCCESS; // 1. 成功 } catch (Exception e) { return ConsumeConcurrentlyStatus.RECONSUME_LATER; // 2. 失败 - 自动重试 } });逐行第 2 行返回RECONSUME_LATERRocketMQ 会把消息重新投递最多 16 次后进死信队列%DLQ%你定期捞死信人工处理即可。我们后来的订单消费全用 RocketMQ就是因为失败自动重试 死信省掉了自己写重试调度的大量代码。顺序消息谁更天然订单状态流转创建→支付→发货要求同订单消息严格有序。Kafka 的顺序体现在单 partition 内有序要同订单进同 partition 得自己按 orderId 哈希路由RocketMQ 的顺序消费MessageListenerOrderly在 SDK 层面保证单队列串行消费。// Kafka同 key 进同 partition 才有序 producer.send(new ProducerRecord(order, orderId, json)); // 1. orderId 相同必进同 partition // 消费端单 partition 单线程消费即可保序 // RocketMQ顺序消费监听器 consumer.subscribe(order-topic, *); consumer.registerMessageListener((MessageListenerOrderly) (msgs, context) - { for (MessageExt m : msgs) handle(m); // 2. 单队列串行天然有序 return ConsumeOrderlyStatus.SUCCESS; });逐行Kafka 第 1 行靠 orderId 哈希保证同订单进同 partition再配合单线程消费才有序配置稍繁琐RocketMQ 第 2 行MessageListenerOrderly直接锁队列串行消费语义上更直白。但顺序消费的代价是吞吐下降队列被串行锁住高并发下要权衡。横向对比别只看 TPS 数字维度KafkaRocketMQ设计目标高吞吐日志/流业务消息可靠投递事务消息有生产者事务无回查有Half Message 回查贴合业务消费失败重试需自写offset 重放内置 16 次重试 死信队列顺序消息单 partition 内有序顺序消费监听器单队列串行定时/延迟消息不支持原生靠外部支持18 个延迟级别 / 任意时间吞吐量极高百万 TPS 级高十万 TPS 级够绝大多数业务运维复杂度中依赖 ZooKeeper/新版 KRaft中NameServer 轻量注意Kafka 吞吐量通常高于 RocketMQ但绝大多数业务系统的瓶颈不在 MQ 吞吐而在下游数据库。我们压测订单链路RocketMQ 十万级 TPS 完全喂不饱瓶颈在 MySQL。所以Kafka 更快对业务系统经常是伪命题。我的取舍业务通知用 RocketMQ日志流用 Kafka我的观点凡是涉及业务一致性、失败要重试、要事务的消息用 RocketMQ凡是日志、埋点、流计算这类丢了能忍、要吞吐的用 Kafka。我们现在的架构就是这么分的——交易链路的订单、库存、券用 RocketMQ吃它的事务消息 重试 死信用户行为日志、监控事件用 Kafka吃它的吞吐和流处理能力。当初那次背锅本质是拿 Kafka 当业务事务队列用硬凑补偿逻辑出了 bug。如果早按这个边界切分根本不会出那档子事。另外提醒一点RocketMQ 的事务消息也有代价——半消息机制增加了一轮往返和回查压力且回查逻辑必须幂等、必须查得到本地状态否则回查也会出问题。别以为用了事务消息就万事大吉本地事务和回查方法照样要写对。思考题假设你用 Kafka 做扣库存后发券的一致性不用 RocketMQ 事务消息你会在本地怎么设计一张事务状态表 补偿任务来兜底它和 RocketMQ 的 Half Message 回查相比多写了哪些代码、多承担了哪些风险写在最后Kafka 和 RocketMQ 没有谁碾压谁是定位不同Kafka 为流、为吞吐RocketMQ 为业务、为可靠。我们切了三次队才落地业务用 RocketMQ、日志用 Kafka的边界。那次背锅的教训是拿错工具做错的事比工具不够快危险得多。选型先看你的消息是流还是业务通知再决定用谁别被 TPS 数字带偏。