
Kafka 事务与幂等生产者跨分区原子写、读已提交与消费-处理-生产闭环1. 从一个真实业务问题说起假设你在一家电商公司做订单系统。用户支付成功后需要同时做三件事把订单状态从“待支付”改成“已支付”、扣减对应商品的库存、给用户发一条通知消息。在早期的实现里这三件事分别写不同的库、发不同的 Kafka 消息。上线之后你开始收到客诉有的订单已经显示支付成功但库存没扣导致超卖有的库存扣了订单状态却没变用户以为没付款成功。排查日志后发现问题的根因不是代码写错了而是没有把“发消息”和“改数据库”放在同一个原子操作里。数据库事务只能保证数据库内的操作一起成功或一起失败它管不了 Kafka 消息。于是出现了一个经典局面数据库提交成功后应用在发送 Kafka 消息之前崩溃消息就永远丢了或者消息发出去了但数据库事务回滚了下游却已经按“支付成功”处理了。Kafka 事务要解决的正是这个问题它让一批消息的发送以及消费位点的提交能够作为一个整体一起提交或一起回滚从而把上游的数据库操作和下游的消息投递串成一条可信的链路。读者先不要急着记 API我们先建立整体框架。2. 一句话模型与全局角色框架先记住一个最小模型Kafka 事务就是给一批跨分区的消息贴上一个统一的“事务编号”由一个专门的协调者决定它们到底算不算数消费者只有选择“读已提交”才只能看到最终算数的那部分。这句话里有三个关键角色生产者、事务协调器、消费者。生产者负责声明事务、发送消息、提交或回滚事务协调器负责给这次事务分配状态、写事务日志、在提交时通知所有涉及的分区消费者负责在读取时根据隔离级别决定“要不要看还没提交的消息”。把整体拆成三部分再看它们怎样串起来生产者(Producer) 事务协调器(TransactionCoordinator) Broker/分区 | | | | 1. 用 transactional.id 找到协调器 | | |-----------------------------------| | | | 2. 初始化事务、分配 PIDepoch | |-----------------------------------| | | 3. beginTransaction | | | 4. send 消息到多个分区 P0/P1 | | |---------------------------------------------------------------| P0 | |---------------------------------------------------------------| P1 | | 5. 把消息注册到事务中(AddPartitionsToTxn) | |-----------------------------------| | | 6. commitTransaction | | |-----------------------------------| 7. 写 PREPARE_COMMIT 到事务日志 | | | 8. 通知 P0/P1 写 COMMIT 标记 | | |----------------------------------| | 9. 返回提交成功 | | |-----------------------------------| | 消费者(隔离级别read_committed)只返回 P0/P1 上“已提交”的消息跳过未提交或已回滚的消息这张图帮助理解的是“谁在什么条件下做什么”。它不能替代真实的网络重试、协调器故障转移和多 Broker 协作细节真实机制比这张图复杂得多但骨架就是这些步骤。3. 幂等生产者事务的第一块基石在讲事务之前必须先讲幂等生产者。原因很简单如果连“同一条消息被重复发送”都搞不定谈跨分区原子写就没有意义。幂等生产者解决的是单个生产者会话内、同一分区上多条消息不因重试而重复写入的问题。它的原理是生产者启动时Broker 给它分配一个生产者编号 PIDProducer ID并对该生产者的每个分区维护一个序列号sequence number。Broker 收到消息时如果序列号不是“上一个已接受序列号 1”就拒绝或去重。这样即使网络抖动导致客户端重试Broker 也只会保留一份。开启方式非常简单但注意acks必须配合设置# 生产者端幂等配置 enable.idempotencetrue acksall retries2147483647 max.in.flight.requests.per.connection5这里最容易误解的是幂等生产者不等于“端到端精确一次”。它只保证同一个生产者会话、同一个分区内的去重。如果生产者重启PID 变了旧的重试可能和新会话的消息被当作两条不同消息所以还需要事务来跨会话、跨分区保证一致性。适用场景日志采集、指标上报这类允许少量重复但不希望重复放大写入的场景。边界跨分区、跨会话、消费-处理-生产链路单靠幂等不够。4. transactional.id把生产者变成“可恢复的事务参与者”transactional.id是事务生产者的身份标识。它的作用不是给消息打标签而是让 Broker 能够在生产者崩溃重启后用同一个身份找回之前未完成的事务状态从而完成“中断检测”和“僵尸隔离”。可以把它理解成同一个业务逻辑比如“订单服务”长期使用同一个transactional.id。当新的实例用这个 id 启动时协调器会认为旧实例是“僵尸”把它上一轮未提交的事务直接标记为中止防止旧实例半途而废的消息被后续消费者读到。配置示例transactional.idorder-service-tx-01 enable.idempotencetrue acksall transaction.timeout.ms60000这段配置放在生产者端。transaction.timeout.ms表示这个事务从开始到提交最多允许多久超过之后协调器会主动中止它。注意这个值不能超过 Broker 端transaction.max.timeout.ms否则 Broker 直接拒绝生产者会报错。这里最容易误解的是transactional.id不是随机的、也不是每次重启换一个。它需要稳定才能让协调器识别“这是同一个逻辑生产者”。如果每次启动都换 id中断恢复能力就失效了。5. 事务协调器与事务日志谁在记账事务协调器是 Kafka Broker 内部的一个组件。每一个transactional.id都会被映射到某一个 Broker 上由这个 Broker 负责该事务的协调。协调器的职责包括分配 PID 和 epoch、维护事务状态、把事务状态变更写入内部主题__transaction_state、在提交或回滚时通知所有涉及的分区。事务的生命周期大致有这些状态状态含义触发条件Empty还没有进行中的事务初始化或上一次事务已结束Ongoing事务正在进行生产者调用 beginTransaction 并发送消息PrepareCommit准备提交正在写入提交标记协调器收到 commitTransactionPrepareAbort准备回滚正在写入中止标记协调器收到 abortTransaction 或超时CompleteCommit提交完成所有分区写入提交标记成功CompleteAbort回滚完成所有分区写入中止标记成功Dead该事务 id 被中止回收长时间未活动或超时处理这张表帮助理解状态变化但它不能替代真实并发下多事务交错、协调器 leader 切换时的复杂时序。有了协调器生产者提交时并不是直接对每个分区写“提交标记”而是把“要提交”这件事委托给协调器协调器再逐个通知涉及的分区写入事务提交标记。这样即使部分分区暂时不可用协调器也会重试直到完成保证最终所有分区要么都看到提交、要么都看不到。6. 跨分区原子写的完整过程现在让一次完整请求走一遍。假设订单服务要在一个事务里向order-topic、inventory-topic、notify-topic三个不同的主题各发一条消息这三个主题各有自己的分区。时间线 T1 生产者 initTransactions() T2 beginTransaction() T3 send(order-topic, 订单已支付) T4 send(inventory-topic, 扣减SKU-1001) T5 send(notify-topic, 通知用户) T6 commitTransaction() T7 协调器写 PREPARE_COMMIT 到 __transaction_state T8 协调器通知三个分区所在 Broker 写入 COMMIT 标记 T9 三个分区都对 read_committed 消费者可见在 T3 到 T5 之间如果生产者进程崩溃协调器在超时后会把事务标记为中止三个分区上的消息都不会对read_committed消费者可见达到“一起失败”的效果。如果 T6 之后崩溃协调器会根据事务日志继续推进提交最终三个分区都可见达到“一起成功”的效果。这就是“原子写”的含义它不是真的在物理上同时写三个分区而是通过事务标记和两阶段式的提交协调让最终可见性保持一致。7. read_committed消费者才能看见的世界消息提交了但消费者默认看到的可能不是“已提交”的视图。Kafka 消费者有一个隔离级别配置取值有read_uncommitted和read_committed两种。隔离级别能读到什么适用场景read_uncommitted所有消息包括未提交、已中止事务中的消息对一致性要求不高、追求最大吞吐read_committed只读已提交事务的消息跳过未提交和中止的消息支付、订单、账务等强一致场景消费者配置isolation.levelread_committed enable.auto.commitfalse这里最容易误解的是read_committed不是“读到最新”而是“读到已提交”。它会让消费者在遇到未提交消息时等待直到该事务的提交或中止标记写入后才继续。这也意味着如果某个事务一直不结束消费者可能被卡住——这就是事务超时必须设置、协调器必须主动中止超时事务的现实原因。另一个常被忽略的点read_committed消费者读到的位点语义会跳过中止事务的消息所以不能假设消息在物理上连续。业务代码不要依赖“消息序号连续”。8. 消费-处理-生产闭环最难也是最有价值的部分前面讲的是“生产端原子写”。但真实链路往往是消费者从上游主题读一条消息处理业务逻辑再把结果写到下游主题。这就是消费-处理-生产闭环。它之所以难是因为要同时解决三件事上游位点提交、下游消息发送、业务处理三者必须一致。Kafka 事务提供了sendOffsetsToTransaction能力把消费者的位点提交也纳入同一个事务。这样下游消息和上游位点要么一起生效要么一起回滚不会出现“下游收到消息但上游位点没提交重启后重复处理”或者“上游位点提交了但下游没收到”的错位。流程草图消费者轮询上游消息 | v 开启事务 | v 处理业务逻辑可选写数据库 | v 发送下游消息 | v sendOffsetsToTransaction(上游位点) | v 提交事务 ---- 成功下游可见上游位点前进 | ---- 失败下游不可见上游位点回退重新消费这个闭环允许我们把“至少一次”的消费者和事务生产者组合逼近“精确一次”的效果前提是下游消费者也用read_committed读。9. 完整示例一最小事务生产者目标用 Java 客户端发送一批跨两个分区的消息观察提交与回滚的效果。前置环境本地启动单节点 KafkaKRaft 模式即可创建两个主题。bin/kafka-topics.sh --bootstrap-server localhost:9092--create--topictx-a--partitions2--replication-factor1bin/kafka-topics.sh --bootstrap-server localhost:9092--create--topictx-b--partitions2--replication-factor1importorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerConfig;importorg.apache.kafka.clients.producer.ProducerRecord;importorg.apache.kafka.common.serialization.StringSerializer;importjava.util.Properties;publicclassMinimalTxProducer{publicstaticvoidmain(String[]args){PropertiespropsnewProperties();props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,localhost: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.ACKS_CONFIG,all);props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,minimal-tx-demo);KafkaProducerString,StringproducernewKafkaProducer(props);producer.initTransactions();try{producer.beginTransaction();producer.send(newProducerRecord(tx-a,k1,订单已支付));producer.send(newProducerRecord(tx-b,k2,库存已扣减));producer.commitTransaction();System.out.println(事务提交成功);}catch(Exceptione){producer.abortTransaction();System.out.println(事务已回滚: e.getMessage());}finally{producer.close();}}}关键步骤initTransactions必须最先调用它负责和协调器建立身份beginTransaction开启新事务两条send分别指向不同主题commitTransaction触发协调器写提交标记。预期结果用read_committed消费者能同时看到两条消息。如果把commitTransaction换成abortTransaction两条都看不到。容易改错的地方忘记initTransactions或者transactional.id在不同运行间随意变化。10. 完整示例二消费-处理-生产闭环目标从order-topic读取订单事件处理后写到notify-topic并把上游位点提交纳入事务。前置环境主题order-topic与notify-topic已存在且notify-topic的消费者将使用read_committed。importorg.apache.kafka.clients.consumer.ConsumerConfig;importorg.apache.kafka.clients.consumer.ConsumerRecord;importorg.apache.kafka.clients.consumer.KafkaConsumer;importorg.apache.kafka.clients.producer.KafkaProducer;importorg.apache.kafka.clients.producer.ProducerConfig;importorg.apache.kafka.clients.producer.ProducerRecord;importorg.apache.kafka.common.serialization.StringDeserializer;importorg.apache.kafka.common.serialization.StringSerializer;importjava.time.Duration;importjava.util.Collections;importjava.util.Properties;publicclassConsumeProcessProduceTx{publicstaticvoidmain(String[]args){PropertiesconsumerPropsnewProperties();consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG,cpp-tx-demo-group);consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class.getName());consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG,StringDeserializer.class.getName());consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG,false);consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG,read_committed);PropertiesproducerPropsnewProperties();producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG,StringSerializer.class.getName());producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG,true);producerProps.put(ProducerConfig.ACKS_CONFIG,all);producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG,cpp-tx-demo-id);KafkaConsumerString,StringconsumernewKafkaConsumer(consumerProps);KafkaProducerString,StringproducernewKafkaProducer(producerProps);producer.initTransactions();consumer.subscribe(Collections.singletonList(order-topic));try{while(true){varrecordsconsumer.poll(Duration.ofMillis(500));if(records.isEmpty())continue;producer.beginTransaction();try{for(ConsumerRecordString,Stringrecord:records){Stringnotify用户通知: record.value();producer.send(newProducerRecord(notify-topic,record.key(),notify));}producer.sendOffsetsToTransaction(consumer.groupMetadata(),records.nextOffsets()null?Collections.emptyMap():java.util.Map.of());producer.commitTransaction();}catch(Exceptione){producer.abortTransaction();System.out.println(处理失败已回滚: e.getMessage());}}}finally{consumer.close();producer.close();}}}说明上面的位点提交部分为便于展示做了简化。生产代码中应使用consumer.committed()或直接构造MapTopicPartition, OffsetAndMetadata按每个分区“下一条要读的位点”传入。简化写法不能直接运行需要替换。关键点在于sendOffsetsToTransaction必须与下游send在同一个事务内否则位点和消息会分家。预期结果正常时下游可见通知且位点推进异常时下游不可见且位点不前进重启后重新消费。11. 事务超时与回滚别让半成品挂太久事务不是永久有效的。生产者一旦开启事务却迟迟不提交协调器会按transaction.timeout.ms主动中止它。这个设计是必要的因为如果某个生产者卡死或网络分区未完成的事务会一直让read_committed消费者阻塞在那些分区上。超时与回滚相关的关键配置配置位置作用常见取值transaction.timeout.ms生产者事务最长持续时间60000transaction.max.timeout.msBroker允许的最大事务超时900000transactional.id.expiration.msBroker事务 id 多久未活动被清理604800000enable.idempotence生产者幂等与事务的前置true回滚的触发路径有三条业务代码显式abortTransaction、异常捕获后回滚、协调器超时中止。三者都会在事务日志和分区上写入中止标记对read_committed消费者表现为“这些消息不存在”。注意回滚不是删除消息而是让它们在read_committed视图中不可见物理消息仍会保留直到日志留存策略清理。12. 设计取舍哪些场景该用事务哪些不该事务带来一致性也带来成本。事务提交需要写事务日志、写分区标记、协调器与分区多次交互延迟高于普通生产read_committed消费者在遇到未提交消息时会等待吞吐也会受影响。条件化结论当链路包含“数据库写入 多主题消息发送”且业务不能容忍不一致时选择事务当只是单主题日志采集、允许少量重复且不需要原子跨分区时选择幂等生产者即可当下游是可容忍重复的离线分析链路时普通生产加幂等去重反而更简单。还需要权衡的是事务只覆盖 Kafka 内部操作如果你在事务里还写了外部数据库那这个数据库写入不在Kafka 事务范围内需要业务自己用本地消息表或幂等键配合否则仍会出现不一致。这是事务边界最容易踩的坑。13. 生产实践建议第一transactional.id要稳定且与业务实例对应推荐按“服务名-业务域”命名不要用随机值。第二transaction.timeout.ms要小于 Broker 的transaction.max.timeout.ms并留出业务处理时间的余量。第三事务内不要做耗时外部调用事务越长协调器压力和消费者阻塞风险越大。第四read_committed消费端要设置合理的fetch.max.wait.ms和max.poll.records避免单次处理太久导致位点事务迟迟不提交。第五监控__transaction_state所在 Broker 的负载协调器热点会拖慢所有相关事务。14. 常见误区这里集中纠正几类高频误解。第一“开了事务就不会重复消费”。事务保证的是发送和位点提交的原子性但消费者如果处理失败重试仍可能重复处理业务必须做幂等。第二“幂等生产者等于事务”。两者是不同层级幂等只覆盖单分区去重事务才覆盖跨分区原子性。第三“read_committed 就是全局一致”。它只保证读不到未提交和已中止的消息不保证读到的消息一定最新。第四“事务回滚会删掉消息”。消息仍在物理日志里只是对read_committed不可见。第五“事务可以覆盖数据库”。Kafka 事务不管理外部数据库事务跨系统一致性要自己设计。15. 排障清单当遇到事务相关问题时按顺序检查生产者是否调用了initTransactions()且transactional.id是否稳定。生产者的target.timeout.ms与 Broker 的transaction.max.timeout.ms是否冲突日志中是否出现InvalidTxnTimeoutException。是否出现ProducerFencedException这通常意味着同一transactional.id有新旧实例同时运行旧实例已被隔离。消费者isolation.level是否为read_committed否则会读到未提交消息。__transaction_state主题是否健康协调器所在 Broker 是否负载过高。事务是否长时间未提交导致read_committed消费者在某个分区上卡住。排查时先用kafka-transactions.sh或 Broker 日志确认协调器状态再结合生产者客户端日志中的TransactionalId和ProducerId定位。16. 面试与复盘问题可以拿这几个问题检验自己是否真的理解为什么幂等生产者不能替代事务事务协调器为什么需要有状态read_committed消费者遇到未提交事务时会等待多久sendOffsetsToTransaction为什么必须和send在同一个事务内当transactional.id相同的新实例启动时旧实例的消息会怎样事务超时后协调器做了哪些动作把这些问题回答清楚基本就掌握了 Kafka 事务的核心脉络。17. 总结回到最初的框架Kafka 事务由生产者、事务协调器、事务日志、分区提交标记和消费者隔离级别共同构成。幂等生产者保证单分区去重transactional.id让事务可识别、可恢复协调器负责状态管理和提交协调跨分区原子写通过事务标记实现最终一致的可见性read_committed决定消费者看到的世界消费-处理-生产闭环把上游位点和下游消息绑在一起。事务超时与回滚则划定了这套机制的安全边界。工程上记住一句话事务解决的是 Kafka 内部的原子可见性业务幂等和外部系统一致性仍要自己负责。把这句话落实到配置、代码和监控上才算真正用对了 Kafka 事务。18. 参考资料Apache Kafka 官方文档Transactional Producer 与 Transactions 章节。Apache Kafka 官方文档Idempotent Producer 与 Producer Configs 章节。Apache Kafka 官方文档Consumer Configs 中isolation.level说明。Apache Kafka 官方文档KRaft 模式与__transaction_state内部主题说明。《Kafka: The Definitive Guide》相关事务与幂等章节。