Skip to content

Kafka Exactly-Once 语义

提出问题

消息系统的交付语义是面试和架构设计中的高频考点。Kafka 从 0.11 版本开始支持 Exactly-Once 语义(EOS),但很多开发者只停留在"知道有幂等生产者和事务"的层面。面试官问 EOS 通常不是考概念,而是追问:你说的 Exactly-Once 到底是怎么做到的?事务能保证什么、不能保证什么?在流处理场景中 consume-process-produce 如何形成完整闭环?这些问题不搞清楚,生产环境要么不敢用事务,要么用了却喝了一壶。

分析问题

三种交付语义的本质

Kafka 的交付语义分三个级别:

语义生产者行为消费者行为典型场景代价
At-most-once发送后不重试消费后立即提交 offset可丢失的监控指标、日志采样丢数据
At-least-once失败后重试处理完再提交 offset绝大多数业务场景下游需幂等
Exactly-once幂等+事务read_committed 隔离金融、流处理精确聚合吞吐下降 15-25%
  • At-most-once:生产者发送后不重试,消息可能丢失,但不会重复。适合纯日志、可丢失的监控数据。
  • At-least-once:生产者重试 + 消费者至少一次提交,消息不丢但可能重复。这是绝大多数场景的默认语义,需要下游幂等配合。
  • Exactly-once:每条消息对系统状态的影响恰好一次。在分布式系统中,这通常意味着"输出幂等 + 事务保证",而不是字面上物理消息只发一次。

Kafka 的 EOS 拆成了两个独立的能力:幂等生产者(防止生产端重复)和 事务(原子写多分区 + 消费者隔离)。

幂等生产者:PID + Sequence Number 去重

开启 enable.idempotence=true 后,每个生产者实例被分配一个唯一的 PID(Producer ID),每个分区维护一个单调递增的 Sequence Number。

java
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("enable.idempotence", true);
// acks 必须为 all,retries 必须大于 0
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);

KafkaProducer<String, String> producer = new KafkaProducer<>(props);

Broker 端以 <PID, 分区, SeqNum> 为 key 缓存最近 5 个请求的 Sequence Number。当生产者重试发送同一条消息时,Broker 发现 SeqNum 已消费,直接返回成功而不重复写入。这保证了单一生产者实例内单次发送的消息不会重复。

时序流程

生产者                              Broker
  │                                   │
  │── ProduceRequest(SeqNum=100) ────→│ 写入成功,记录 SeqNum=100
  │←─────── acks=all 响应 ───────────│
  │                                   │
  │── ProduceRequest(SeqNum=100) ────→│ 重试,发现 SeqNum 已存在
  │←─────── 返回成功(去重) ────────│ 不重复写入
  │                                   │

幂等生产者的边界:解决不了跨分区原子性、也解决不了生产者重启后 PID 变更带来的重复问题。PID 重启后重置,Broker 端的 <PID, 分区, SeqNum> 缓存失效。KIP-588 引入了 Producer Fencing 机制:新启动的生产者通过 InitProducerId 请求获取新 PID,同时 Broker 会 fence 掉旧 PID 的未完成请求,但这只能防止旧 PID 的并发写入,无法防止旧 PID 已写入但 Broker 还没来得及响应就崩溃的 case——那部分消息对客户端算"发送失败",重启后客户端重试会用新 PID 发,Broker 端看在两条不同的 <PID, 分区, SeqNum> 记录上,没法去重。

这就是为什么单靠幂等生产者不构成 Exactly-Once,需要事务层来解决。

事务:原子写多分区 + 消费者隔离

Kafka 事务在幂等的基础上,引入了 Transaction CoordinatorTransaction Log(内部 Topic __transaction_state)。事务的生命周期:

java
// 1. 初始化事务
producer.initTransactions();

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

// 3. 发送消息到多个分区(原子写入)
producer.send(new ProducerRecord<>("topic-a", "key1", "value1"));
producer.send(new ProducerRecord<>("topic-b", "key2", "value2"));

// 4. 提交事务
producer.commitTransaction();

事务时序流程

