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。
二、链路拆解:丢与重分别发生在哪
| 环节 | 丢失的典型原因 | 重复的典型原因 |
|---|---|---|
| ① 业务写库与发消息 | 库提交了,发消息前进程崩溃 | 发了消息,库事务回滚 |
| ② 生产者发送 | acks=0/1;回调里的失败没处理;超时后放弃 | 网络超时后重试,而 Broker 实际已写入 |
| ③ Broker 复制 | 只写入 Leader 就确认,Leader 随后宕机;不干净的 Leader 选举 | —— |
| ④⑤⑥ 消费处理 | 先提交 offset,后处理失败 | 处理成功,提交 offset 前崩溃或发生重平衡 |
下面按段展开。
三、生产端与 Broker:已确认的消息不丢
3.1 三个配置要一起看
| 配置 | 作用于 | 默认值 | 建议 |
|---|---|---|---|
acks | 生产者 | all(3.0 起) | 保持 all |
enable.idempotence | 生产者 | true(3.0 起,无冲突配置时) | 保持开启 |
replication.factor | Topic | 自动创建时取 default.replication.factor,默认为 1 | 生产环境至少 3 |
min.insync.replicas | Topic / Broker | 1 | 与 3 副本搭配设为 2 |
unclean.leader.election.enable | Topic / Broker | false | 需要强一致时保持 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,如果回调里的异常没人处理,丢了也不知道:
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 的做法是不在业务事务里发消息,而是把消息写进同一个数据库(投递进程、重放与消费方去重的完整实验见 上下文集成):
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 幂等消费:去重记录与业务写入同一个事务
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)
);@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。
// 消费者:关闭自动提交,只读取已提交事务中的消息
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 边界在哪里
上面的 enrich() 如果调用了外部接口、写了数据库、发了邮件,这些副作用不在 Kafka 事务里。事务中止并重试时,它们会再执行一次。
也就是说,Kafka 的 Exactly Once 保证的是「输出 Topic 中的结果与恰好处理一次相同」,而不是「处理函数只执行一次」。
5.3 输出在外部系统时
官方文档给出的思路是:把 offset 和输出存进同一个外部系统,让两者在一次写入中原子地更新,就不需要两阶段提交。
以写入 MySQL 为例:
- 关闭自动提交,在同一个数据库事务中写入业务结果和
(topic, partition, offset); - 分区分配给消费者时(
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 事务传播。
配套实验
- codesphere-labs/messaging/kafka-delivery:事务的中止与提交、
transactional.id隔离旧实例、自动提交 + 异步处理丢消息、min.insync.replicas与 ISR 收缩(验证记录) - 消费方按事件 ID 去重、与业务写入放在同一事务:见 上下文集成 的配套实验
参考资料