Spring Kafka事务异常:IN_TRANSACTION状态非法转换问题排查
兄弟,你遇到的这个Invalid transition attempted from state IN_TRANSACTION to state IN_TRANSACTION异常我太熟了!本质就是重复开启事务搞出来的坑。咱们用的Spring Kafka 2.1.4版本有个小缺陷:当监听器方法已经处于消费者事务上下文里时,如果你用事务性的KafkaTemplate发消息,默认配置下框架会“傻傻地”再开一个事务,可Kafka的事务状态机根本不允许从IN_TRANSACTION状态再次进入该状态,直接就把事务回滚了。
给你三个解决方案,按需选择
1. 让KafkaTemplate复用现有事务(最省心)
改一下KafkaTemplate的配置,让它乖乖复用监听器已经启动的事务,别自己瞎折腾。你可以在配置类里把事务管理器绑定给KafkaTemplate:
@Bean public KafkaTemplate<String, Object> kafkaTemplate(ProducerFactory<String, Object> producerFactory, KafkaTransactionManager<String, Object> transactionManager) { KafkaTemplate<String, Object> template = new KafkaTemplate<>(producerFactory); template.setTransactionManager(transactionManager); return template; }
这样KafkaTemplate会自动识别当前已存在的事务上下文,不会重复开启事务,完美避开状态冲突。
2. 手动控制监听器的事务传播(不推荐,太繁琐)
如果实在不想改配置,你可以把监听器的事务传播属性设为NEVER,然后自己手动管理事务的提交和回滚,但这种方式容易出错,除非万不得已别用:
@KafkaListener(topics = "your-topic", containerFactory = "kafkaListenerContainerFactory") @Transactional(propagation = Propagation.NEVER) public void listen(ConsumerRecord<String, Object> record) { TransactionStatus status = transactionManager.getTransaction(new DefaultTransactionDefinition()); try { kafkaTemplate.send("another-topic", "your-message"); transactionManager.commit(status); } catch (Exception e) { transactionManager.rollback(status); } }
3. 升级Spring Kafka版本(最推荐,一劳永逸)
其实Spring Kafka 2.2.x及以后的版本已经修复了这个问题,框架会自动检测当前是否存在事务上下文,不会盲目开启新事务。如果你的项目允许升级,直接升到2.2.x或更高的兼容版本,不仅能解决这个问题,还能拿到一堆新特性和bug修复,香得很。
最后再提醒几个配置细节
别忘把生产者和消费者的事务参数配到位,比如application.properties里得有这些:
# 生产者事务前缀,必须配置 spring.kafka.producer.transaction-id-prefix=tx- # 关闭消费者自动提交,交给事务管理 spring.kafka.consumer.enable-auto-commit=false # 消费者隔离级别设为读已提交,符合事务语义 spring.kafka.consumer.isolation-level=read_committed
还有监听器容器工厂一定要绑定事务管理器,不然监听器的事务根本起不来:
@Bean public ConcurrentKafkaListenerContainerFactory<String, Object> kafkaListenerContainerFactory( ConsumerFactory<String, Object> consumerFactory, KafkaTransactionManager<String, Object> transactionManager) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); factory.getContainerProperties().setTransactionManager(transactionManager); return factory; }
内容的提问来源于stack exchange,提问作者Eugene Khyst

