Skip to content

Kafka 不丢、不重与 Exactly Once:先划清系统边界 ​

Kafka 的 Exactly Once 是真实存在的,但它覆盖的是「从 Kafka 读、往 Kafka 写」这一段。你关心的「扣款恰好执行一次」发生在数据库里,不是打开一个配置就能得到的。

「Kafka 怎么保证消息不丢」「怎么避免重复消费」「Exactly Once 是怎么实现的」,这组问题的难点不在配置,而在于不先划定边界,就无法判断谁来保证:是生产者发出去的消息不丢,还是 Broker 宕机后已确认的消息不丢,还是业务数据库里的效果不丢?每一段的保证机制完全不同。

本文基于 Apache Kafka 4.3 的官方文档。示例代码使用 kafka-clients,文中用到的事务 API 在 3.x 与 4.x 中一致。ISR 收缩、事务中止、旧实例隔离和自动提交丢消息这几条,在 Kafka 4.3.1 的 3 broker 集群上实测过,见文末配套实验。

一、先说结论 ​

  • 把链路拆成三段分别回答:生产者到 Broker、Broker 内部存储与复制、消费者到业务副作用。
  • 生产端:acks=all 加幂等生产者,配合 min.insync.replicas ≥ 2 和足够的副本数。Kafka 3.0 起,生产者默认就是 acks=all 与 enable.idempotence=true,但 Topic 侧的 min.insync.replicas 默认仍为 1。
  • 消费端:先处理、后提交 offset,得到至少一次;这意味着必然可能重复,所以业务处理要幂等。
  • Kafka 的 Exactly Once 依靠事务性生产者把「输出消息」和「消费 offset」放进同一个 Kafka 事务,下游用 read_committed 读取。它只覆盖 Kafka 内部的读—处理—写链路,Kafka Streams 就建立在这个机制上。
  • 写到外部系统时,要么把 offset 和结果存进同一个外部系统的同一个事务,要么「至少一次 + 幂等」。数据库变更与发消息的一致性,用事务型 Outbox。

二、链路拆解:丢与重分别发生在哪 ​

业务库orders生产者ProducerBroker 集群LeaderFollower · ISR消费者Consumer下游库points__consumer_offsets1234561写库与发送 → Outbox2发送与重试 → acks=all + 幂等生产者3副本复制 → min.insync.replicas ≥ 24拉取 → read_committed(事务场景)5处理副作用 → 幂等消费6提交位置 → 先处理、后提交
图 1 · 一条消息经过的六个环节,以及每个环节上防丢、防重的主要手段
环节丢失的典型原因重复的典型原因
① 业务写库与发消息库提交了,发消息前进程崩溃发了消息,库事务回滚
② 生产者发送acks=0/1;回调里的失败没处理;超时后放弃网络超时后重试,而 Broker 实际已写入
③ Broker 复制只写入 Leader 就确认,Leader 随后宕机;不干净的 Leader 选举——
④⑤⑥ 消费处理先提交 offset,后处理失败处理成功,提交 offset 前崩溃或发生重平衡

下面按段展开。

三、生产端与 Broker:已确认的消息不丢 ​

3.1 三个配置要一起看 ​

配置作用于默认值建议
acks生产者all(3.0 起)保持 all
enable.idempotence生产者true(3.0 起,无冲突配置时)保持开启
replication.factorTopic自动创建时取 default.replication.factor,默认为 1生产环境至少 3
min.insync.replicasTopic / Broker1与 3 副本搭配设为 2
unclean.leader.election.enableTopic / Brokerfalse需要强一致时保持 false

acks=all 的含义是「Leader 等 ISR(同步副本集合) 中所有副本都写入后才确认」。问题在于 ISR 会收缩:如果另外两个副本都掉队了,ISR 里只剩 Leader 自己,acks=all 就退化成了 acks=1。

min.insync.replicas=2 补上了这个漏洞:ISR 数量低于 2 时,Broker 直接拒绝 acks=all 的写入,生产者收到 NotEnoughReplicas 类错误。这是用可用性换持久性——3 副本、min.insync.replicas=2 可以容忍 1 个副本故障而不影响写入,也不丢已确认的数据。

