当前位置: 首页 > news >正文

把 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 在「功能丰富」(延迟消息、事务、顺序、重试)上更全,但极端吞吐略逊。

维度KafkaRocketMQ
默认可靠性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,这条半消息会怎样?你会在业务侧怎么兜底?

http://www.cnnetsun.cn/news/4159957.html

相关文章:

  • KMS_VL_ALL_AIO 本地KMS快速激活指南:从零到一,一个批处理搞定 Windows 和 Office 激活
  • ol-ext 上手实操手册:OpenLayers 地图扩展库核心能力拆解
  • `import fnmatch` 是 Python 中导入标准库模块 `fnmatch` 的语句
  • 主题公园移动供电案例:从环球影城场景,看户外频繁插拔工况下工业连接器选型思路
  • 读懂eas.json:expo-react-native-cicd中dev、prod-apk、prod-aab三大构建Profile配置详解
  • PyULog:8条命令解析PX4 ULog日志,导出CSV、KML与SQLite
  • RDMA数据传输操作:Send/Recv与Read/Write全解析
  • AnythingLLM 教程:10 分钟搭建一个本地私有知识库问答应用
  • Element Tiptap富文本编辑器:Vue3项目5分钟接入带菜单的WYSIWYG编辑器
  • SPA 刷新 404 难题终结者:boot-react SinglePageAppConfig pushState 资源解析器深度剖析
  • 我的价值观
  • 如何使用 draw.io 桌面版:离线绘图与批量导出完整指南
  • CEdev 图形编程完全指南:graphx 库调色板、精灵动画与 Tilemap 实战教程
  • iOS跨平台位置模拟实战:基于WebKit调试协议实现GeoPort方案
  • Weasis:内建 2D/3D 影像分析的开源 DICOM 查看器
  • res-downloader完全教程:免费的跨平台资源嗅探器,一键下载视频音乐图片
  • PhoneProfilesPlus新手必学的8个实用场景:会议自动静音、通勤一键飞行模式
  • 让Claude Code、Codex与Gemini协同工作:Agent Relay Harnesses完整指南
  • 训练UniDetector前必看的20+个关键超参数:完整配置项逐条解读
  • 快速上手solid-dnd:10分钟从零搭建你的第一个拖拽应用,新手友好教程
  • 如何系统掌握高级数据结构?AlgorithmsAndDataStructuresInAction官方代码库入门指南
  • Vue-preview 图片预览:新手安装与上手完整指南
  • ShawzinBot:免费把 MIDI 变成游戏按键
  • 如何重置 Navicat 试用期:3 条命令跑通 navicat-key 注册表清理工具
  • ODC 生产环境部署最佳实践:MetaDB、Docker 与高可用架构配置全解析
  • 让联邦查询提速10倍:aws-athena-query-federation谓词下推、分区裁剪与TopN优化实战
  • MobilityDB高精度建模:tpose四元数姿态类型如何描述自动驾驶与机器人运动
  • TeslaLogger新功能MCP Server详解:用自然语言向AI查询你的特斯拉数据
  • B站视频下载完全指南:5 分钟跑通 BilibiliDown,把喜欢的内容存进本地
  • Rollbar.js Node.js 接入实战:Express 服务端错误追踪的 5 步配置法