You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

寻找已废弃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;
    }
}

关键说明:

  1. 替换事务管理器类型:将原ChainedKafkaTransactionManager替换为Spring TX的ChainedTransactionManager,Spring Boot项目默认已包含spring-tx依赖,无需额外引入。
  2. 保持事务顺序一致:链式事务管理器的顺序直接影响提交/回滚逻辑,需和原代码保持一致(先Kafka后JPA),确保行为与之前完全相同。
  3. 错误尝试的问题:给事务管理器Bean添加@Transactional注解无效,该注解用于标记需要事务管理的业务方法,而非定义事务管理器本身,无法实现多事务联动。

内容的提问来源于stack exchange,提问作者Suchit Khadtar

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 18:44:52