【今天下午正准备摸鱼,测试突然在群里拍我:“订单状态乱跳,先成功又回滚,但消费者明明说处理过了?”我一看监控,Kafka消费量正常,对账却有缺口,脑壳嗡的一声——多半是读到了不该看的‘鬼消息’。}


事故现场

  • 现象:
    • 下游“订单状态同步”消费者处理了几条后来被撤销的事件,导致库里出现“已支付->又被覆盖成未支付”的离谱状态;
    • 同一时间段内,生产者日志里有事务相关异常;
    • 但消费者没报错,还挺勤奋地消费了一堆。

日志随手一截:

2026-03-21 14:02:31 INFO  OrderProducer - begin txn for order=1008611
2026-03-21 14:02:33 ERROR OrderProducer - txn aborted
org.apache.kafka.common.errors.TransactionAbortedException: Transaction aborted

消费者侧却在嘀咕:

2026-03-21 14:02:34 INFO  OrderSyncListener - handled event {orderId=1008611, status=PAID}

这就诡异了:生产者把事务_abort_了,消费者却把消息当真的处理了。


排查脑回路

  1. 先看生产者:这次我们为了“精准一次”(EOS)给生产者开了事务(transactional.id),链路里有多条topic需要“要么都成功,要么都撤销”。
  2. 生产阶段某个数据库校验失败触发了TransactionAbortedException,按理整笔事务内的消息都应被过滤掉,不该被消费者看见。
  3. 然后我翻了消费者配置,当场裂开:消费者默认isolation.level=read_uncommitted,也就是“读未提交”。这玩意会把失败事务里产生的record也投递给你,下游还一本正经落库,直接制造数据屎山。
  4. 侧证:把同分区的相同事件抓出来对比,生产端那条对应的producerId在事务状态日志里是aborted,但消费者依然收到了。板上钉钉:开了生产事务,却忘了把消费者设成read_committed。

问题代码(反例)

生产端(Spring Kafka)开了事务,但配置不完整:

@Bean
public ProducerFactory<String, String> pf() {
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
    props.put(ProducerConfig.ACKS_CONFIG, "all");
    props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
    props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-001");
    return new DefaultKafkaProducerFactory<>(props);
}

@Bean
public KafkaTemplate<String, String> kt(ProducerFactory<String, String> pf) {
    KafkaTemplate<String, String> kt = new KafkaTemplate<>(pf);
    kt.setTransactionIdPrefix("order-tx-");
    return kt;
}

// 业务里用事务发多条消息
public void publishOrderEvents(Order order) {
    kafkaTemplate.executeInTransaction(ops -> {
        ops.send("order-status", key(order), toJson(order));
        ops.send("stock-reserve", key(order), toJson(order));
        // 期间抛异常 -> 整个事务 abort
        validateOrThrow(order);
        return true;
    });
}

消费者(反例)默认读未提交:

spring:
  kafka:
    consumer:
      group-id: order-sync
      enable-auto-commit: false
      # 没设置isolation-level,默认read_uncommitted(坑)

Listener:

@KafkaListener(topics = "order-status", groupId = "order-sync")
public void onMsg(ConsumerRecord<String, String> r, Acknowledgment ack) {
    OrderEvent e = fromJson(r.value());
    syncService.apply(e); // 这里把“已被abort的消息”也处理了
    ack.acknowledge();
}

根因

  • 生产者开启了事务(EOS路径),但消费者没设置read_committed,导致读到了被abort的事务内消息;
  • 一些主题混用了“事务生产者”和“非事务生产者”,数据边界更乱;
  • 少量分区里min.insync.replicas偏低,网络抖动时NotEnoughReplicas触发重试,进一步诱发事务abort。

解决方案(一次到位)

1) 消费者必须read_committed + 手动提交

spring:
  kafka:
    consumer:
      enable-auto-commit: false
      isolation-level: read_committed   # 关键!只读已提交的事务消息
      max-poll-interval-ms: 600000
      max-poll-records: 200
@KafkaListener(topics = "order-status", groupId = "order-sync")
public void onMsg(ConsumerRecord<String, String> r, Acknowledgment ack) {
    try {
        OrderEvent e = fromJson(r.value());
        syncService.applyWithIdempotency(e); // 幂等兜底
        ack.acknowledge();
    } catch (Exception ex) {
        // 失败交给错误处理器/重试+DQL
        throw ex;
    }
}

2) 生产者完整配置EOS(事务+幂等)

@Bean
public ProducerFactory<String, String> txProducerFactory() {
    Map<String, Object> props = new HashMap<>();
    props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka:9092");
    props.put(ProducerConfig.ACKS_CONFIG, "all");
    props.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, true);
    props.put(ProducerConfig.RETRIES_CONFIG, 5);
    props.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, 5); // 避免乱序
    props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, "order-tx-${HOSTNAME}");
    return new DefaultKafkaProducerFactory<>(props);
}

public void publishInTxn(List<ProducerRecord<String, String>> records) {
    kafkaTemplate.executeInTransaction(ops -> {
        records.forEach(ops::send);
        // 可选:在同一事务里写本地Outbox(见下一节)
        return true;
    });
}

3) 主题治理与Broker侧

  • 不要在同一个topic/partition里混用事务生产者和非事务生产者;
  • Broker参数:
    • min.insync.replicas >= 2(配合acks=all防单副本写入);
    • transaction.state.log.replication.factor >= 3,transaction.state.log.min.isr >= 2;
  • 为关键topic开log.cleanup.policy=compact或严格保留策略,避免幂等键的墓碑过早清理。

4) 幂等与补偿(消费侧/生产侧都要有)

  • 消费侧:按bizId/订单号做唯一键去重,重复消息也不怕;
ALTER TABLE t_order_sync ADD UNIQUE KEY uk_biz (biz_id);
  • 生产侧:即使事务提交失败或模板外异常,也要有Outbox兜底,再投递(隔离出站可靠性)。

验证

  • 配置read_committed后,复现同样的事务abort,消费者不再收到那些record;
  • 对账恢复一致,无“先成功后回滚还被处理”的鬼畜状态;
  • 压测1小时,消费速率稳定,重试/死信量在可控范围。

踩坑总结

  • 只开生产事务不改消费者=自欺欺人;read_committed是下游的第一道门神。
  • EOS不等于万无一失:幂等与Outbox依然要上,Broker副本/ISR也要兜住尾部风险。
  • 看见“事务abort但消息被处理”,第一时间检查消费者隔离级别和主题混用情况,少走弯路。
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