---
title: Kafka 事务消息
type: concept
domain: mq
tags: [kafka, transaction, transactional-id, exactly-once, kafka-streams]
status: learning
created: 2026-09-22
last_reviewed: 2026-09-22
method: feynman
related_questions:
  - "Kafka 支持事务消息吗？如何实现的？"
related_knowledge:
  - ./kafka-reliability.md
  - ./kafka-message-structure.md
visualizations: []
anki_cards: 5
interview_rounds:
  - "2026-09-22-round-2-Q8"
---

# Kafka 事务消息

> Kafka 从 0.11 起**原生支持事务消息**，核心目标是让"多个 Topic 写入 + Consumer offset 提交"作为一个原子操作。

## 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，跨系统仍需外部存储配合

## 关键结论

- **结论 A**：Kafka 从 0.11 起**原生支持**事务消息，不是"业务自行实现"
- **结论 B**：Kafka 事务的核心目标是**跨 Topic 原子写 + offset 提交原子化**
- **结论 C**：事务消息 ≠ 消息不丢——不丢靠 acks+ISR+手动 commit；事务靠 transactional.id + `__transaction_state`
- **结论 D**：Kafka 事务**只能保证 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

```java
// 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.id` | Producer | 事务标识，必须设置（幂等也依赖它） |
| `enable.idempotence` | Producer | 开启幂等（默认 false） |
| `isolation.level` | Consumer | `read_committed`（只读已提交）/ `read_uncommitted`（默认） |
| `transaction.timeout.ms` | Producer | 事务超时时间，默认 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

## 常见误区

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

## 延伸追问

1. **Kafka 事务消息和 RocketMQ 事务消息有什么区别？**
   - Kafka：靠 `transactional.id` + PID + `__transaction_state` Topic，事务范围限于 Kafka 内部
   - RocketMQ：Half 消息 + 回查机制，事务消息与业务服务交互，可感知业务侧确认
2. **Kafka 事务能保证外部存储（MySQL）也原子提交吗？**
   - 不能。需要业务侧事务或最终一致性方案（本地消息表、SAGA）
3. **`transaction.timeout.ms` 超时会怎样？**
   - 事务自动 abort，Producer 抛异常；需处理重启或降低事务范围
4. **事务消息能保证 Exactly-Once 吗？**
   - 只能保证 Kafka 内部 exactly-once；端到端 Exactly-Once 需要 Producer 幂等 + Broker 副本可靠 + Consumer 事务 + 外部存储事务
5. **`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，跨系统需外部配合
```

## 关联题目（题库）

- ⚠️ 《Kafka 支持事务消息吗？如何实现的？》— Round 2 Q8, ⭐（核心概念错误，说"Kafka 没有事务"）

## 关联知识

- [Kafka 消息可靠性](./kafka-reliability.md)
- [Kafka 消息结构 + Offset](./kafka-message-structure.md)
- [主题地图](./_moc.md)

## Anki 候选卡片

1. **正**：Kafka 从哪个版本开始支持事务消息？**反**：0.11（Confluent 2016 年推出）
2. **正**：Kafka 事务的 4 个核心 API？**反**：`initTransactions` / `beginTransaction` / `sendOffsetsToTransaction` / `commitTransaction`
3. **正**：Kafka 事务底层靠什么实现？**反**：`transactional.id` → PID → `__transaction_state` Topic → begin/commit 记录
4. **正**：Consumer 端怎么只读已提交事务？**反**：`isolation.level=read_committed`
5. **正**：Kafka 事务能保证跨系统 Exactly-Once 吗？**反**：不能，只保证 Kafka 内部；跨系统需业务侧事务或最终一致性

---

*最后更新：2026-09-22*
