Administrator
发布于 2022-02-20 / 8430 阅读
142

Kafka 精准一次语义的实现与代价

对账群里跳出一条消息:同一笔订单扣了两次款

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 精准一次语义的实现与代价》本身不算多难,难的是线上真出问题那十分钟里的判断。经验都是这么来的。

参考