KafkaNotes

第 14 章:交付语义:幂等、事务与 Exactly Once

zjc 于 2026-01-14 发布

这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 “消息会不会丢、会不会重”是 Kafka 工程的灵魂问题。本章把三种交付语义的定义、实现手段与适用边界彻底讲清楚,并给出可以落地的代码。

14.1 三种交付语义

语义 含义 典型实现 风险
At most once 至多一次 先提交位移,后处理消息 可能丢,不会重
At least once 至少一次 先处理消息,后提交位移 可能重,不会丢
Exactly once 精确一次 幂等 + 事务 + 消费端幂等存储 实现复杂,范围有限制

注意一个常见误解:Exactly once 不是魔法,而是若干机制组合出的端到端属性,而且它的“端”有明确边界。Kafka 事务保证的是“Kafka 到 Kafka”的处理链路精确一次;一旦涉及外部系统(MySQL、Redis),仍需要业务侧幂等或事务配合。

14.2 至少一次:默认的工程底线

标准做法(第 6 章):

poll -> 处理 -> commit offset

处理完成但未来得及提交时宕机,重启后会重复消费,所以是“至少一次”。重复交给下游幂等解决:

绝大多数业务系统应该停在这里:至少一次 + 业务幂等。 简单、可控、可观测。

14.3 幂等生产者:解决发送端重复

问题:为什么发送会重复

Producer 发消息给 Broker,Broker 写入成功,但响应在网络中丢了。Producer 超时重试,Broker 收到两条相同消息。网络层面无法区分“没写进去”和“写进去了但响应丢了”。

原理:PID + Sequence Number

开启幂等(enable.idempotence=true,3.0 起默认开启)后:

  1. Producer 启动时从 Broker 获取一个 PID(producer id);
  2. 对每个 (PID, partition) 维护单调递增的 sequence number
  3. 消息(batch)带着 PID + epoch + sequence 写入 Broker;
  4. Broker 为每个分区缓存最近 5 个 batch 的序号:
    • 序号正好是预期值:正常写入;
    • 序号小于预期(重复):丢弃并返回成功
    • 序号大于预期(乱序缺口):返回 OutOfOrderSequenceException
Producer 重试:  seq=7 -> Broker 已有 seq=7 -> 丢弃,返回成功
               客户端视角:发送成功,且日志中只有一条

约束与注意:

14.4 事务:跨分区原子写

幂等解决“单分区单会话不重”,事务解决“多条消息原子写入多个分区”。

完整 API

Properties props = new Properties();
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-producer-1");
props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
// 其余配置同普通生产者

KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();      // 1. 注册到事务协调器,获取 epoch(会 fence 旧实例)

try {
    producer.beginTransaction();  // 2. 开启事务

    producer.send(new ProducerRecord<>("order-events", "o1", "created"));
    producer.send(new ProducerRecord<>("audit-events", "o1", "audit"));

    producer.commitTransaction(); // 3. 提交:两条要么都可见,要么都不可见
} catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
    producer.close();             // 不可恢复,直接关闭
} catch (KafkaException e) {
    producer.abortTransaction();  // 可恢复错误则回滚
}

要点:

14.5 LSO 与 read_committed

事务消息写入后并非立即可见。Broker 为每个分区维护 LSO(Last Stable Offset)

日志:  [普通消息] [事务T1消息] [普通消息] [事务T2消息] ...
                          ^                 ^
                        LSO 在这里         T2 未提交

read_committed 消费者最多读到 LSO。事务提交前,T1 的消息不可见;一旦提交,marker 之后的 LSO 前移,消费者才能读到。

监控事务相关指标:

14.6 Consume-Transform-Produce:流式精确一次

最常见的 exactly-once 形态:消费上游 topic,处理后写下游 topic,把“输入位移提交”和“输出写入”放进同一个事务

props.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");
props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "etl-job-1");

consumer.subscribe(List.of("input-topic"));
producer.initTransactions();

while (true) {
    var records = consumer.poll(Duration.ofMillis(500));
    if (records.isEmpty()) continue;

    producer.beginTransaction();
    try {
        for (var r : records) {
            producer.send(new ProducerRecord<>(
                    "output-topic", r.key(), transform(r.value())));
        }
        // 把消费位移也纳入事务
        Map<TopicPartition, OffsetAndMetadata> offsets = buildOffsets(records);
        producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());

        producer.commitTransaction();
    } catch (Exception e) {
        producer.abortTransaction();
        // 重新 poll 会从上一次提交位移开始
    }
}

这段代码保证了:输出消息与位移提交原子生效。要么下游看到结果且位移前移,要么两者都回滚,重放后再次写入。配合输入输出都在 Kafka,就构成了 Kafka 端到端的 exactly-once。

Kafka Streams 的 exactly_once_v2 本质就是把这套机制框架化,业务代码无需手写事务边界。

14.7 消费端的精确一次

如果下游是外部系统,事务管不到那里,需要自己设计幂等:

数据库方案

-- messages 表主键 = topic + partition + offset 或业务 request_id
INSERT INTO consumed_messages(message_key, topic_name, partition_id, offset_id)
VALUES (?, ?, ?, ?)
ON CONFLICT DO NOTHING;

插入成功才处理业务;插入冲突说明已处理过,直接跳过。把业务写和去重表放同一个本地事务,就是完整的端到端精确一次。

状态机方案

CREATED -> PAID -> SHIPPED
重复的 PAID 事件到来时,当前状态已是 SHIPPED -> 直接忽略

代价与选择

事务会带来吞吐下降(两阶段提交与 marker 写入)和延迟上升。问自己:

  1. 重复会造成资金/库存错误吗?必须严格幂等;
  2. 重复只是多记一条日志/多推一次通知吗?至少一次 + 去重可能就够;
  3. 输入输出是否都在 Kafka 内?是则优先用事务/Streams。

14.8 常见误区

  1. 开了幂等就不丢消息:幂等只防重,不丢消息靠 acks=all + 副本配置;
  2. 事务保证外部系统原子性:不能。MySQL 写入失败不会被 Kafka 事务回滚;
  3. read_committed 能去重业务消息:它只过滤未提交/已中止事务,不处理业务层重试造成的逻辑重复;
  4. 换个 transactional.id 重启:必须保持稳定,否则失去 fence 能力;
  5. 把所有业务都上事务:事务有真实成本,大多数场景“至少一次 + 幂等”更划算。

本章小结

思考题

  1. 为什么幂等生产者不能防止“进程重启后重发同一业务消息”?该用什么解决?
  2. sendOffsetsToTransaction 如果漏掉,会发生什么?
  3. 设计一个“Kafka -> Redis -> MySQL”链路的端到端精确一次方案,指出每一段由谁保证。