实测停掉两个 follower、ISR 收缩到只剩 leader 之后:min.insync.replicas=2 的 topic 上,acks=all 被拒绝,acks=1 照常成功;min.insync.replicas=1(默认值)的 topic 上,acks=all 也照常成功,此时只有 leader 一份数据。还有一个排查时容易看错的地方:服务端返回的是 NotEnoughReplicasException,但它属于可重试错误,生产者默认会一直重试到 delivery.timeout.ms,回调最终拿到的是 TimeoutException。看到投递超时,要同时查 ISR。

unclean.leader.election.enable=false 是另一侧的取舍:所有 ISR 副本都宕机时,Kafka 默认等待 ISR 副本恢复,而不是选一个数据落后的副本当 Leader,后者会丢掉已确认的消息。

3.2 幂等生产者解决了什么,没解决什么 ​

网络超时后,生产者不知道 Broker 是否已经写入,只能重试。没有幂等时,重试会在分区里产生重复消息。

幂等生产者的做法:Broker 为每个生产者分配一个 Producer ID,生产者为发往每个分区的批次附带递增的序列号,Broker 发现序列号重复就丢弃。所以:

  • 能解决:同一个生产者实例因重试导致的分区内重复;
  • 不能解决:应用层重复发送。业务代码自己调用了两次 send(),那就是两条不同的消息;
  • 不能解决:进程重启后的重复。新进程拿到新的 Producer ID,之前未确认的消息被业务重新发送时无法去重。跨重启的去重需要配置 transactional.id。

3.3 生产者的失败要被处理 ​

send() 是异步的,即使 acks=all,如果回调里的异常没人处理,丢了也不知道:

java
producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        // 重试已在客户端内部完成(受 delivery.timeout.ms 约束,默认 2 分钟)
        // 到这里说明最终失败:记录、告警,并交给补偿流程
        log.error("send failed, key={}", record.key(), exception);
        failedCounter.increment();
    }
});

3.4 业务数据与消息的一致性:Outbox ​

即使上面全部配置正确,也挡不住环节 ①:订单写库成功后,进程在 send() 之前崩溃,消息就永远不会发出。反过来先发消息再写库,又可能发出一条「订单已创建」但库事务回滚了。

事务型 Outbox 的做法是不在业务事务里发消息,而是把消息写进同一个数据库(投递进程、重放与消费方去重的完整实验见 上下文集成):

sql
BEGIN;
INSERT INTO orders (id, user_id, amount, status) VALUES (1001, 42, 99.00, 'CREATED');
INSERT INTO outbox (event_id, aggregate_id, topic, payload, created_at)
VALUES ('evt-7f3a', 1001, 'order-events', '{"type":"ORDER_CREATED", ...}', NOW());
COMMIT;

之后由独立的投递进程(轮询 outbox 表,或用 CDC 工具读取 Binlog)把消息发到 Kafka,发送成功后标记已投递。投递进程崩溃重启会重发,所以这一段是至少一次,下游仍需幂等。event_id 在写入 Outbox 时生成,重发时保持不变,这是下游去重的依据。

四、消费端:至少一次与幂等 ​

4.1 offset 提交时机决定语义 ​

官方文档对消费端语义的描述很直接:

  • 先保存位置、再处理:处理中崩溃,这批消息不会再被消费,是至多一次;
  • 先处理、再保存位置:保存前崩溃,这批消息会被重新消费,是至少一次。

业务系统几乎总是选择至少一次,因为「重复」可以通过幂等消除,「丢失」通常无法挽回。

需要注意 enable.auto.commit=true(默认值)的行为:自动提交发生在 poll() 调用中,提交的是上一次 poll() 返回的位置。如果在 poll 循环里同步处理完消息再进入下一次 poll(),它也是至少一次;但如果把消息丢进线程池异步处理,下一次 poll() 就可能提交还没处理完的消息,变成至多一次。实测 100 条消息交给线程池、每条处理 50ms,2 秒后进程退出:只处理了 40 条,offset 已经提交到 100,重启后另外 60 条不会再被消费。异步处理时要关闭自动提交,自己管理 offset。

4.2 幂等消费:去重记录与业务写入同一个事务 ​