Producer              Transaction Coordinator (TC)        Broker (Topic-A)    Broker (Topic-B)
  │                              │                              │                   │
  │── FindCoordinator ──────────→│                              │                   │
  │←─── Coordinator 地址 ───────│                              │                   │
  │                              │                              │                   │
  │── InitProducerId ───────────→│ 分配 PID + Epoch             │                   │
  │←─── PID + Epoch ────────────│                              │                   │
  │                              │                              │                   │
  │── AddPartitionsToTxn ───────→│ 注册分区                     │                   │
  │                              │                              │                   │
  │── ProduceRequest(txn) ──────│─────────────────────────────→│ 写入(标记 pending)│
  │── ProduceRequest(txn) ──────│───────────────────────────────────────────────→│ 写入(标记 pending)│
  │                              │                              │                   │
  │── EndTxn(Commit) ──────────→│ 写 __transaction_state       │                   │
  │                              │── CommitTxnRequest ─────────→│ 标记 committed    │
  │                              │── CommitTxnRequest ──────────────────────────→│ 标记 committed    │
  │                              │                              │                   │

其核心机制:

  • Transaction Coordinator 负责维护事务状态(Ongoing → PrepareCommit → Committed 或 Aborted),通过 __transaction_state Topic 持久化状态。
  • 写入时,消息先标记为 ABORTEDCOMMITTED 状态,Broker 正常存储但默认不让未提交事务的消费者读到。
  • 消费者设置 isolation.level=read_committed 后,只会消费已提交事务的消息,跳过被中止事务的写入。默认 read_uncommitted 则能看到未提交消息(包括最终会被 abort 的)。
  • 最终提交时,Coordinator 将事务状态写入 Transaction Log,Broad 收到 CommitTxnRequest 后将消息的 COMMITTED 标记写入。

事务的坑(面试高频):

  1. 事务超时transaction.timeout.ms 默认 60000ms。超时后 Coordinator 自动 abort 事务,生产者继续 commit 会收到 ProducerFencedException。生产环境如果遇到批量处理耗时较长,必须调大这个值或拆分成多个小事务。
  2. 事务 Fencing:同一个 transactional.id 下,后启动的 Producer 会 fenced 掉旧的。如果 Producer 网络分区后 reconnect,旧的还未超时的事务会被强制 abort。这是保证幂等的关键机制,但也是"事务 failover 后一定有还没处理完的消息被 abort"的代价。
  3. 事务隔离的误解:Kafka 事务保证的是跨分区原子写入,而不是分布式事务中的 ACID。事务中的消息在 commit 前对其他生产者完全可见(因为它们没有事务标记)。只有设置了 read_committed 的消费者才看不到未提交消息。说"Kafka 事务就是分布式事务"在面试里会被扣分。

一句话总结:Kafka 事务 = 跨分区原子写入 + 消费者隔离,不是 ACID 事务。

流处理中的 EOS 完整闭环:consume-process-produce

Kafka Streams 或自己实现的 consume-process-produce 循环,是 EOS 最经典的应用场景:

java
// 使用 Kafka Streams 的 Exactly-Once 语义
Properties props = new Properties();
props.put(StreamsConfig.APPLICATION_ID_CONFIG, "word-count-app");
props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, 
         StreamsConfig.EXACTLY_ONCE_V2);  // KIP-732 优化版

StreamsBuilder builder = new StreamsBuilder();
KStream<String, String> source = builder.stream("input-topic");
source
    .flatMapValues(value -> Arrays.asList(value.toLowerCase().split("\\W+")))
    .groupBy((key, word) -> word)
    .count(Materialized.as("counts-store"))
    .toStream()
    .to("output-topic");

consume-process-produce 闭环的原子性难题

消费消息 ──→ 业务处理 ──→ 写入输出 ──→ 提交 offset
  │              │            │            │
  └──────────────┴────────────┴────────────┘
           必须同时成功或同时失败

这个闭环的难点在于:消费消息、处理、写入结果、提交 offset 这四个动作必须原子化。如果没有事务,会出什么情况?

现实案例:某金融公司做实时风控聚合,消费交易流水 → 计算用户累计交易额 → 写入结果 Topic。业务用 enable.auto.commit=true + 轮询间隔 5s。某次流量突增,处理耗时超过 5s,offset 自动提交了但结果还没写入。应用重启后,offset 已提交的旧消息不再消费,那些交易直接丢失了,当天风控规则漏判了 20 多笔大额交易。

