Kafka 的 Exactly-Once 语义听着很美,在跨系统场景下我还是老老实实上了业务去重表

Kafka 的 Exactly-Once 语义(EOS)只覆盖 Kafka 自身的生产-消费闭环,一旦消费端需要往外部数据库、缓存或第三方 API 写数据,那个“恰好一次”的承诺就失效了。你必须在业务层自己兜底,而业务去重表是跨系统场景下最靠谱的兜底方案。

我刚踩完这个坑。一个订单系统从 Kafka 消费支付结果,然后写入 MySQL 订单表并调微信支付的分账接口。开了 Kafka EOS,以为万事大吉,结果一次 broker 切换导致 consumer rebalance,同一笔支付回调被处理了两遍,分账接口被调了两次,财务对账全乱了。事后复盘,根因很清楚:Kafka 的 Exactly-Once 语义压根没管到下游外部系统。

Kafka EOS 到底管到哪一步

Kafka 的 Exactly-Once 语义是 2017 年 6 月在 0.11.0 版本引入的,核心靠三个机制协同:幂等生产者、事务性生产者和事务性消费者。

幂等生产者给每个 producer session 分配一个 producer ID,每条消息带一个单调递增的 sequence number。Broker 发现重复的 PID + sequence number 就直接丢弃,能防止网络重试导致的消息重复写入。事务性生产者在此基础上加了跨分区原子写入的能力,通过 transaction coordinator 管理两阶段提交。事务性消费者则配合 isolation.level=read_committed,只读取已提交的事务消息,并且可以把消费位移的提交也纳入同一个事务。

问题就出在最后一步。Kafka 事务的边界止于 consumer offset 的提交。你在 poll() 拿到消息之后,把数据处理逻辑写进 KafkaProducer.send() 的事务块里,Kafka 能保证的是:消息处理成功 → offset 提交,消息处理失败 → offset 不提交,消息会被重新消费。这个语义确实消除了“消息丢了但 offset 已提交”或者“消息重复但 offset 只提交了一次”的中间状态。

但你的数据库写入、API 调用、缓存更新,这些外部操作 Kafka 根本感知不到。事务 coordinator 不知道你往 MySQL 插了条数据,也不知道你调了微信支付的 HTTP 接口。如果消费者在数据库写入成功后、offset 提交前崩溃了,重启后消息会被重新消费,数据库写入再来一次——Kafka EOS 对此毫无办法,因为它的“Exactly-Once”保证只存在于 Kafka 内部。

跨系统场景下的真实失效路径

拿我那个订单系统拆开看。Consumer 的逻辑大致是这样:

@KafkaListener(topics = "payment_result")
@Transactional(transactionManager = "kafkaTransactionManager")
public void onMessage(ConsumerRecord<String, PaymentResult> record) {
    PaymentResult result = record.value();
    // 1. 更新订单状态
    orderService.updateStatus(result.getOrderId(), result.getStatus());
    // 2. 如果是支付成功,调分账接口
    if ("SUCCESS".equals(result.getStatus())) {
        paymentService.profitSharing(result.getOrderId(), result.getAmount());
    }
    // 3. 发送确认消息到下游
    kafkaTemplate.send("order_confirmed", buildConfirmEvent(result));
}

这段代码在 Kafka 的事务保护下运行,EOS 能保证步骤 3 的 send 和 offset 提交是原子的。但步骤 1 和 2 是跨系统的。如果步骤 1 成功、步骤 2 也成功了,然后 JVM 在步骤 3 执行前被 OOM Killer 干掉,Kafka 事务回滚,消息被重新投递。下一次消费时,步骤 1 和 2 再来一遍,重复扣款或者重复分账。

还有个更隐蔽的场景。步骤 2 调分账接口,HTTP 请求发出去了,微信那边也处理成功了,但网络抖动导致响应超时,consumer 这边以为失败了,抛异常触发事务回滚。消息重试,分账又调一次。微信那边没有幂等保证的话,就是实打实的重复分账。

Kafka EOS 解决的是 Kafka 内部的重复问题,但跨系统交互中的重复,根因在外部系统的非幂等性上。你没法要求微信支付的分账接口支持幂等(实际上它确实不支持),也没法让 MySQL 的 UPDATE 天然识别重复的业务请求。这些都需要业务层自己处理。

业务去重表为什么是正解

业务去重表的思路简单到粗暴:在消费消息之前,先把消息的唯一标识记下来,处理完标记为已完成。下次同样的消息来了,查一下表就知道已经处理过了,直接跳过。

但简单不代表简陋。恰恰因为它把幂等控制的逻辑完全放在业务数据库里,和你的业务数据在同一事务中,才能做到真正意义上的“恰好一次处理”。

具体实现上,我通常建一张 idempotency_records 表:

CREATE TABLE idempotency_records (
    idempotency_key VARCHAR(128) PRIMARY KEY,
    status ENUM('PROCESSING', 'COMPLETED') NOT NULL DEFAULT 'PROCESSING',
    created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
    completed_at TIMESTAMP NULL,
    INDEX idx_created_at (created_at)
);

消息的唯一标识怎么取?Kafka 消息本身有一个复合唯一键:topic + partition + offset。你可以直接用这三个字段拼成 idempotency_key。但要注意,如果你开了日志压缩(log compaction),offset 可能会变,更稳妥的做法是让生产者把业务唯一 ID 放进消息体里,比如支付回调的 transaction_id 或者订单的 order_id

