Kafka 事务消息

mq 📚 learning kafka-transaction · kafka · transaction · transactional-id · exactly-once · kafka-streams

TL;DR(30 秒扫完)

  • Kafka 0.11+ 原生支持事务(Confluent 2016 年推出)
  • 核心目标:让"多 Topic 写入 + Consumer offset 提交"原子化
  • 典型场景:CDC(Debezium)、Kafka Streams、消息路由
  • 4 个核心 API:initTransactions / beginTransaction / sendOffsetsToTransaction / commitTransaction
  • 底层:transactional.id → PID → __transaction_state Topic → begin/commit 记录
  • Consumer 端:isolation.level=read_committed 只读已提交事务
  • 局限:只能保证 Kafka 内部 exactly-once,跨系统仍需外部存储配合

关键结论

结论 AKafka 从 0.11 起原生支持事务消息,不是"业务自行实现"
结论 BKafka 事务的核心目标是跨 Topic 原子写 + offset 提交原子化
结论 C事务消息 ≠ 消息不丢——不丢靠 acks+ISR+手动 commit;事务靠 transactional.id + __transaction_state
结论 DKafka 事务只能保证 Kafka 内部 exactly-once,跨系统(如 Kafka+MySQL)仍需业务侧事务或最终一致性

完整讲解(费曼四步)

STEP 1 · 概念

Kafka 事务消息是 Kafka 0.11+ 引入的原生能力,让"多个 Topic 写入 + Consumer offset 提交"作为一个原子操作——要么全部成功,要么全部回滚。

STEP 2 · 大白话

餐厅下单 + 结账比喻:
  • 消费者点了 3 道菜(写 3 个 Topic)+ 结账(提交 offset)
  • 没有事务:点了 3 道菜但结账失败 → 下次来重复点单(重复消费)
  • 有事务:要么 3 道菜 + 结账都成功,要么都回滚——下次来还是干净的起点

STEP 3 · 底层

典型场景

场景说明
CDC(Change Data Capture)Debezium 消费数据库变更写入 Kafka,同时提交消费位点,防止重复处理
Kafka Streams流处理读 A Topic、写 B Topic,两个操作必须原子
消息路由消费一个 Topic、按条件分发到多个 Topic

4 个核心 API

// 1. 初始化(只需一次)
producer.initTransactions();

// 2. 开始事务
producer.beginTransaction();

// 3. 写入消息
producer.send(record1);
producer.send(record2);

// 4. 把消费 offset 加入事务(关键!)
producer.sendOffsetsToTransaction(offsets, groupId);

// 5. 提交事务(要么 commit 要么 abort)
producer.commitTransaction();
// 或 producer.abortTransaction();
关键:sendOffsetsToTransaction 把消费 offset 加入事务,提交事务时 offset 也会一并提交。

底层实现

1. Producer 用 transactional.id 从 Broker 换取 PID(Producer ID)
2. PID 保存在 __transaction_state Topic
3. 事务开始 → 向 __transaction_state 写 begin 记录
4. 事务提交 → 向 __transaction_state 写 commit 记录
5. Consumer 端 isolation.level=read_committed 只读已提交事务

关键配置

参数位置说明
transactional.idProducer事务标识,必须设置(幂等也依赖它)
enable.idempotenceProducer开启幂等(默认 false)
isolation.levelConsumerread_committed(只读已提交)/ read_uncommitted(默认)
transaction.timeout.msProducer事务超时时间,默认 60s

与 Exactly-Once 的关系

端到端 Exactly-Once = 
  Producer 幂等(enable.idempotence=true)
+ Broker 副本可靠(replication + min.insync)
+ Consumer 事务(transactional.id + isolation.level=read_committed)
+ 外部存储事务(MySQL 唯一约束或事务)
Kafka 事务只解决 Kafka 内部 exactly-once,跨系统仍需业务配合。

STEP 4 · 简化

一句话总结:Kafka 0.11+ 原生支持事务,靠 transactional.id + __transaction_state Topic 实现"多操作原子化"。 记忆口诀:
  • 4 步:init → begin → send + sendOffsets → commit
  • 关键 API:sendOffsetsToTransaction(把消费 offset 加入事务)
  • 底层:transactional.id → PID → __transaction_state Topic
  • 局限:只保证 Kafka 内部 exactly-once