Kafka Streams 的做法是:

  1. 将 offset 提交也视为一个事务写入(写到 __consumer_offsets 或内部 Topic)。
  2. 整个 consume-process-produce 在一个事务内完成:消费 input-topic 的消息、写入 output-topic、记录 offset 提交。
  3. 如果事务提交失败,三者全部回滚,重启后从上一个已提交的 offset 重新消费。

这本质上是通过事务将 offset 和输出绑定在一起,做到了"消费一次,输出一次"。

KIP-732(EXACTLY_ONCE_V2) 进一步优化了性能:

  • V1 版本对每个 task 都创建一个独立的事务,Task 数多时 Coordinator 压力大。
  • V2 版本改用一个共享事务实例,减少事务 fencing 和协调的开销,吞吐提升约 20-30%。
  • V2 已经成为默认值(Kafka 3.0+),除非显式设置 EXACTLY_ONCE(等同于 V1)。

性能代价

EOS 不是免费的。事务引入了额外的 RTT(与 Coordinator 交互、写 Transaction Log),每个事务的 commit 都需要一次刷盘。实测对比(3 分区、3 副本,acks=all,Kafka 3.5,生产者线程数=4,消息体 1KB):

配置吞吐(msg/s)相对基准适用场景
acks=1(非幂等)185,000基准可丢失日志
acks=all(幂等生产者)172,000-7%大多数业务
事务(每 1000 条提交)148,000-20%批量流处理
事务(每 100 条提交)130,000-30%对时延敏感
事务(每 1 条提交)80,000-57%不要这么用

所以实践中,大多数场景用幂等生产者 + 下游幂等消费就够了,只有需要精确一次 stream processing 时才开事务。

常见面试追问

Q1:幂等生产者重启后,为什么还会丢消息?

PID 重启后重置,Broker 端 <PID, 分区, SeqNum> 缓存对旧 PID 的 SeqNum 记录会失效。如果应用发送消息后 Broker 回复了成功但响应在网络中丢失,应用判定为"发送失败"并重试——此时用新 PID 发送,Broker 端看到的是全新的 SeqNum 序列,无法去重,导致重复。所以幂等生产者只保证"同一连接生命周期内不重复",不保证应用级 Exactly-Once。

Q2:Kafka 事务能保证"消费一次、处理一次、输出一次"吗?

能,但前提是消费、处理、输出、offset 提交都在同一个事务中(即 Kafka Streams 模式)。如果消费和输出不在同一个事务内——比如你用普通 Consumer 消费消息,然后手动调用 Producer 的事务——offset 提交不在事务范围内,你的业务逻辑崩溃后,offset 可能已经提交了但输出还没写入,消息就丢了。

Q3:read_committed 消费者能读到什么?

  • 读到所有已提交事务写入的消息。
  • 跳过已中止事务写入的消息(这些消息物理上仍然在分区日志里,但 Consumer 不会返回给业务代码)。
  • 对于非事务写入的消息,照常读取。

注意:read_committed 不影响事务内消息对其他事务生产者的可见性——事务中的消息在 commit 前,Broker 已经存储了,只是事务消费者不返回而已。非事务生产者/消费者不受限制。

总结

从面试角度,可以这样组织回答:Kafka 的 Exactly-Once 是分层实现的——底层是幂等生产者(PID + Sequence Number 去重),上层是事务(Transaction Coordinator + 状态 Topic + 原子提交)。幂等解决"单连接不重复",事务解决"跨分区原子写 + 消费 offset 绑定"。在流处理场景中,Kafka Streams 通过事务将 consume-process-produce 做成原子闭环。但 EOS 有性能代价,生产环境按需启用,不要无脑开事务。

生产上的一条经验:先保证下游幂等,再考虑生产者事务。多数场景下幂等生产者 + 幂等消费已经够用,事务留着给需要精确一次流处理的地方。

参考

参考:Kafka 官方文档 Transactions 章节;KIP-98(Exactly Once Delivery and Transactional Messaging);KIP-588(Producer Fencing);KIP-732(EXACTLY_ONCE_V2);《Kafka: The Definitive Guide》第 4 章

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