我把 Kafka 消费改成 EOS 事务后,资损从月均 30 万压到 0:幂等键 + Transactional Producer 实战
我把 Kafka 消费改成 EOS 事务后资损从月均 30 万压到 0幂等键 Transactional Producer 实战凌晨 3 点 12 分财务给我打电话说上个月对账差了 31 万。我当时第一反应是数据库出问题了结果查了一晚上发现不是数据库是 Kafka。订单事件被重复消费了。每条订单事件平均被下游 3 个服务消费有的服务 at-least-once 业务没做幂等导致同一个扣款事件触发了两次。这事儿说实话不稀奇但金额一上来就是真金白银。这篇文章我把我后来怎么把 Kafka 消费改成 EOSExactly-Once Semantics事务的完整过程写下来包括代码、踩坑和最终效果对比。如果你也在做金融/电商类业务并且对资损零容忍那这篇文章可能对你有点用。问题到底出在哪我们订单链路的拓扑大致是这样订单服务 → Kafka topic order.created → [3 个下游服务] ├─ 账务服务扣款 ├─ 库存服务扣库存 └─ 营销服务发券消费模式是默认的 at-least-onceKafkaListener(topicsorder.created)publicvoidonMessage(ConsumerRecordString,Stringrec){OrderEventevtparse(rec.value());accountService.deduct(evt.userId,evt.amount);// 扣款}这套代码的问题我一开始没看出来因为单条消息是幂等的——同一个订单号扣两次业务层有兜底金额校验。但问题在于下游回写 重平衡 重投递这三种场景叠在一起时幂等键会失效。举个真实的例子账务服务在扣款成功后写了一个本地事务表但这个事务表写完的瞬间Kafka 消费者 offset 还没提交broker 就把 consumer 踢出了消费组rebalance。等这个 consumer 重新加入时会从上一个提交的 offset 重新拉这条消息然后再次调accountService.deduct但这次因为某些原因具体后面讲业务幂等键没生效。月均资损 30 万就是这么来的。为什么我一开始没选 EOSKafka 的 EOS 不是什么新东西从 0.11 版本就有了。但我之前一直没上主要有 3 个顾虑性能损耗。开启transactional.id和enable.idempotence后Producer 吞吐会下降社区里传的是 10%~30%。我们订单 topic 的峰值 QPS 2.5 万开 EOS 我担心被打爆。复杂度。Transactional API 用起来比普通 Producer 复杂得多光initTransactions()、beginTransaction()、commitTransaction()这三个方法就要在每个 Producer 端写对位置写错一个就数据不一致。生态兼容性。我们下游有 3 个语言栈Java/Go/PythonEOS 必须 Producer Consumer 端都用 Kafka 原生客户端才能严格保证 exactly-once跨语言理论上能互通但有坑。后面我重新算了一笔账性能损耗可以靠扩容消化钱能解决的问题都不是问题跨语言我们 3 个服务其实都用官方库剩下就是复杂度——这个只能硬啃。落地方案四步走我把改造拆成了四步每步都有可验证的产物。第一步Producer 端上事务先在订单服务加 Transactional Producer。这里有个关键点transactional.id必须每个 Producer 实例唯一并且重启后保持不变——Kafka 用它来检测僵尸实例。PropertiespropsnewProperties();props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,kafka-1:9092,kafka-2:9092);props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG,true);props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,order-tx-hostname-pid);props.put(ProducerConfig.ACKS_CONFIG,all);props.put(ProducerConfig.RETRIES_CONFIG,Integer.MAX_VALUE);props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION,5);KafkaProducerString,StringproducernewKafkaProducer(props);producer.initTransactions();// 必须在第一次发消息前调用发消息时用sendOffsetsToTransaction把消费 offset 一起提交到事务里这样消费 生产就是原子的try{producer.beginTransaction();// 业务逻辑扣款 写本地事务表accountService.deduct(evt.userId,evt.amount);txLogRepo.markProcessed(evt.orderId);// 发到下游 topicproducer.send(newProducerRecord(order.deducted,evt.userId,evt.toJson()));// 把消费 offset 也提交进这个事务MapTopicPartition,OffsetAndMetadataoffsetsnewHashMap();offsets.put(newTopicPartition(rec.topic(),rec.partition()),newOffsetAndMetadata(rec.offset()1));producer.sendOffsetsToTransaction(offsets,consumerGroupMetadata);producer.commitTransaction();}catch(Exceptione){producer.abortTransaction();throwe;// 让 consumer 重试}注意一个细节sendOffsetsToTransaction第二个参数是consumerGroupMetadata必须从当前正在用的 consumer取不能用上一次的值。我们一开始图省事复用了一个静态的 metadata结果事务提交后 offset 没前进又重投了一次。第二步Consumer 端读已提交消息Consumer 端必须设置isolation.levelread_committed否则会读到未提交的事务消息。props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG,read_committed);这个参数加上之后Consumer 会自动跳过未提交的事务和已中止的事务。看起来很简单但这步让我踩了一个大坑——后面踩坑记录里讲。第三步业务层补幂等键EOS 在 Kafka 链路内是 exactly-once但跨服务的边界HTTP/RPC 调用、数据库写还是 at-least-once。事务只能保证 Kafka 内部的原子性业务层的幂等必须自己做。我们的账务服务幂等键是这样设计的CREATETABLEaccount_tx(idBIGINTPRIMARYKEYAUTO_INCREMENT,order_idVARCHAR(64)NOTNULL,user_idVARCHAR(64)NOTNULL,amountDECIMAL(18,2)NOTNULL,statusTINYINTNOTNULLDEFAULT0,-- 0pending, 1successcreated_atDATETIMEDEFAULTCURRENT_TIMESTAMP,UNIQUEKEYuk_order_id(order_id))ENGINEInnoDB;扣款流程Transactionalpublicvoiddeduct(StringuserId,StringorderId,BigDecimalamount){// 1. 查事务表OptionalAccountTxexistingtxRepo.findByOrderId(orderId);if(existing.isPresent()existing.get().getStatus()1){return;// 已处理过直接返回}// 2. 扣余额accountRepo.deduct(userId,amount);// 3. 写事务表如果重复插入会抛唯一约束异常if(existing.isEmpty()){txRepo.insert(userId,orderId,amount);}else{txRepo.markSuccess(orderId);}}这里UNIQUE KEY uk_order_id是关键——MySQL 唯一约束会兜底所有并发场景业务层不需要再加分布式锁。第四步监控 对账EOS 上完之后资损能不能真的压到 0得靠对账来验证不是靠感觉。我们加了一个每日凌晨的对账任务# 每天 02:00 跑defreconcile():kafka_countquery_kafka_offset(order.deducted,groupreconcile-group)db_countquery_db(SELECT COUNT(*) FROM account_tx WHERE status1 AND created_at CURDATE())diffkafka_count-db_countifabs(diff)100:# 允许 100 条以内延迟alert(订单对账差异过大,diffdiff)这个对账任务一周内帮我们抓到了 3 次事务提交失败但没报警的场景都是commitTransaction()抛了ProducerFencedException但被 catch 后吞掉了。踩坑记录这 4 个坑真的别踩坑 1transactional.id 用 IP 拼的重启后僵尸锁死我一开始图省事transactional.id直接用本机 IPprops.put(TRANSACTIONAL_ID_CONFIG,tx-InetAddress.getLocalHost().getHostAddress());结果容器编排里 Pod IP 变了旧的 IP 对应的事务状态还在 broker 端新 Producer 用相同的 ID 一启动就被 fencedProducerFencedException。正确做法是用稳定的 hostname pid或者直接用UUID.randomUUID()但写到一个持久化目录里。我后来改成用 StatefulSet 的 pod name 加 sequence稳了。坑 2read_committed 配上 auto.offset.resetearliest 把整库重放了一遍我同事在测试环境改了isolation.levelread_committed但忘了改auto.offset.reset结果重启后从最早的消息开始读把一个月前的订单事件全部重放。账务服务因为有幂等键没出事但营销服务发券的接口没幂等那时候还没补一下发出去了 2 万张重复券直接被业务部门找上门。教训改 Consumer 配置前先确认 group 的 offset 位置或者干脆用一个新的 consumer group。坑 3事务里调用了慢 RPC撑爆 transaction.timeout.ms我一开始把扣款逻辑整个包在事务里producer.beginTransaction();accountService.deduct(...);// 同步 HTTP 调用riskService.checkFraud(...);// 同步 HTTP 调用可能 800msproducer.send(...);producer.commitTransaction();结果风险检查超时整个事务被 abortoffset 没提交下次重投再走一遍风险检查永远卡死。正确做法业务逻辑和事务边界要分开。事务里只放必须原子提交的部分写数据库 发送到 Kafka慢 RPC 在事务外做完再进事务。// 先做慢 RPCbooleanfraudOkriskService.checkFraud(evt);// 再进事务producer.beginTransaction();accountService.deduct(...);producer.send(...);producer.commitTransaction();坑 4commitTransaction 抛异常被 catch 后没 aborttry{producer.commitTransaction();}catch(Exceptione){log.error(提交失败,e);// 忘了 abort}如果 commit 失败但没 abort这个 Producer 实例会卡死——下次 beginTransaction 会抛IllegalStateException。正确做法try{producer.commitTransaction();}catch(Exceptione){producer.abortTransaction();// 一定要 abortthrowe;}最终效果改造前后对比QPS 峰值没掉反而因为减少了重复消息处理broker I/O 下降了一点指标改造前改造后变化月均重复扣款笔数42000-100%月均资损金额~31 万0-100%Kafka 集群写入吞吐2.5 万 QPS2.3 万 QPS-8%P99 订单处理延迟85ms92ms8%Consumer CPU 使用率45%52%7%吞吐和延迟的损耗完全在可接受范围内——我们当时还留了 50% 容量余量这点损耗对账务服务的影响基本可以忽略。但最关键的指标是月度对账差异改造前每月 30 万左右改造后连续 4 个月差异都是 0。财务再没半夜给我打过电话。写在最后如果你问我 EOS 是不是银弹我的答案是不是。它解决的是 Kafka 链路内的 exactly-once但跨服务、跨存储系统的边界还是需要业务层幂等来兜底。我们最终方案是 “EOS 业务幂等键 每日对账” 三件套少哪个都不行。另外 EOS 的复杂度是真的高团队里至少要有 1-2 个对 Kafka 内部机制比较熟的人才能 hold 住。如果你们业务对资损没那么敏感比如只是日志、推荐特征这类其实用 at-least-once 业务幂等键就够了没必要硬上 EOS。最后说一句可能得罪人的话很多团队出问题不是技术选型错了是监控没跟上。我们最初出问题的时候资损了 3 个月才发现。如果对账任务早点上可能根本不需要上 EOS。行了今天就写到这。如果你有 Kafka 相关的踩坑经历欢迎评论区聊聊。

相关新闻

进程互斥和进程同步

进程互斥和进程同步

一、进程互斥由于进程具有独立性和异步性等并发特征,计算机的资源有限,导致了进程之间的资源竞争和共享,也导致了对进程执行过程的制约。1、临界资源和临界区(临界部分)临界资源:一次只能供一个进程访问的资…

2026/7/28 17:06:45 阅读更多 →
我用 Nix 重写 Docker 镜像构建,把镜像从 1.2GB 压到 80MB 且可复现

我用 Nix 重写 Docker 镜像构建,把镜像从 1.2GB 压到 80MB 且可复现

我用 Nix 重写 Docker 镜像构建,把镜像从 1.2GB 压到 80MB 且可复现 说实话,我一开始是拒绝用 Nix 的。 那天我们项目的 Node.js 镜像在生产环境又出问题了——同样是 node:20-slim 基础镜像,同样是 package.json 没变,但 CI 流水…

2026/7/28 17:06:45 阅读更多 →
Spring Boot + Shiro 等保三级复测实战:12行代码修复高危漏洞

Spring Boot + Shiro 等保三级复测实战:12行代码修复高危漏洞

1. 项目背景与核心挑战最近,我们团队负责维护的一套基于Spring Boot和Shiro的医疗信息系统,迎来了等保三级(网络安全等级保护第三级)的年度复测。对于非技术出身的同事可能不太清楚,等保三级是国内非银行机构的最高安全…

2026/7/28 17:05:45 阅读更多 →

最新新闻

MongoDB 4.2——在生产环境中设置 MongoDB

MongoDB 4.2——在生产环境中设置 MongoDB

在生产环境中设置 MongoDB 1、从命令行启动2、停止MongoDB3、安全性3.1、数据加密3.2、SSL连接 4、日志 1、从命令行启动 MongoDB 服务器是以 mongod 可执行文件来启动的。mongod 有很多可配置的启动选项,要查看这些选项,可以在命令行中运行 mongod --h…

2026/7/28 17:15:47 阅读更多 →
Ubuntu smba 重启

Ubuntu smba 重启

sudo /etc/init.d/samba restart

2026/7/28 17:15:47 阅读更多 →
MongoDB 4.2——持久性

MongoDB 4.2——持久性

持久性1、使用日志机制的成员级别持久性2、使用写关注的集群级别持久性2.1、writeConcern的w和wtimeout选项2.2、writeConcern的j(日志)选项3、使用读关注的集群级别持久性4、使用写关注的事务持久性5、MongoDB不能保证什么6、检查数据损坏1、使用日志机…

2026/7/28 17:15:47 阅读更多 →
EncodingChecker:3步解决文件乱码问题的终极指南

EncodingChecker:3步解决文件乱码问题的终极指南

EncodingChecker:3步解决文件乱码问题的终极指南 【免费下载链接】EncodingChecker A GUI tool that allows you to validate the text encoding of one or more files. Modified from https://encodingchecker.codeplex.com/ 项目地址: https://gitcode.com/gh_m…

2026/7/28 17:15:47 阅读更多 →
springMVC定义拦截器判断用户是否为管理员

springMVC定义拦截器判断用户是否为管理员

首先是数据库的设计CREATE table tb_user( uId int(11) primary key auto_increment,uName varchar(20) not null,uPassword varchar(20) not null,uPhone varchar(30) not null,uAddress VARCHAR(100) not null,isManager int(1) default 0 -- 1:为管理员 0:不为…

2026/7/28 17:15:47 阅读更多 →
终极指南:浏览器端离线语音识别技术深度解析与Vosk-Browser实战

终极指南:浏览器端离线语音识别技术深度解析与Vosk-Browser实战

终极指南:浏览器端离线语音识别技术深度解析与Vosk-Browser实战 【免费下载链接】vosk-browser A speech recognition library running in the browser thanks to a WebAssembly build of Vosk 项目地址: https://gitcode.com/gh_mirrors/vo/vosk-browser 在…

2026/7/28 17:14:47 阅读更多 →

日新闻

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生

告别臃肿!3步让你的暗影精灵笔记本重获新生 【免费下载链接】OmenSuperHub Control Omen laptop performance, fan speeds, and keyboard lighting, and unlock power limits. 项目地址: https://gitcode.com/gh_mirrors/om/OmenSuperHub 你是否也曾为官方Om…

2026/7/28 0:00:43 阅读更多 →
RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

RAG必踩坑!财报法规检索不准?这款开源工具让答案浮出水面,准确率飙升98.7%!

做 RAG 的人应该都踩过这个致命的坑:把几百页的财报、法规、技术手册扔给向量库,问一个具体问题,搜出来的全是沾边但没用的内容 —— 关键信息要么被硬切块拆碎了,要么藏在几十条结果的最下面。语义相似≠真正相关,这个…

2026/7/28 0:00:43 阅读更多 →
抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

抖音视频文案提取工具全指南:免费2026版、手机App、在线工具一网打尽

2026年做短视频运营,从抖音上扒文案早就不是偷偷抄笔记的事了。我刚开始做内容的时候,每天刷半小时抖音,手动把爆款视频的口播敲进备忘录,一条2分钟的视频得花十来分钟,碰到语速快的还要反复回听。后来试了一圈工具&am…

2026/7/28 0:00:43 阅读更多 →

周新闻

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 道路桥梁裂缝检测数据集 道路桥梁病害识别检测数据集

深度学习道路桥梁裂缝检测系统 数据集6000张 完整源码已标注数据集训练好的模型环境配置教程程序运行说明文档,可以直接使用!系统支持图片、视频、摄像头等多种方式检测裂缝,功能强大实用。 1数据集6000张 8各类别

2026/7/28 12:04:22 阅读更多 →
深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

深度学习YOLO模型如何训练 PUBG 绝地求生目标检测数据集

pubg数据集 精选原图1.42万数据 1.49万标签 无任何重复、算法增强或冗余图像! pubg绝地求生目标检测数据集 1分类:e_body,14905个标签,txt格式 共计14244张图,99%为640*640尺寸图像 适合yolo目标检测、AI训练关键词&am…

2026/7/28 8:29:16 阅读更多 →
Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex英雄目标检测数据集 深度学习框架YOLO如何训练APEX数据集

Apex检测数据集数据集详情检测类别: allies enemy tag图片总量:7247张训练集:5139张验证集:1425张测试集:683张标注状态:全部已标注,即拿即用数据格式:支持YOLO格式及其他格式&#…

2026/7/28 5:03:42 阅读更多 →

月新闻