寻找已废弃ChainedKafkaTransactionManager的替代实现方案
替代已弃用的ChainedKafkaTransactionManager实现Kafka与JPA事务联动
我之前基于ChainedKafkaTransactionManager实现了同时管理KafkaTransactionManager和JpaTransactionManager的功能,但该类已被标记为@Deprecated,希望找到同等功能的替代方案。
原实现代码:
@EnableKafka @Configuration @RequiredArgsConstructor public class DunningCycleKafkaConfiguration { private final KafkaConfigurationProperties kafkaConfigurationProperties; @Bean public ConsumerFactory<String, String> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, kafkaConfigurationProperties.getConsumer().getEnableAutoCommit()); props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaConfigurationProperties.getConsumer().getGroupId()); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, kafkaConfigurationProperties.getConsumer().getKeyDeserializer()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, kafkaConfigurationProperties.getConsumer().getValueDeserializer()); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfigurationProperties.getBootstrapServers()); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, kafkaConfigurationProperties.getConsumer().getMaxPollRecords()); return new DefaultKafkaConsumerFactory<>(props); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory( AfterRollbackProcessor<Object, Object> processor, ChainedKafkaTransactionManager<Object, Object> chainedKafkaTransactionManager) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setRecordInterceptor(new KafkaConsumerInterceptor()); ContainerProperties containerProps = factory.getContainerProperties(); containerProps.setAckMode(ContainerProperties.AckMode.valueOf(kafkaConfigurationProperties.getListener().getAckMode())); factory.setAfterRollbackProcessor(processor); factory.getContainerProperties().setTransactionManager(chainedKafkaTransactionManager); return factory; } @Bean public JpaTransactionManager transactionManager(EntityManagerFactory entityManagerFactory) { return new JpaTransactionManager(entityManagerFactory); } @Bean public ChainedKafkaTransactionManager<Object, Object> chainedKafkaTransactionManager( JpaTransactionManager transactionManager, KafkaTransactionManager<?, ?> kafkaTransactionManager) { return new ChainedKafkaTransactionManager<>(kafkaTransactionManager, transactionManager); } @Bean @Primary public KafkaTransactionManager<Object, Object> kafkaTransactionManager(ProducerFactory<Object, Object> producerFactory) { return new KafkaTransactionManager<>(producerFactory); } @Bean public AfterRollbackProcessor<Object, Object> processor() { DefaultAfterRollbackProcessor<Object, Object> processor = new DefaultAfterRollbackProcessor<>( new FixedBackOff(1000L, 3L)); processor.addNotRetryableExceptions(DataIntegrityViolationException.class); processor.addNotRetryableExceptions(IllegalStateException.class); processor.addNotRetryableExceptions(RestClientException.class); processor.addNotRetryableExceptions(NullPointerException.class); processor.addNotRetryableExceptions(NumberFormatException.class); processor.addNotRetryableExceptions(IllegalArgumentException.class); processor.addNotRetryableExceptions(NoSuchMethodException.class); processor.addNotRetryableExceptions(JsonParseException.class); processor.addNotRetryableExceptions(MessageConversionException.class); return processor; } }
我尝试过给事务管理器Bean添加@Transactional注解,但不确定这种方式是否正确:
@Bean @Primary @Transactional public KafkaTransactionManager<Object, Object> kafkaTransactionManager(ProducerFactory<Object, Object> producerFactory) { return new KafkaTransactionManager<>(producerFactory); } @Bean @Transactional public JpaTransactionManager transactionManager(EntityManagerFactory entityManagerFactory) { return new JpaTransactionManager(entityManagerFactory); }
正确的替代方案:使用Spring官方的ChainedTransactionManager
Kafka的ChainedKafkaTransactionManager被弃用后,官方推荐使用Spring TX模块提供的org.springframework.transaction.support.ChainedTransactionManager来实现多事务管理器的联动。这个类可以将多个事务管理器按顺序链式管理,保证事务的一致性(要么全部提交,要么全部回滚)。
修改后的完整配置代码:
@EnableKafka @Configuration @RequiredArgsConstructor public class DunningCycleKafkaConfiguration { private final KafkaConfigurationProperties kafkaConfigurationProperties; @Bean public ConsumerFactory<String, String> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, kafkaConfigurationProperties.getConsumer().getEnableAutoCommit()); props.put(ConsumerConfig.GROUP_ID_CONFIG, kafkaConfigurationProperties.getConsumer().getGroupId()); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, kafkaConfigurationProperties.getConsumer().getKeyDeserializer()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, kafkaConfigurationProperties.getConsumer().getValueDeserializer()); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaConfigurationProperties.getBootstrapServers()); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, kafkaConfigurationProperties.getConsumer().getMaxPollRecords()); return new DefaultKafkaConsumerFactory<>(props); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory( AfterRollbackProcessor<Object, Object> processor, ChainedTransactionManager chainedTransactionManager) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.setRecordInterceptor(new KafkaConsumerInterceptor()); ContainerProperties containerProps = factory.getContainerProperties(); containerProps.setAckMode(ContainerProperties.AckMode.valueOf(kafkaConfigurationProperties.getListener().getAckMode())); factory.setAfterRollbackProcessor(processor); factory.getContainerProperties().setTransactionManager(chainedTransactionManager); return factory; } @Bean public JpaTransactionManager jpaTransactionManager(EntityManagerFactory entityManagerFactory) { return new JpaTransactionManager(entityManagerFactory); } @Bean @Primary public KafkaTransactionManager<Object, Object> kafkaTransactionManager(ProducerFactory<Object, Object> producerFactory) { return new KafkaTransactionManager<>(producerFactory); } @Bean public ChainedTransactionManager chainedTransactionManager(KafkaTransactionManager<Object, Object> kafkaTransactionManager, JpaTransactionManager jpaTransactionManager) { // 注意顺序:先Kafka事务管理器,后JPA事务管理器,和原ChainedKafkaTransactionManager保持一致 return new ChainedTransactionManager(kafkaTransactionManager, jpaTransactionManager); } @Bean public AfterRollbackProcessor<Object, Object> processor() { DefaultAfterRollbackProcessor<Object, Object> processor = new DefaultAfterRollbackProcessor<>( new FixedBackOff(1000L, 3L)); processor.addNotRetryableExceptions(DataIntegrityViolationException.class); processor.addNotRetryableExceptions(IllegalStateException.class); processor.addNotRetryableExceptions(RestClientException.class); processor.addNotRetryableExceptions(NullPointerException.class); processor.addNotRetryableExceptions(NumberFormatException.class); processor.addNotRetryableExceptions(IllegalArgumentException.class); processor.addNotRetryableExceptions(NoSuchMethodException.class); processor.addNotRetryableExceptions(JsonParseException.class); processor.addNotRetryableExceptions(MessageConversionException.class); return processor; } }
关键说明:
- 替换事务管理器类型:将原
ChainedKafkaTransactionManager替换为Spring TX的ChainedTransactionManager,Spring Boot项目默认已包含spring-tx依赖,无需额外引入。 - 保持事务顺序一致:链式事务管理器的顺序直接影响提交/回滚逻辑,需和原代码保持一致(先Kafka后JPA),确保行为与之前完全相同。
- 错误尝试的问题:给事务管理器Bean添加
@Transactional注解无效,该注解用于标记需要事务管理的业务方法,而非定义事务管理器本身,无法实现多事务联动。
内容的提问来源于stack exchange,提问作者Suchit Khadtar
相关产品推荐
相关产品推荐