常见误区

认为 Kafka 没有事务消息
Kafka 0.11+ 原生支持事务消息
把"事务消息"和"消息不丢失"混淆
不丢靠 acks+ISR+手动 commit;事务靠 transactional.id + __transaction_state
认为 Kafka 事务能保证跨系统原子性
只保证 Kafka 内部 exactly-once,跨系统需业务配合
忘记 sendOffsetsToTransaction
不调用这个 API,消费 offset 不会加入事务
忘记设置 isolation.level=read_committed
Consumer 端不设这个,会读到未提交事务的消息

延伸追问

Kafka 事务消息和 RocketMQ 事务消息有什么区别?
Kafka:靠 transactional.id + PID + __transaction_state Topic,事务范围限于 Kafka 内部
Kafka 事务能保证外部存储(MySQL)也原子提交吗?
不能。需要业务侧事务或最终一致性方案(本地消息表、SAGA)
transaction.timeout.ms 超时会怎样?
事务自动 abort,Producer 抛异常;需处理重启或降低事务范围
事务消息能保证 Exactly-Once 吗?
只能保证 Kafka 内部 exactly-once;端到端 Exactly-Once 需要 Producer 幂等 + Broker 副本可靠 + Consumer 事务 + 外部存储事务
sendOffsetsToTransaction 什么时候调用?
消费完一批消息后,处理业务前;把消费 offset 加入事务,提交事务时 offset 一并提交

速查表

4 API: initTransactions → beginTransaction → sendOffsetsToTransaction → commitTransaction
关键: transactional.id 必须设置;sendOffsetsToTransaction 是关键
底层: transactional.id → PID → __transaction_state Topic
Consumer: isolation.level=read_committed
局限: 只保证 Kafka 内部 exactly-once,跨系统需外部配合

Anki 候选卡片

Q: Kafka 从哪个版本开始支持事务消息?
A: 0.11(Confluent 2016 年推出)
Q: Kafka 事务的 4 个核心 API?
A: initTransactions / beginTransaction / sendOffsetsToTransaction / commitTransaction
Q: Kafka 事务底层靠什么实现?
A: transactional.id → PID → __transaction_state Topic → begin/commit 记录
Q: Consumer 端怎么只读已提交事务?
A: isolation.level=read_committed
Q: Kafka 事务能保证跨系统 Exactly-Once 吗?
A: 不能,只保证 Kafka 内部;跨系统需业务侧事务或最终一致性

关联题目

  • ⚠️ 《Kafka 支持事务消息吗?如何实现的?》— Round 2 Q8, ⭐(核心概念错误,说"Kafka 没有事务")
  • 关联知识

    Kafka 0.11+ 原生支持事务,让多 Topic 写入 + offset 提交原子化
    餐厅下单+结账:要么 3 道菜+结账都成功,要么都回滚
    ✦ 记 忆 口 诀 ✦
    init → begin → send + sendOffsets → commit;transactional.id 必须设
    关键可视化
    4 个核心 API
    flowchart TD
      A[initTransactions 初始化] --> B[beginTransaction 开始事务]
      B --> C[send 写入消息]
      B --> D[sendOffsetsToTransaction 把消费 offset 加入事务]
      C --> E[commitTransaction 提交]
      D --> E
      E --> F[要么全部成功 要么全部回滚]
    底层实现
    flowchart TD
      A[Producer 设置 transactional.id] --> B[从 Broker 换取 PID]
      B --> C[PID 保存在 __transaction_state Topic]
      C --> D[事务开始 写 begin 记录]
      D --> E[事务提交 写 commit 记录]
      E --> F[Consumer isolation.level=read_committed 只读已提交]
    知识关系

    ⬆️ 前置(Prerequisite)

    kafka-reliability

    🔄 延伸(Extends)

    暂无

    ⚡ 对比(Contrast)

    消息可靠性— 不丢靠 acks+ISR+手动 commit;事务靠 transactional.id + __transaction_state
    🎯 概念 📏 规则 ⚠️ 误区 🔍 追问 ✨ 口诀 共 0 张卡,点击翻面