Skip to content

Kafka 事务原理与分布式事务场景

提出问题

Kafka 从 0.11 版本开始引入事务机制,但在实际面试和线上使用中,很多人对它的理解停留在"Kafka 也支持事务了"这个层面,搞不清楚它到底解决了什么问题、和 RocketMQ 的事务消息有什么区别。

最常见的误解有三个:把 Kafka 事务理解成"跨数据库的分布式事务";以为开了事务就能保证端到端 Exactly-Once;或者在生产环境直接照搬事务配置,结果发现吞吐量暴跌、消费延迟飙升。

Kafka 事务的核心应用场景是流式 ETL 中的消费-生产模式:从 Kafka 的某个 Topic 消费消息,经过处理后写入另一个 Kafka Topic(或同一个 Topic 的不同分区),需要保证"消费 offset 提交"和"生产消息写入"这两个操作原子完成——要么都成功,要么都回滚。这个场景在 Kafka Streams、Flink 等流处理框架中非常常见。

举个例子:一个实时风控系统收到订单事件(Topic: orders),经过规则引擎计算后,把命中策略的订单发到告警 Topic(Topic: alerts)。如果"消费 orders offset"和"生产 alerts 消息"不一致,要么重复告警,要么漏告警。Kafka 事务解决的就是这个"消费-生产"原子性问题。

分析问题

事务机制的三个核心组件

Kafka 事务的实现依赖三个基础设施:

1. 事务协调器(Transaction Coordinator)

每个 Producer 调用 initTransactions() 时,会向 Broker 集群中某个节点注册为事务协调器。协调器负责分配 Producer ID(PID)和事务 Epoch,并维护事务的生命周期状态。协调器的高可用通过 __transaction_state 主题的多副本机制保证。

2. 事务日志(__transaction_state 主题)

这是一个内部主题,默认 3 个副本、50 个分区。记录每个事务的状态变迁:Ongoing → PrepareCommit → CompletedCommit,或 Ongoing → Abort。协调器故障时,新的协调器可以从这个主题恢复事务状态,继续推进。每条事务日志记录包含:事务 ID、PID、Epoch、超时时间、涉及的分区列表、当前状态。

3. 控制消息(Control Batches)

事务提交或中止时,Kafka 会在目标分区写入一个特殊的控制批次(Control Batch)。Consumer 端通过 isolation.level=read_committed 模式,根据控制消息来判断哪些消息已经提交、哪些是未提交的(跳过)。这也意味着,Consumer 在 read_committed 模式下读取时,会停在一个"最后稳定偏移量"(Last Stable Offset, LSO)之后不再前进,直到事务完成——这就是消费延迟的来源。

事务的完整工作流程

完整的时序流程如下:

Producer                     Coordinator                    Target Partition          Consumer
   |                              |                              |                      |
   |-- initTransactions() ------->|                              |                      |
   |                              |-- 分配 PID + Epoch           |                      |
   |<---- PID + Epoch ------------|                              |                      |
   |                              |                              |                      |
   |-- beginTransaction() ------->|                              |                      |
   |                              |-- 状态: Ongoing              |                      |
   |                              |                              |                      |
   |-- send(record) -------------->|                              |                      |
   |                              |-- 记录分区信息到事务日志     |                      |
   |                              |-- 转发消息至目标分区 ------->|                      |
   |                              |                              |-- 写入(未提交标记) |
   |                              |                              |                      |
   |-- sendOffsetsToTransaction() |                              |                      |
   |                              |-- 记录 offset 到事务日志     |                      |
   |                              |                              |                      |
   |-- commitTransaction() ------>|                              |                      |
   |                              |-- 状态: PrepareCommit         |                      |
   |                              |-- 写入控制消息 ------------>|-- 写入 Commit Batch  |
   |                              |-- 状态: CompletedCommit       |                      |
   |<---- commit 成功 ------------|                              |-- 消息变为可见       |
   |                              |                              |                      |
   |                              |                              |         (read_committed 模式)
   |                              |                              |<-- 消费已提交的消息   |

这个流程的关键点在于:事务真正提交时,协调器会先向事务涉及的所有分区写入 Commit Marker(控制消息),然后才将事务状态标记为 CompletedCommit。如果写入 Commit Marker 后协调器崩溃,Consumer 会看到部分分区有 Marker、部分没有,此时 Consumer 会等待所有分区都有 Marker 后才继续消费——这就是幂等和原子性的保证。

典型的事务使用代码:

java
producer.initTransactions();
try {
    producer.beginTransaction();
    // 从某个 Topic 消费消息
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000));
    for (ConsumerRecord<String, String> record : records) {
        // 处理业务逻辑
        String processed = process(record.value());
        // 生产到另一个 Topic
        producer.send(new ProducerRecord<>("output-topic", processed));
    }
    // 关键:将消费 offset 也纳入事务
    producer.sendOffsetsToTransaction(getOffsets(records), consumer.groupMetadata());
    producer.commitTransaction();
} catch (Exception e) {
    producer.abortTransaction();
}

sendOffsetsToTransaction() 是事务的精髓:它将当前批次的消费 offset 提交也放在同一个事务中。如果后续 commitTransaction() 成功,则消费 offset 和生产的消息一起对外可见;如果失败回滚,则 offset 不前进,消息也不会被 Consumer 读到。

幂等 Producer 与事务的关系

