
Kafka 事务消息实现详解Kafka 从 0.11 版本开始引入事务配合幂等生产者实现了Exactly-Once语义。一、为什么需要事务1.1 典型的消费-处理-生产场景Consumer ──读取──► Topic A ──处理──► 写入 Topic B如果处理成功但写入失败或者写入成功但 offset 未提交就会产生重复消费写完 B 后崩溃重启后再次消费同一条 A 的消息消息丢失处理完 offset 已提交但 B 写入失败1.2 事务解决了什么Kafka 事务提供的是原子多分区写入 消费位移提交保证要么全部成功消息写入 B offset 提交 要么全部失败消息不写入 offset 不回滚保持原样不是传统数据库的 ACID 事务——不涉及多行回滚、隔离级别等概念。二、核心概念2.1 四个关键角色┌──────────────────────────────────────────────────────────┐ │ │ │ ┌──────────┐ ┌──────────────────┐ ┌─────────────┐ │ │ │ Producer │ │ Transaction │ │ __transaction│ │ │ │ (事务ID) │ │ Coordinator │ │ _state │ │ │ │ │ │ (TC, Broker上) │ │ (内部Topic) │ │ │ └──────────┘ └──────────────────┘ └─────────────┘ │ │ │ │ │ │ │ │ 注册/提交 │ 持久化状态 │ │ │ ├────────────────►├─────────────────────►│ │ │ │ │ │ │ │ ┌──────────┐ │ │ │ │ │ Consumer │ │ │ │ │ │ (隔离级别)│ │ │ │ │ └──────────┘ │ │ │ │ │ └──────────────────────────────────────────────────────────┘角色说明Transactional Producer配置transactional.id的生产者Transaction Coordinator (TC)运行在 Broker 上的事务协调器每个transactional.id会 hash 到固定的 TC__transaction_state内部 Topic50 分区TC 将事务状态持久化到这里Consumer (isolation.level)消费者通过隔离级别控制是否读取未提交的事务消息2.2 关键 IDID含义生命周期transactional.id用户配置的稳定事务标识跨 Producer 重启不变PID(Producer ID)Broker 分配的 64 位数字Producer 重启后重新分配但transactional.id不变时会复用旧的 epoch 过期Producer Epoch单调递增区分同 PID 的不同实例每次 Producer 重新 initTransactions 时 1Transactional Sequence Number每个分区内的消息序号幂等生产的基础transactional.id: order-service-tx-001 ← 用户配置稳定 ↓ initTransactions 请求 PID: 123456789, Epoch: 0 ← Broker 分配 ↓ Producer 重启再次 initTransactions PID: 123456789, Epoch: 1 ← PID 不变Epoch 递增三、事务完整流程时间线3.1 创建事务生产者PropertiespropsnewProperties();props.put(bootstrap.servers,localhost:9092);props.put(transactional.id,order-tx-001);// ← 关键稳定的事务 IDprops.put(enable.idempotence,true);// ← 事务自动开启幂等// 幂等要求acksall, retries0, max.in.flight.requests.per.connection≤5KafkaProducerString,StringproducernewKafkaProducer(props);3.2 完整流程图Producer Transaction Coordinator __transaction_state │ │ │ │ ① initTransactions() │ │ ├───────────────────────────────►│ FindCoordinator │ │ │ 分配 PID Epoch │ │◄───────────────────────────────┤ │ │ ② beginTransaction() │ │ │ (仅本地标记不发网络请求) │ │ │ │ │ │ ③ send(topic-A, msg) │ │ ├───────────────────────────────►│ 消息写入分区但标记为未提交 │ │ │ │ │ ④ send(topic-B, msg) │ │ ├───────────────────────────────►│ 同上 │ │ │ │ │ ⑤ sendOffsetsToTransaction() │ │ ├───────────────────────────────►│ 将消费位移也纳入事务 │ │ │ │ │ ⑥ commitTransaction() │ │ ├───────────────────────────────►│ 写入 PrepareCommit │ │ ├─────────────────────────────►│ │ │ 写入 Committed │ │ ├─────────────────────────────►│ │ │ 向涉及的分区 Leader 发送 │ │ │ TransactionMarker(COMMIT) │ │ │ │ │◄───────────────────────────────┤ 返回成功 │3.3 Java 完整示例// 事务生产者 PropertiespropsnewProperties();props.put(bootstrap.servers,localhost:9092);props.put(transactional.id,order-tx-001);// enable.idempotence 在配置 transactional.id 时自动为 trueKafkaProducerString,StringproducernewKafkaProducer(props);producer.initTransactions();// ① 初始化注册 PIDtry{producer.beginTransaction();// ② 开启事务// ③ 发送业务消息producer.send(newProducerRecord(order-topic,order-123,created));// ④ 将消费位移也纳入事务原子绑定MapTopicPartition,OffsetAndMetadataoffsetsnewHashMap();offsets.put(newTopicPartition(source-topic,0),newOffsetAndMetadata(100));producer.sendOffsetsToTransaction(offsets,consumer-group-1);producer.commitTransaction();// ⑤ 提交}catch(ProducerFencedExceptione){// PID 被 epoch 更新的实例抢占当前实例已僵尸producer.close();}catch(KafkaExceptione){producer.abortTransaction();// ⑥ 异常时回滚}3.4 消费者端PropertiesconsumerPropsnewProperties();consumerProps.put(bootstrap.servers,localhost:9092);consumerProps.put(group.id,consumer-group-1);consumerProps.put(isolation.level,read_committed);// ← 关键配置// read_committed: 只读已提交的事务消息默认是 read_uncommittedKafkaConsumerString,StringconsumernewKafkaConsumer(consumerProps);四、内部实现原理4.1 事务消息在分区中的物理形态Partition 0 中的消息序列逻辑视图 ┌──────────────────────────────────────────────────────────────┐ │ offset │ key │ value │ 事务标记 │ ├────────┼───────────┼─────────────┼───────────────────────────┤ │ 0 │ order-123 │ created │ (无) │ │ 1 │ order-456 │ created │ (无) │ │ 2 │ order-789 │ created │ ┌─ PID100, Epoch0 │ │ 3 │ order-789 │ paid │ │ 事务未提交 │ ← read_committed 读不到 │ 4 │ │ COMMIT │ └─ TransactionMarker │ ← 控制消息消费者不可见 │ 5 │ order-999 │ created │ ┌─ PID100, Epoch0 │ │ 6 │ │ ABORT │ └─ TransactionMarker │ ← 回滚标记 │ 7 │ order-111 │ created │ (无) │ └──────────────────────────────────────────────────────────────┘关键点事务消息照常写入分区和普通消息混在一起区别在于消息头部带有事务元信息PID、Epoch、Sequence Number事务结束时TC 会在每个涉及的分区末尾追加一条TransactionMarker控制消息Consumerread_committed模式读到 ABORT 标记后会跳过该事务的所有消息4.2 __transaction_state 内部 Topic┌──────────────────────────┐ │ __transaction_state │ │ 分区数: 50 (固定) │ │ 副本数: 3 (建议) │ │ 压缩策略: compaction │ ├──────────────────────────┤ │ Key: transactional.id │ │ Value: 事务状态消息 │ │ - Empty / Ongoing │ │ - PrepareCommit │ │ - PrepareAbort │ │ - CompleteCommit │ │ - CompleteAbort │ └──────────────────────────┘TC 启动时通过读取__transaction_state恢复所有处于Ongoing状态的事务然后继续推进。4.3 事务提交流程深入Producer Transaction 分区 Leader │ Coordinator │ │ commitTransaction() │ │ ├───────────────────────►│ │ │ │ 1. 写 PrepareCommit │ │ │ 到 __transaction │ │ │ _state │ │ │ │ │ │ 2. 对各分区 Leader │ │ │ 发送 WriteTxnMarker│ │ ├──────────────────────►│ │ │ │ 3. 追加 COMMIT │ │ │ Control Message │ │◄──────────────────────┤ │ │ │ │ │ 4. 写 Committed 到 │ │ │ __transaction_state │ │◄───────────────────────┤ │ │ 返回成功 │ │两阶段提交不完全是。Kafka 事务是1.5 阶段提交——TC 先持久化PrepareCommit然后并行发 Marker 到各分区全部成功后再持久化Committed。如果 TC 在中间宕机重启后从__transaction_state恢复重试发 Marker。4.4 事务回滚producer.abortTransaction();流程与提交对称TC 写PrepareAbort到__transaction_state向各分区发送 ABORT TransactionMarkerread_committed消费者读到 ABORT 后跳过该事务的消息仿佛从未发送五、Consumer 隔离级别5.1 read_uncommitted默认读到所有消息包括未提交和已回滚的事务消息 → 可能读到最终被回滚的脏数据5.2 read_committed只读到已提交事务的消息 → 实现 Exactly-Once 的前提 Consumer 会在内存中维护一个已中止事务的缓存 ┌──────────────────────────────────┐ │ 正在消费 offset5 │ │ 读到 ABORT TransactionMarker │ │ → 把 (PID100, Epoch0) 记入缓存 │ │ → 回溯跳过该事务的所有消息 │ │ offset 2, 3 │ │ → 继续从 offset7 消费 │ └──────────────────────────────────┘六、幂等生产者 (Idempotent Producer) —— 事务的前提事务依赖于幂等生产幂等又是独立可用的功能。6.1 幂等原理普通生产者 send(msg) → 网络超时 → 重试 → Broker 收到两条 msg重复 幂等生产者 send(msg, PID100, Seq5) → 网络超时 → 重试 send(msg, PID100, Seq5) Broker 收到第二条时 PID100 的 Seq5 我已经有了忽略 → 只保留一条Broker 端每个分区维护PID → 最近 5 个 Seq 的去重窗口 ┌──────┬─────────────────────────┐ │ PID │ 已收到的 Seq 号 │ ├──────┼─────────────────────────┤ │ 100 │ [1, 2, 3, 4, 5] │ │ 101 │ [1, 2, 3] │ └──────┴─────────────────────────┘6.2 幂等 vs 事务对比幂等事务开启条件enable.idempotencetruetransactional.id作用域单分区内消息不重复跨分区原子写入跨分区原子性❌✅原子绑定消费位移❌✅sendOffsetsToTransactionPID 生成由 Broker 随机分配关联到transactional.id可恢复七、关键配置清单Producer# 事务 ID设为非空即开启事务支持 transactional.idmy-app-tx-001 # 事务超时默认 60 秒最大 15 分钟 transaction.timeout.ms60000 # 幂等事务自动开启无需显式设置 enable.idempotencetrue # 幂等要求自动设置 acksall retries2147483647 # Integer.MAX_VALUE max.in.flight.requests.per.connection5Consumer# 隔离级别 isolation.levelread_committed # read_uncommitted | read_committed # 会话超时必须大于事务超时 session.timeout.ms transaction.timeout.msBroker# 事务状态日志副本数 transaction.state.log.replication.factor3 # 事务状态日志分区数默认 50不可通过配置修改 transaction.state.log.num.partitions50 # 事务状态日志分段大小 transaction.state.log.segment.bytes104857600 # 事务超时最大值 transaction.max.timeout.ms900000 # 15 分钟八、僵尸实例隔离 —— Producer Fencing场景 1. Producer-A (Epoch0) 开启事务网络分区 2. Producer-A 被认为宕机 3. Producer-A 重启 → initTransactions → PID 复用Epoch1 4. 老的 Producer-A (Epoch0) 网络恢复尝试 commitTransaction Broker 收到 Epoch0 的提交请求 你的 Epoch 已经过期当前 Epoch1拒绝 → ProducerFencedException → 防止僵尸实例写入脏数据这就是transactional.id必须稳定的原因同一个 ID 恢复后 PID 不变但 Epoch 递增旧的 Epoch 的所有操作自动失效。九、适用场景与局限性适用场景场景示例消费-转换-生产从 topic-A 读 → 处理 → 写 topic-B 提交 offset多分区原子写入订单创建同时写订单主题和通知主题Kafka Streams内部大量使用事务保证 exactly-once局限性局限说明不跨系统只能保证 Kafka 内部的原子性不涉及数据库、Redis 等外部系统性能开销事务提交有额外网络往返 __transaction_state写入吞吐下降约 10-20%消费者需配合必须isolation.levelread_committed且消费者会缓冲未提交消息超时限制默认 60 秒最长 15 分钟不适合长事务不覆盖 consumer 的非 Kafka 操作如果在消费-处理-生产中间写了数据库数据库写入不受事务保护十、常见问题Q1: commitTransaction 没收到响应是成功还是失败调用producer.commitTransaction()时如果超时或抛异常不要重试 commit。正确做法是producer.close()然后重建 Producer。Broker 会自行完成或超时回滚。Q2: 同一个 transactional.id 能多实例并发吗绝对不能。同一时刻只能有一个transactional.id的活跃 Producer。第二个initTransactions会触发 fencing第一个被踢出。Q3: 事务中的消息什么时候对消费者可见read_committed模式下当 TransactionMarkerCOMMIT被追加到分区后消费者才能读到该事务的消息。注意消费者可能滞后于 LSOLast Stable Offset。十一、总结Kafka 事务的本质 幂等生产者 (PID Seq) │ ├── 单分区不重复 │ ▼ 事务生产者 (PID Epoch transactional.id) │ ├── 跨分区原子写入 ├── 原子绑定消费位移 ├── 僵尸实例隔离 (Fencing) │ ▼ 实现 Exactly-Once 语义 (Kafka Streams EOS)一句话Kafka 事务 幂等发送 原子多分区提交 消费位移原子绑定通过 Transaction Coordinator 和__transaction_state内部 Topic 协调消费者配合read_committed实现端到端的 Exactly-Once。