Spring Kafka事务已启用,消费者仍收到回滚消息问题排查
嘿,我仔细看了你的代码,发现核心问题是生产者端同时用了两种冲突的事务管理方式,导致Kafka消息的事务没有和数据库事务绑定,最终数据库回滚时,Kafka消息已经被提前提交了,所以即使消费者设置了read_committed还是能收到这条消息。
问题到底出在哪?
你的生产者服务里同时用了@Transactional(transactionManager = "chainedTransactionManager")和kafkaTemplate.executeInTransaction(...),这俩是互斥的:
@Transactional搭配链式事务管理器,是想让Kafka和数据库事务统一由Spring管理,要么一起提交,要么一起回滚- 但
kafkaTemplate.executeInTransaction()会单独开启一个独立的Kafka事务上下文,它内部会自己完成提交/回滚,完全不受外层链式事务的控制
所以当你调用executeInTransaction发送消息后,这个Kafka事务已经提交了,之后数据库保存失败触发回滚时,Kafka那边早就把消息落盘了,消费者自然能读到。
修复步骤(一步步来)
1. 删掉executeInTransaction,改用普通send
把生产者服务里的kafkaTemplate.executeInTransaction(...)换成普通的kafkaTemplate.send(),让链式事务管理器来统一控制Kafka事务的生命周期:
修改后的生产者服务代码:
@Component @Slf4j public class KafkaProducerService { private final KafkaTemplate<String, Customer> kafkaTemplate; private final CustomerConverter customerConverter; private final CustomerRepository customerRepository; // 用构造器注入(建议加final,更安全) public KafkaProducerService(KafkaTemplate<String, Customer> kafkaTemplate, CustomerConverter customerConverter, CustomerRepository customerRepository) { this.kafkaTemplate = kafkaTemplate; this.customerConverter = customerConverter; this.customerRepository = customerRepository; } @Transactional(transactionManager = "chainedTransactionManager", rollbackFor = Exception.class) public void sendEvents(String topic, CustomerModel customer) { LOGGER.info("Sending to Kafka: topic: {}, key: {}, customer: {}", topic, customer.getKey(), customer); // 移除executeInTransaction,改用普通send,交给链式事务管理 kafkaTemplate.send(topic, customer.getKey(), customerConverter.convertToAvro(customer)); // 这里抛出异常时,整个链式事务会回滚,包括Kafka消息 customerRepository.saveToDb(); } }
2. 确认链式事务管理器的顺序没问题
你的链式事务管理器配置是new ChainedTransactionManager(kafkaTransactionManager, jpaTransactionManager),这个顺序是对的!
链式事务的提交顺序是反向的(先提交JPA,再提交Kafka),回滚顺序是正向的(先回滚Kafka,再回滚JPA),这样能避免“Kafka提交成功但数据库提交失败”的不一致场景,这个配置不用改。
3. 其他配置的验证(已经没问题的部分)
- 生产者的
transactionIdPrefix配置正确,这是Kafka事务的必备项,确保每个生产者实例有唯一的事务ID - 消费者的
isolation.level=read_committed配置正确,它会让消费者只读取已提交的事务消息,事务回滚的消息会被过滤掉 - 生产者的
acks=all、enable.idempotence=true这些基础配置也都没问题,是事务的保障
额外提醒
- 确保你的Kafka集群版本是0.11.0.0及以上,因为Kafka事务是从这个版本开始支持的
- 绝对不要混合使用
executeInTransaction和外层的@Transactional,这会导致事务上下文混乱,出现你现在遇到的问题
这样修改后,当saveToDb()抛出异常时,链式事务管理器会同时回滚Kafka事务和数据库事务,Kafka消息不会被提交,消费者自然就收不到这条消息了。
内容的提问来源于stack exchange,提问作者ttt

