对账群里跳出一条消息:同一笔订单扣了两次款
2 月 19 号下午,财务在对账时发现一笔订单出现了两次支付成功记录,金额都是 199 元。看起来是消费者把消息重复处理了。
我们用的 Kafka 是 3.1.0,消费者是 enable.auto.commit=true 默认配置。同事问我:不是说 Kafka 是"至少一次"吗,那重复消费不是正常的?为什么会扣两次款而不是一次都没扣?
这个问题把我问住了。我决定把 Kafka 的语义彻底理一遍,顺便把"精准一次"在我们这套系统里到底能不能做到,写清楚。
先搞清楚三个语义
Kafka 的投递语义有三个层级,容易混:
| 语义 | 含义 | 典型后果 |
|---|---|---|
| 最多一次(at-most-once) | 消息可能丢,不会重复 | 下单成功但没记日志 |
| 至少一次(at-least-once) | 消息不丢,可能重复 | 同一笔订单被处理两遍 |
| 精准一次(exactly-once) | 精确处理一次 | 既不多也不少 |
默认配置下,Spring Kafka 的 @KafkaListener 用的是 auto.commit:poll 回来之后,还没处理完就提交了 offset。一旦处理到一半进程挂了,这批消息的 offset 已经提交了,重启后从新 offset 开始,中间的就丢了——这是最多一次。
我们当时为了不丢消息,关了自动提交、改成处理完手动 ack。这下倒过来了:处理完了但还没 ack 就挂,重启会重新消费这批——这是至少一次,重复就来自这里。
幂等生产者:先把"发重"这条路堵死
重复其实有两层。第一层是生产者自己重发导致的重复。比如 broker 写成功了,但返回 ack 的网络包丢了,生产者重试,同一条消息就落了两份。
Kafka 的幂等生产者解决的就是这个。开启方式就一行:
Properties props = new Properties();
props.put("bootstrap.servers", "10.0.0.41:9092");
props.put("enable.idempotence", true); // 关键
props.put("acks", "all");
props.put("retries", Integer.MAX_VALUE);
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
开启后,broker 给每个生产者分配一个 producerId(PID),每条消息带一个单调递增的 sequence。broker 端按 (PID, partition) 维护最新 sequence,收到重复(sequence 不连续或回退)的消息直接丢弃。
我本地压测验证了一次:故意在发送中途把网络掐断,让生产者触发重试。开幂等前,topic 里出现 1042 条重复;开幂等后,消费端去重后数量与实际发送数完全一致,0 重复。
但注意:幂等只保证"单分区、单会话"内不重复。换分区、或者生产者重启(PID 变了,旧 PID 的状态在 broker 端只有短保留)就保不住了。我们那笔重复扣款不是生产者重发造成的,是消费者重处理,所以幂等没帮上忙。
事务 API:consume-transform-produce 的原子性
我们的场景是经典的"读一个 topic、处理、写另一个 topic":从 PAY_ORDER 读支付消息,校验后写 PAY_RESULT。这条链路里,要保证"读到的 offset 提交"和"写出结果"是原子的——要么都成功,要么都失败重来。
用事务 API 才能做到:
props.put("enable.idempotence", true);
props.put("transactional.id", "pay-sink-01"); // 必须全局唯一
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.initTransactions();
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
if (records.isEmpty()) continue;
producer.beginTransaction();
try {
for (ConsumerRecord<String, String> r : records) {
String result = process(r.value());
producer.send(new ProducerRecord<>("PAY_RESULT", r.key(), result));
}
// 把消费 offset 纳入事务一起提交
producer.sendOffsetsToTransaction(
currentOffsets(records), consumer.groupMetadata());
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction(); // 回滚,offset 不提交,下次重读
}
}
关键点有两个。一是 transactional.id 必须稳定且唯一,它让 broker 能把"这个事务"和"这个生产者实例"绑定,恢复时知道哪些事务该补提交、哪些该回滚。二是 sendOffsetsToTransaction 把消费位移和产出消息放进同一个事务,外部读者用 isolation.level=read_committed 才看得到已提交的结果。
我用 3.1.0 跑了一个 kill -9 测试:在 commitTransaction 之前杀进程。重启后,那批 PAY_RESULT 消息对 read_committed 消费者不可见,offset 也没前进,重读重处理,落库次数正好是一次。这才是真正的 exactly-once 在 Kafka 内部。
真正难的是:外部系统怎么一起原子
上面那套只在"消息进、消息出"都发生在 Kafka 里时成立。我们扣款是写数据库的,问题在这里:
事务提交的是"Kafka 内部的 offset 和产出消息"。数据库的那笔 UPDATE,不在这个事务里。两者之间没有分布式事务协调,无法保证原子。
我画了时间线给同事看:
写 DB(扣款成功)
→ 还没 commitTransaction
→ 进程挂了
→ DB 已落(199 元扣了)
→ Kafka 事务回滚,offset 没前进
→ 重启重读这条消息
→ 又扣一次 ← 重复扣款就是这个
这就是 hint 里说的"外部系统写入的原子性难题"。Kafka 的 EOS 管不了 Kafka 之外的事。
业界通用的解法是幂等消费或事务性发件箱(Transactional Outbox):
- 幂等消费:在业务表上加唯一约束(比如
pay_no唯一索引),重复处理时数据库报DuplicateKeyException,catch 掉当成功。我们最后用的就是这个,改动最小。 - 发件箱模式:不在消费者里直接写 DB,而是把"要做的变更"作为一条消息写进数据库同一事务,再由一个可靠投递者把消息发到 Kafka。这样 DB 写入和"待发消息"原子,下游靠 Kafka 事务保证。
我们选了幂等消费,因为支付表本来就有 pay_no 唯一索引。改造后的消费代码:
@KafkaListener(topics = "PAY_ORDER", groupId = "pay-group")
public void consume(ConsumerRecord<String, String> r, Acknowledgment ack) {
try {
payService.deduct(r.value()); // 唯一索引兜底
ack.acknowledge();
} catch (DuplicateKeyException e) {
log.warn("重复支付消息, payNo={}, 跳过", parsePayNo(r.value()));
ack.acknowledge(); // 重复也提交,别卡住
}
}
写在后面
现在回头看,《Kafka 精准一次语义的实现与代价》本身不算多难,难的是线上真出问题那十分钟里的判断。经验都是这么来的。