sql
CREATE TABLE consumed_event (
    consumer_group VARCHAR(64)  NOT NULL,
    event_id       VARCHAR(64)  NOT NULL,
    consumed_at    DATETIME     NOT NULL,
    PRIMARY KEY (consumer_group, event_id)
);
java
@Transactional
public void onOrderCreated(OrderCreatedEvent event) {
    int inserted = jdbc.update("""
            INSERT IGNORE INTO consumed_event (consumer_group, event_id, consumed_at)
            VALUES (?, ?, NOW())
            """, "points-service", event.eventId());
    if (inserted == 0) {
        return;                                   // 已处理过,直接跳过
    }
    pointRepository.add(event.userId(), event.points());
}

要点:

  • 去重记录和业务更新必须在同一个本地事务中。分开写的话,两者之间崩溃又会出现不一致。
  • 幂等键来自事件本身,比如 Outbox 中生成的 event_id 或业务上的「订单号 + 事件类型」,不能用 Kafka 的 offset(重新投递到 Outbox 会产生新 offset),也不能在消费时生成。
  • 如果业务操作天然幂等(如「把订单状态设为已支付」配合状态机校验),可以不用去重表,但要确认乱序重放时的结果依然正确。
  • INSERT IGNORE 是 MySQL 语法,会把一些其他错误降级为警告;对数据校验要求严格时,可以改为普通 INSERT 并捕获唯一键冲突。
  • 去重表需要按时间清理,保留时长应长于消息可能被重放的最长时间(包括 Topic 保留期和人工回放的窗口)。

五、Kafka 的 Exactly Once 覆盖了什么 ​

5.1 读—处理—写链路 ​

当输入和输出都在 Kafka 中时,可以用事务性生产者把两件事变成原子操作:写入输出 Topic,以及提交输入 Topic 的消费 offset。

java
// 消费者:关闭自动提交,只读取已提交事务中的消息
consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
consumerProps.put(ConsumerConfig.ISOLATION_LEVEL_CONFIG, "read_committed");

// 生产者:配置 transactional.id 后自动开启幂等
producerProps.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-enricher-0");

consumer.subscribe(List.of("orders"));
producer.initTransactions();

while (running) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    if (records.isEmpty()) {
        continue;
    }
    producer.beginTransaction();
    try {
        Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
        for (ConsumerRecord<String, String> r : records) {
            producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
            offsets.put(new TopicPartition(r.topic(), r.partition()),
                    new OffsetAndMetadata(r.offset() + 1));
        }
        // offset 作为事务的一部分提交,而不是由消费者单独提交
        producer.sendOffsetsToTransaction(offsets, consumer.groupMetadata());
        producer.commitTransaction();
    } catch (ProducerFencedException e) {
        throw e;                                 // 有同 transactional.id 的新实例,本实例必须退出
    } catch (KafkaException e) {
        producer.abortTransaction();             // 已被隔离时这里会抛 ProducerFencedException,同样退出
        rewind(consumer, records);               // 回退到本批次起点,重新处理
    }
}

这样一来:

  • 事务提交前失败,输出消息和 offset 都不生效,重启后重新处理,下游 read_committed 消费者看不到中止事务中的消息。实测中止一个已经写出 10 条的事务:read_committed 读到 0 条,read_uncommitted 读到 10 条(消息已经在日志里,只是被标记为中止),输入 topic 的 offset 没有前进;
  • transactional.id 让重启后的新实例能「隔离」(fence)旧实例,避免僵尸实例继续写入。旧实例不一定先收到 ProducerFencedException:实测它继续 send 时拿到的是 InvalidProducerEpochException,要到随后 abortTransaction() 才抛出 ProducerFencedException;直接 commitTransaction() 则立即得到 ProducerFencedException。所以上面的代码在 abortTransaction() 失败时也要退出,不能只靠第一个 catch;
  • Kafka Streams 设置 processing.guarantee=exactly_once_v2 后,框架替你完成了上面这些工作。

5.2 边界在哪里 ​

poll()ordersKafka 事务 · 一起提交或一起中止send → orders-enrichedsendOffsetsToTransaction下游消费者read_committed外部副作用接口 · 数据库 · 邮件事务中止后重试,会再执行一次
图 2 · Kafka 事务把输出消息与消费位置绑在一起提交;处理过程中调用的外部系统不在事务边界内

