在分布式消息系统中,「至少一次(at-least-once)」是默认语义:生产者重试会导致重复消息,消费者崩溃重启会导致重复消费。大多数场景下靠幂等消费兜底即可。但订单支付、账户转账这类强一致业务,需要真正的 Exactly-Once(精确一次)。Kafka 从 0.11 起引入事务机制,让「多分区原子写入」「消费-处理-产出原子化」成为可能。本文从幂等生产者讲起,深入事务协调器与两阶段提交,再到 read-process-write 端到端模式与 Outbox 集成,帮你判断何时用、怎么用、代价是什么。
1. 为什么需要事务:从至少一次到精确一次
1.1 消息系统的语义光谱
| 语义 | 含义 | 代价 |
|---|---|---|
| At-most-once(至多一次) | 消息可能丢,但绝不重复 | 最省,适合可容忍丢失的遥测 |
| At-least-once(至少一次) | 不丢,但可能重复 | 默认,靠幂等消费兜底 |
| Exactly-once(精确一次) | 不丢也不重复 | 需要事务/幂等,性能有代价 |
Kafka 的 acks=all + 重试 保证的是不丢(至少一次),重试本身就会造成重复——同一批消息被发送两次,broker 存两份。
1.2 为什么「至少一次」不够
看一个经典的转账场景:
① 用户发起转账 → ② 扣款事件写入 Kafka → ③ 消费者把扣款写入账户库
(发送时网络抖动,重试成功)
④ 同一事件又发一次 → ⑤ 账户库重复扣款
如果生产端多写了一条、消费端多处理了一次,账户就被扣两次钱。Kafka 事务要解决的核心问题正是:让「写入多个分区」与「消费+产出」变得原子。
1.3 Kafka 事务能保证什么(边界)
需要先明确边界——Kafka 事务保证的是消息层面的精确一次:
- ✅ 同一事务内的消息要么全部可见,要么全部不可见(跨分区、跨 Topic 原子);
- ✅ 结合幂等生产者,重试不产生重复;
- ✅ 结合 read-process-write,流处理应用可做到输入输出一致;
- ❌ 不是分布式事务——不管理 MySQL/Redis 等外部系统(那需要 Outbox 或 Saga 配合)。
一句话:事务把 Kafka 从「高吞吐队列」升级为「可靠的原子写入系统」——但它只对 Kafka 内部负责,外部系统的原子性要靠你自己拼。
2. 幂等生产者:事务的地基
2.1 幂等解决的问题:重试重复
普通生产者重试时,broker 无法区分「这条是新消息」还是「重发上一条」,于是同一消息被存多次。幂等生产者(Idempotent Producer) 用三要素让 broker 能识别重复:
<Producer ID(PID), 分区, 序列号>
→ broker 按(PID, 分区)记录最大连续序列号
→ 收到重复序列号 → 直接返回成功(已存过)
→ 收到乱序/缺失 → 报 OutOfOrderSequenceException
2.2 开启幂等
Properties props = new Properties();
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true); // 开启幂等
props.put(ProducerConfig.ACKS_CONFIG, "all"); // 幂等要求 acks=all
props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5);
props.put(ProducerConfig.RETRIES_CONFIG, Integer.MAX_VALUE);
开启幂等的硬性约束:
| 配置 | 幂等强制值 |
|---|---|
acks | 必须为 all |
max.in.flight.requests.per.connection | 需 ≤5(默认 5,可兼容) |
retries | 需 >0,官方建议 MAX_VALUE |
delivery.timeout.ms | 需 ≥ 60000 |
从 Kafka 3.0 起幂等成为默认开启。单独用幂等即可消除单个生产者的重复,但无法跨多个生产者实例去重——这正是事务要补的一层。
2.3 幂等 vs 事务的关系
幂等生产者:单个生产者实例,重试不重复
↓ 升级
事务生产者:同一 producer,多批次原子提交 + 故障恢复(fencing)
事务 = 幂等 + 事务协调器 + 事务标记。幂等是地基,事务是上层建筑。
一句话:幂等生产者靠「PID+序列号」让 broker 识别重复,解决单实例重试重复;事务在此基础上叠加跨批次的原子性与故障时的旧实例隔离。
3. 事务机制核心:事务协调器与两阶段提交
3.1 事务协调器(Transaction Coordinator)
Kafka 在内部 Topic __transaction_state 中维护每个事务的状态,由一组 broker(默认 transaction.state.log.replication.factor=3)扮演 事务协调器:
生产者 → FindCoordinator 定位协调器 → 协调器在 __transaction_state 记录事务状态
状态机:Empty → Ongoing → PrepareCommit → CompleteCommit
└→ PrepareAbort → CompleteAbort
3.2 两阶段提交的过程
一个完整事务的生命周期:
① 生产者发起事务(initTransactions)
② 发送消息到各分区(此时消息「未提交」,其他事务消费者不可见)
③ 发送 offset 提交到 __consumer_offsets(可选项,绑定消费位移)
④ 生产者调用 commitTransaction → 协调器协调所有参与分区写入事务标记
⑤ 协调器向 __transaction_state 写 PrepareCommit → 完成后写 CompleteCommit
消息被事务标记(Transaction Marker)标记为已提交后,开启 read_committed 的消费者才能读到。
3.3 关键 API 示例
// 事务生产者:转账扣款 + 通知 两条消息原子提交
Properties props = new Properties();
props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "payment-producer-1"); // 必须全局唯一
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
props.put(ProducerConfig.ACKS_CONFIG, "all");
props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions(); // 与协调器注册事务 ID
try {
producer.beginTransaction();
// 业务事件 + 通知事件原子写入
producer.send(new ProducerRecord<>("payments", "txn-001", "{\"debit\":100}"));
producer.send(new ProducerRecord<>("notifications", "txn-001", "{\"msg\":\"扣款成功\"}"));
producer.commitTransaction();
} catch (ProducerFencedException e) {
// 被新实例接管(fencing),必须放弃
producer.close();
} catch (Exception e) {
producer.abortTransaction(); // 回滚
}
3.4 僵尸防护(Zombie Fencing)
致命场景:旧生产者分区被接手后恢复,继续「提交事务」,与新生产者冲突。Kafka 的解法是 fencing:
每个事务 ID 有递增的 epoch
协调器收到更老 epoch 的生产者请求 → 抛 ProducerFencedException 拒绝
旧实例被「隔离」,无法再写 → 保证同一事务 ID 只有一个活跃生产者
一句话:事务 = 协调器状态机 + 两阶段提交 + 事务标记 + epoch 隔离——既保证跨分区原子,又防止旧实例「诈尸」乱写。
4. read-process-write:端到端精确一次
4.1 模式拆解
流处理最经典的可靠模式 read-process-write:
消费 Topic A(读取输入)
↓ 处理(去重、聚合、转换)
写入 Topic B(产出输出)
↓
消费 Topic C(下一个处理阶段)
要做到端到端精确一次,必须让「消费偏移量提交」与「产出消息写入」在同一个事务里原子完成——要么都成功,要么都回滚。
4.2 Kafka Streams 的 EOS 实现
Kafka Streams 内置支持,一行配置开启:
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "order-enricher");
props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG,
StreamsConfig.EXACTLY_ONCE_V2); // 精确一次(V2 起效率更高)
EXACTLY_ONCE_V2 相比 V1 的改进:V1 为每个任务分配独立事务 ID,任务重平衡时需要 initTransactions 同步等待;V2 把多个任务合并到一个事务,用固定数量的事务 ID 避免协调器热点,吞吐更高。
4.3 手动实现 read-process-write
不用 Streams 时,用原生客户端手动组合:
// 消费者:绑定到事务(提交位移参与事务)
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
consumer.subscribe(List.of("orders"));
// 生产者:开启事务
producer.initTransactions();
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
if (records.isEmpty()) continue;
producer.beginTransaction();
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
for (ConsumerRecord<String, String> r : records) {
// 业务处理
String enriched = enrich(r.value());
// 产出到下游 Topic
producer.send(new ProducerRecord<>("enriched-orders", r.key(), enriched));
offsets.put(new TopicPartition(r.topic(), r.partition()),
new OffsetAndMetadata(r.offset() + 1));
}
// 关键:把消费位移一起提交到事务里
producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
producer.commitTransaction();
}
要点:sendOffsetsToTransaction 把位移提交绑定到事务——下游产出成功则位移提交,失败则回滚,重启后从旧位移重读。这就是「输入-输出」原子化的原理。
一句话:read-process-write 的关键是
sendOffsetsToTransaction把消费位移和产出写入绑进同一事务——重启后两者一起回滚重来,才能做到端到端精确一次。
5. 事务消费:隔离级别与可见性
5.1 两个隔离级别
| 隔离级别 | 行为 | 适用 |
|---|---|---|
read_uncommitted | 读到所有消息,含未提交事务 | 默认,吞吐高 |
read_committed | 只读已提交事务消息,且只读完整事务 | 事务消费端 |
设置方式:
props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
5.2 read_committed 的两个语义细节
- 未提交的消息不可见:事务进行中写入的消息,
read_committed消费者读不到,直到CompleteCommit; - 中断的事务也会被读到(但以「已回滚」标记结束):Kafka 用事务终止标记(Transaction Abort Marker) 告诉消费者「这一段是作废的」,消费者需跳过
AbortedTransaction中的消息。
用 Java 消费时,ConsumerRecord 不带 abort 标记,但底层 fetch 响应会携带;多数高级 API 已自动跳过。手动处理需注意:读到 abort marker 时丢弃该事务内的所有消息。
5.3 什么时候必须 read_committed
- 事务生产者的下游消费者;
- Kafka Streams / Flink 的精确一次管道;
- 任何「不能读到半截事务」的强一致业务。
一句话:
read_committed让消费者只见已提交的事务、自动跳过回滚段——事务生产端配上它,才形成完整的精确一次链路。
6. Outbox 模式:把 Kafka 事务与业务数据库打通
6.1 问题的本质
Kafka 事务管不了 MySQL。典型场景:用户下单 → 写订单库 → 发 Kafka 事件。如果「写库成功、发消息失败」,下游永远不知道有新订单。Outbox 模式用「同一本地事务写业务表 + outbox 表」解决:
订单服务(本地事务):
① INSERT INTO orders (...) -- 业务数据
② INSERT INTO outbox (id, topic, key, payload) -- 待发事件
两行写同一事务 → 要么都成功,要么都失败
然后:
③ Debezium/自研 worker 读 outbox → 发 Kafka → 标记已发送
6.2 配合 Kafka 事务的双保险
- Outbox 保证「不丢」:事件与业务数据同事务落库,绝不一致;
- Kafka 事务保证「不重复」:CDC worker 发送用事务/幂等生产者,重试不重复;
- 位移绑定:worker 把 outbox 游标与发送绑定,故障恢复不重发。
6.3 Outbox 表设计
CREATE TABLE outbox (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
aggregate_id VARCHAR(64) NOT NULL, -- 业务聚合根 ID
topic VARCHAR(255) NOT NULL, -- 目标 Topic
key VARCHAR(255), -- 分区键
payload JSON NOT NULL, -- 事件体
status TINYINT DEFAULT 0, -- 0=待发 1=已发
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP
) ENGINE=InnoDB;
6.4 何时不需要 Outbox
- 允许丢事件的非关键日志/遥测;
- 已有可靠的消息最终一致方案(如 Saga + 补偿);
- 事件与数据库无强一致要求。
一句话:Outbox 让业务库与 Kafka 之间也做到「同事务不丢」——CDC 管道把 outbox 表搬到 Kafka,Kafka 事务再保证不重,双保险补齐「外部系统」缺口。
7. 事务的性能代价与适用场景
7.1 代价来自哪里
| 开销来源 | 说明 |
|---|---|
| 额外 RPC | begin/commit/abort 各需多轮协调器交互 |
| 元数据写放大 | 每个事务在 __transaction_state 写状态 |
| 吞吐下降 | 实测事务生产者吞吐约为普通生产者的 70%~85% |
| 延迟上升 | 每事务多 2~3 次往返,端到端延迟显著增加 |
7.2 优化的工程手段
- 合并事务:多条业务消息放同一事务提交,摊薄固定开销;
- 合理事务大小:事务过大会占用内存、加长回滚窗口,控制在数百条/事务量级;
- EXACTLY_ONCE_V2:Kafka Streams 用共享事务池,替代 V1 的逐任务事务;
- 只在关键路径用:批处理管道、日志聚合用 at-least-once + 幂等即可。
7.3 适用性判断
| 场景 | 是否用事务 |
|---|---|
| 订单支付、账务、库存扣减 | ✅ 用 |
| 用户行为埋点、日志聚合 | ❌ 幂等 + at-least-once 即可 |
| 流式 ETL、特征计算 | ⚠️ 视一致性需求,多数可用 at-least-once |
| 高吞吐压测基准 | ❌ 事务会拉低指标 |
一句话:事务是**「用吞吐换一致」的机制——只在强一致关键路径使用,并用合并事务、共享事务池**摊薄成本;能靠幂等解决的场景别上事务。
8. 常见坑与最佳实践
8.1 常见坑
| 坑 | 现象 | 对策 |
|---|---|---|
transactional.id 不唯一 | 生产端互相 fencing 失败 | 用主机名/Pod 名/业务前缀保证唯一 |
| 事务里发大消息 | 内存暴涨、回滚慢 | 拆分小事务或先写对象存储 |
| 消费者没开 read_committed | 读到未提交事务 | 事务下游一律开 read_committed |
| 事务模式 max.in.flight>5 | 序列号乱序报错 | 保持默认 ≤5 |
| 滥用事务做日志管道 | 吞吐腰斩 | 评估是否真需要精确一次 |
| abort 后不关生产者 | 连接泄漏 | catch 里关闭/重置生产者 |
8.2 最佳实践清单
- 事务 ID 命名规范:
{app}-{instance},如payment-producer-3; - 幂等默认开启,事务只在关键业务用;
- 捕获
ProducerFencedException并重建生产者,切勿无视; - read_committed 与 read_uncommitted 不要混用在同一 Topic 生态;
- 监控
__transaction_state大小与协调器负载; - Outbox + 事务双保险 构建「数据库→Kafka」端到端可靠。
9. 总结
本文从消息语义光谱出发,讲了 Kafka 事务的完整体系:
| 层级 | 机制 | 解决的问题 |
|---|---|---|
| 幂等生产者 | PID + 序列号 | 单生产者重试不重复 |
| 事务协调器 | 状态机 + 两阶段提交 | 跨分区原子写入 |
| Zombie Fencing | epoch 隔离 | 旧实例不干扰新实例 |
| read-process-write | sendOffsetsToTransaction | 消费-产出原子化 |
| read_committed | 隔离级别 | 只见已提交事务 |
| Outbox 模式 | 本地同事务落表 | 业务库与 Kafka 一致 |
一句话记住:Kafka 事务 = 幂等地基 + 协调器两阶段提交 + fencing 隔离 + 位移绑定,让「多分区写入」与「消费-产出」变成原子操作,实现端到端精确一次——但它只对 Kafka 内部负责,跨系统的强一致需要 Outbox/Saga 补齐。用吞吐换一致,只在关键路径用。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。