我把 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 相关的踩坑经历欢迎评论区聊聊。