上面的 enrich() 如果调用了外部接口、写了数据库、发了邮件,这些副作用不在 Kafka 事务里。事务中止并重试时,它们会再执行一次。

也就是说,Kafka 的 Exactly Once 保证的是「输出 Topic 中的结果与恰好处理一次相同」,而不是「处理函数只执行一次」。

5.3 输出在外部系统时 ​

官方文档给出的思路是:把 offset 和输出存进同一个外部系统,让两者在一次写入中原子地更新,就不需要两阶段提交。

以写入 MySQL 为例:

  1. 关闭自动提交,在同一个数据库事务中写入业务结果和 (topic, partition, offset);
  2. 分区分配给消费者时(ConsumerRebalanceListener.onPartitionsAssigned),从数据库读出已保存的 offset,调用 consumer.seek() 从那里开始消费。

这个方案严格但侵入性较强。多数业务选择前面的「至少一次 + 幂等消费」,效果相同,实现更简单。

六、顺序与重复的关系 ​

Kafka 只保证分区内有序。需要同一订单的事件有序时,用订单号作为消息 key,让它们进入同一个分区。

有几种情况会破坏这个前提:

  • 消费者内部并行:同一分区的消息被多线程并发处理,顺序就交给了线程调度。可以按 key 哈希到固定的处理线程。
  • 增加分区:key 到分区的映射改变,扩容前后同一 key 的消息分布在两个分区中,过渡期可能乱序。
  • 重试与重放:至少一次语义下,旧消息可能在新消息之后再次出现。幂等处理要能识别「过期」事件,例如比较事件中的版本号或状态机的当前状态,而不仅仅是判断是否处理过。

七、上线前的可靠性清单 ​

配置

  • [ ] 核心 Topic 的 replication.factor ≥ 3,min.insync.replicas = 2
  • [ ] 生产者 acks=all、幂等开启,没有被其他配置意外关闭
  • [ ] unclean.leader.election.enable=false
  • [ ] 异步处理的消费者关闭了自动提交

代码

  • [ ] 生产者回调处理了最终失败,并有告警
  • [ ] 数据库变更与发消息通过 Outbox 保证一致
  • [ ] 消费者的去重记录与业务更新在同一个事务中
  • [ ] 幂等键来自事件内容,重放时保持不变
  • [ ] 死信队列(DLQ)中的消息可以修复后重放,且重放是幂等的

监控与演练

  • [ ] 监控 under-replicated partitions、ISR 收缩、生产错误率、consumer lag
  • [ ] 演练过 Broker 宕机、消费者在处理中被强制终止、重平衡频繁发生

八、常见误区 ​

  • 「开了幂等生产者就不会重复」:它只消除同一实例内的重试重复,消费端的重复与之无关。
  • 「acks=all 就不会丢」:ISR 收缩到只剩 Leader 时等同 acks=1,要配合 min.insync.replicas。
  • 「Exactly Once 意味着处理函数只执行一次」:它保证 Kafka 中的输出结果等价于处理一次,外部副作用可能执行多次。
  • 「用 offset 做幂等键」:同一业务事件被重新投递后 offset 会变。
  • 「DLQ 就是垃圾桶」:进了死信队列的消息需要有人看、能修复、可重放。

小结 ​

讨论 Kafka 的可靠性,第一步是划边界:生产端靠 acks=all、幂等和 min.insync.replicas 保证已确认的消息不丢;消费端靠「先处理后提交」做到至少一次,再用幂等消除重复;Kafka 的 Exactly Once 只覆盖 Kafka 内部的读—处理—写链路。至于数据库里的效果,要靠 Outbox 把业务写入和发消息绑定起来,靠幂等消费把重复变成无害。

消费组扩缩容时哪些消息会被重复处理,见 Kafka 消费组重平衡;acks 对吞吐和延迟的影响,见 Kafka 为什么快。消息队列的整体选型与延时消息,见 消息队列与延时任务;Spring 中本地事务与发消息的时机,见 Spring 事务传播。


配套实验

参考资料

文章以 CC BY-NC-SA 4.0 授权 · 代码片段以 MIT 授权