很多人搞不清幂等 Producer 和事务的关系。简单说:

  • 幂等 Producerenable.idempotence=true):解决 Producer 重试导致的消息重复。通过 PID + Sequence Number 实现,Broker 端做去重。保证单个分区内、单次会话的 Exactly-Once。
  • 事务:在幂等 Producer 基础上,通过事务协调器保证跨分区、跨会话的原子性。幂等是事务的前置条件——开启事务时系统自动开启幂等,反之不成立。

线上场景区分:如果只是"确保日志不丢不重",幂等 Producer 就够了;如果做"消费-生产"的流式 ETL,必须用事务。

事务的坑与局限

坑 1:事务超时

transaction.timeout.ms 默认 60 秒。如果事务内的业务处理超过这个时间,协调器会主动中止事务,Producer 端会抛出 TransactionTimeoutException。遇到超时问题时,不要惯性调大超时时间,而是应该先分析事务内为什么耗时这么长——是不是批量处理太大?外部依赖调用是否超时?

真实案例:某公司的 Kafka Streams 任务在高峰期每小时挂一次,日志报 TransactionTimeoutException。排查发现,每个事务内处理了 10000 条消息,处理过程中调用了 Elasticsearch 做结果写入,ES 集群 GC 导致响应变慢,事务超时。解决方案:不是调大超时时间,而是把批量大小从 10000 降到 2000,且给 ES 写入加超时熔断。

坑 2:事务堆积

read_committed 模式下的 Consumer 会跳过 LSO 之后的消息。如果某个事务长期未提交(比如事务内的外部 API 调用阻塞),Consumer 端会卡在 LSO 位置,所有后续消息都无法消费,导致消费延迟骤增。线上故障排查时,如果发现 Consumer Lag 正常但消费延迟却很高,优先检查 __transaction_state 主题中是否有 Ongoing 状态的事务残留。

排查命令

bash
# 查看事务状态
kafka-transactions.sh --bootstrap-server localhost:9092 describe --transaction-state

# 查看 __transaction_state 主题中未完成的事务
kafka-console-consumer.sh --bootstrap-server localhost:9092 \
  --topic __transaction_state --from-beginning \
  --formatter "kafka.coordinator.transaction.TransactionLog\$TransactionLogMessageFormatter"

坑 3:吞吐量下降

事务模式下,每条消息的生产都需要与协调器通信来维护事务状态,吞吐量相比非事务模式下降约 20-30%。对于高吞吐场景(日志收集、监控数据),不要轻易开启事务,非事务模式配合幂等 Producer 已经足够。

坑 4:事务协调器成为瓶颈

如果你的集群有大量事务 Producer,所有的 initTransactions()commitTransaction() 请求都要经过协调器。协调器所在的 Broker 可能成为热点。线上遇到过一个场景:Flink 作业的并行度是 128,每个 TaskManager 都开启了一个事务,128 个事务共享同一个协调器,该 Broker 的 CPU 冲到 90%。解决方案:通过 transactional.id 的前缀做哈希,让不同的事务分配不同的协调器。

坑 5:不是跨系统分布式事务

这是最大的误解。Kafka 的事务只保证"Kafka 内的多个分区写入"的原子性,不保证和数据库、Redis 等其他系统的操作原子性。如果业务需要"扣减数据库库存 + 发送 Kafka 消息"两个操作原子完成,Kafka 事务帮不上忙,需要用本地消息表或 TCC/Saga 模式。

总结

维度Kafka 事务RocketMQ 事务消息
核心能力跨分区原子写入 + 消费-生产 EOS半消息 + 回调反查,保证本地事务与消息发送一致性
典型场景流式 ETL(Kafka Streams / Flink)订单-库存-支付等跨系统分布式事务
是否跨系统否,仅限 Kafka 内部否,仅保证 Producer 端与 MQ 端的一致性
吞吐影响下降 20-30%下降约 10-20%
运维复杂度较高(协调器、事务日志、超时配置、协调器瓶颈)中等(回查接口、半消息 Topic 监控)
超时默认值60 秒(transaction.timeout.ms)6 秒(checkImmunityTime)
事务可见性read_committed 模式跳过未提交事务半消息对 Consumer 不可见,直到二次确认

面试话术示例:"Kafka 事务解决的是流处理场景下的消费-生产原子性问题,它和幂等 Producer 结合才能实现端到端 Exactly-Once。但它是 Kafka 内部的事务,不是跨数据库和 MQ 的分布式事务。如果业务需要跨系统事务,我会优先考虑本地消息表方案,或者用 Seata TCC 配合事务消息——不是因为 Kafka 事务不好,而是它解决的不是那个问题。"

关键要点

  • 事务 = 幂等 Producer(精确一次写入)+ 事务协调器(原子性)+ 控制消息(可见性控制)
  • sendOffsetsToTransaction() 是消费-生产原子性的关键,没有它事务就是半成品
  • read_committed 模式下 Consumer 会跳过未提交事务的消息,注意 LSO 卡住导致的消费延迟
  • 事务超时和事务堆积是线上最常见的事故,监控 __transaction_state 主题是基本操作
  • 不要把 Kafka 事务当跨系统事务用,选型前先搞清楚要解决的是"MQ 内部原子性"还是"跨系统一致性"
  • 高并行度场景(Flink、Kafka Streams)注意协调器热点问题,合理设计 transactional.id

参考:Apache Kafka 官方文档 — Transactions;《Kafka 权威指南(第 2 版)》第 7 章;Confluent 博客 — Exactly-Once Semantics;KIP-98 — Exactly Once Delivery and Transactional Messaging

手撕 → 框架 → 生产化,一步步把 AI Agent 工程化搞透。