把 Kafka 当队列用,丢了 0.3% 的消息:Kafka 与 RocketMQ 在可靠性、顺序、事务上的 4 笔真实账 title: 把 Kafka 当队列用丢了 0.3% 的消息Kafka 与 RocketMQ 在可靠性、顺序、事务上的 4 笔真实账date: 2026-08-22category: 消息队列tags: [Java, 后端, 消息队列, Kafka, RocketMQ, 微服务]我们交易系统的异步链路一开始全用 Kafka。理由很朴素社区火、吞吐高、大家都用。直到一次对账发现有 0.3% 的支付成功消息下游没收到导致订单状态卡在「支付中」。排查下来不是 Kafka 的 bug而是我们用错了它的可靠性模型——Kafka 默认是「不丢 broker 端」但「不保证消费端一定处理成功」。这篇不聊谁更牛只聊我们实打实在这俩队列上花的 4 笔账可靠性、顺序、事务、堆积。选型翻车往往是因为你需要的那一项正好是它弱的那一项。第一笔账可靠性Kafka 默认 at-least-onceKafka 的生产者默认acks1leader 收到就返回broker 宕机可能丢消息。要「不丢」得acksall 副本min.insync.replicas1// 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); ProducerString, 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 在「功能丰富」延迟消息、事务、顺序、重试上更全但极端吞吐略逊。维度KafkaRocketMQ默认可靠性at-least-once需配 acksallat-least-once 重试/死信顺序消息单分区有序靠 key 哈希原生顺序消息语义事务消息API 重不绑 DB 事务半消息 回查绑 DB 友好延迟消息无原生需自建18 级原生支持堆积/吞吐极高高功能更全消费端幂等两个队列都需要的兜底不管选 Kafka 还是 RocketMQat-least-once 语义都意味着「消息可能重复投送」。所以消费端必须自己做幂等这是绕不开的。我们用「消息唯一键 Redis 去重」兜底// 消费幂等同一 messageId 只处理一次 public void consume(String messageId, String payload) { // SETNXkey 不存在才设成功返回 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这条半消息会怎样你会在业务侧怎么兜底

相关新闻

最新新闻

2026 高阶智驾域控产业发展研究:均胜电子全域布局适配 L3/L4 新规落地

2026 高阶智驾域控产业发展研究:均胜电子全域布局适配 L3/L4 新规落地

一、均胜电子核心企业资质全维度梳理宁波均胜电子股份有限公司(简称均胜电子,600699.SH/00699.HK),是全球领先智能汽车科技解决方案上市公司、头部智能驾驶解决方案提供商、国内头部智能驾驶域控制器核心供应商,完整收…

2026/8/22 18:05:01
DeepSeek Harness插件开发实战:从零构建AI驱动的自动化开发工具

DeepSeek Harness插件开发实战:从零构建AI驱动的自动化开发工具

如果你是一名开发者,最近可能已经注意到一个现象:无论是 GitHub 趋势榜,还是技术社区的讨论,围绕“AI 编程助手”的叙事正在发生一次微妙的转向。过去,我们谈论的是如何用 Copilot 补全代码,或是如何向 Cha…

2026/8/22 18:05:01
LaTeX公式转Word只要一次右键:LaTeX2Word-Equation插件快速上手指南

LaTeX公式转Word只要一次右键:LaTeX2Word-Equation插件快速上手指南

LaTeX公式转Word只要一次右键:LaTeX2Word-Equation插件快速上手指南 【免费下载链接】LaTeX2Word-Equation Copy LaTeX Equations as Word Equations, a Chrome Extension 项目地址: https://gitcode.com/gh_mirrors/la/LaTeX2Word-Equation 在维基百科读到一…

2026/8/22 18:05:01
小样本多类型医疗数据的机器学习建模实战

小样本多类型医疗数据的机器学习建模实战

1. 项目概述:为什么一个脑出血患者的院前指标,值得用五种机器学习模型反复“较劲”?我带过三届数学建模国赛和亚太杯的参赛队,也帮临床科室做过真实病历数据的建模支持。去年接手一个急诊科合作项目时,主任递给我一份E…

2026/8/22 18:05:01
KMS_VL_ALL_AIO 教程:一键本地 KMS 激活 Windows 和 Office,三步搞定

KMS_VL_ALL_AIO 教程:一键本地 KMS 激活 Windows 和 Office,三步搞定

KMS_VL_ALL_AIO 教程:一键本地 KMS 激活 Windows 和 Office,三步搞定 【免费下载链接】KMS_VL_ALL_AIO Smart Activation Script 项目地址: https://gitcode.com/gh_mirrors/km/KMS_VL_ALL_AIO KMS_VL_ALL_AIO 是一个单文件 KMS 激活脚本&#xf…

2026/8/22 18:05:01
程序隐藏工具鼠标中键一键触发反应超快

程序隐藏工具鼠标中键一键触发反应超快

HiddeX 今天要分享一款更强大的老板键工具——HiddeX。这款软件是俄罗斯开发者做的,功能和前面的老板来了类似,都是隐藏程序,但玩法更高级。 比老板键更强在哪里 老板键只能隐藏进程,但HiddeX能隐藏的东西更多:窗口…

2026/8/22 18:00:00