消费逻辑变成这样:

@Transactional(transactionManager = "dataSourceTransactionManager")
public void processPaymentResult(PaymentResult result, String topic, int partition, long offset) {
    String idempotencyKey = topic + ":" + partition + ":" + offset;
    
    // 插入去重记录,利用唯一索引保证原子性
    try {
        jdbcTemplate.update(
            "INSERT INTO idempotency_records (idempotency_key, status) VALUES (?, 'PROCESSING')",
            idempotencyKey
        );
    } catch (DuplicateKeyException e) {
        // 检查是否已完成
        String status = jdbcTemplate.queryForObject(
            "SELECT status FROM idempotency_records WHERE idempotency_key = ?",
            String.class, idempotencyKey
        );
        if ("COMPLETED".equals(status)) {
            log.info("消息已处理,跳过: {}", idempotencyKey);
            return;
        }
        // 状态是 PROCESSING 说明上一次处理卡住了,需要判断是否重试
        // 这里根据业务决定是等待、重试还是人工介入
        throw new IllegalStateException("消息处理中,可能存在并发冲突: " + idempotencyKey);
    }
    
    // 执行业务逻辑,和去重表在同一个数据库事务里
    orderService.updateStatus(result.getOrderId(), result.getStatus());
    if ("SUCCESS".equals(result.getStatus())) {
        paymentService.profitSharing(result.getOrderId(), result.getAmount());
    }
    
    // 标记完成
    jdbcTemplate.update(
        "UPDATE idempotency_records SET status = 'COMPLETED', completed_at = NOW() WHERE idempotency_key = ?",
        idempotencyKey
    );
}

整个逻辑包在 Spring 的 DataSourceTransactionManager 管理的事务里。如果数据库写入成功但 HTTP 调用失败,事务回滚,去重记录也被回滚,消息重试时能正常进入处理流程。如果全部成功,去重记录标记为 COMPLETED,后续重复消息被唯一索引挡住,直接跳过。

这里有一个关键细节:去重表的插入必须在业务逻辑之前,并且利用数据库唯一索引的原子性做并发控制。如果你先查后插,在并发场景下查和插之间会有竞态窗口。唯一索引 + INSERT + DuplicateKeyException 捕获,是数据库层面能提供的最强保证。

两种方案的取舍边界

Kafka EOS 和业务去重表不是二选一的对立关系,它们的覆盖范围不一样。

Kafka EOS 的适用场景非常窄:当你只需要把 Kafka 里的数据做流式处理后写回 Kafka,整个链路不涉及外部系统时,EOS 是首选。典型的比如 Kafka Streams 做聚合计算,或者从一个 topic 消费后转换格式写入另一个 topic。这种场景下 EOS 开箱即用,不需要额外维护去重状态。

但一旦消费逻辑里包含对数据库、缓存、第三方 API 的写入,EOS 的边界就到了。此时你必须上业务层幂等,而业务去重表是最通用的解法。

成本上,业务去重表需要额外维护一张表,占用存储空间,并且每次消费多一次 INSERT 和一次 UPDATE。这些开销在大部分业务场景下完全可以接受——去重表的写入量级和业务表本身在一个数量级上,MySQL 的 INSERT + 唯一索引的性能瓶颈远在业务逻辑的瓶颈之下。真正需要担心的是去重表的数据膨胀问题,我的做法是保留 7 天的去重记录,超过 7 天的用定时任务清理。Kafka 消息重试通常不会超过几小时,7 天已经是相当保守的窗口了。

还有个容易被忽略的点:去重表的 PROCESSING 状态可以作为一种分布式锁,防止同一消息被多个 consumer 实例并发处理。Kafka 在 rebalance 期间可能会出现短暂的分区所有权重叠,两个 consumer 同时拿到同一条消息。去重表的唯一索引天然防止了这种情况。

常见问题

去重表用 Redis 行不行?

不建议。Redis 和业务数据库是两个独立的存储,你没法把 Redis 的 SETNX 操作和 MySQL 的业务写入放在同一个事务里。如果在 Redis 标记成功、MySQL 写入失败时回滚 Redis 标记,中间窗口期就会有数据不一致的风险。如果你非要减轻数据库压力,可以用 Redis 做前置过滤,命中了直接跳过,但最终的一致性兜底还是得靠数据库去重表。

Kafka 事务性能开销有多大,值不值得开?

Kafka 事务有实实在在的开销。每次 commitTransaction 都需要和 transaction coordinator 做一轮网络往返,吞吐量大概下降 20% 到 30%。如果你的消费逻辑里已经上了业务去重表,Kafka 侧的 Exactly-Once 语义其实可以关掉,用普通的 at-least-once 消费 + 业务去重表就够了。Kafka 事务在纯 Kafka 链路的流处理场景下才有不可替代的价值。

去重表的 idempotency_key 用业务 ID 还是 Kafka offset?

优先用业务 ID,比如订单号、交易流水号。这样即使消息从不同的 topic 或分区进来,只要业务 ID 相同就能去重。Kafka offset 的问题是,如果消息被重新投递到不同的分区(比如因为分区数变更或者生产者重试策略),offset 会变,去重就失效了。业务 ID 的稳定性远高于 Kafka 的内部标识。