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

如何在Spring Integration XML中使用Apache Kafka的Java配置类(生产者/消费者)

Spring Integration Kafka:XML配置迁移至Java配置+XML引用方案

一、生产者配置改造

1. 编写Java配置类

创建包含生产者工厂和Kafka模板的配置类,指定Bean名称方便XML引用,同时注入配置参数:

@Configuration
public class KafkaProducerConfig {

    @Value("${broker.list}")
    private String brokersList;

    @Bean(name = "producerFactory")
    public ProducerFactory<String, String> producerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokersList);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return new DefaultKafkaProducerFactory<>(props);
    }

    @Bean(name = "kafkaTemplate")
    public KafkaTemplate<String, String> kafkaTemplate() {
        return new KafkaTemplate<>(producerFactory());
    }
}

2. 修改XML配置

删除原XML中pf和template的Bean定义,引入Java配置类,并更新出站适配器的模板引用:

<!-- 方式1:扫描配置类所在包 -->
<context:component-scan base-package="com.your.config.package" />

<!-- 方式2:直接注册配置类Bean -->
<!-- <bean class="com.your.config.package.KafkaProducerConfig" /> -->

<!-- 保留原有Integration业务链,替换kafka-template引用 -->
<int:chain input-channel="kafka-output-channel">
    <int:object-to-json-transformer/>
    <int-kafka:outbound-channel-adapter id="kafkaOutboundChannelAdapter"
                                        kafka-template="kafkaTemplate"
                                        topic="${topic.name}" />
</int:chain>

二、消费者配置改造

1. 编写Java配置类

创建包含消费者工厂和监听容器工厂的配置类:

@Configuration
public class KafkaConsumerConfig {

    @Value("${broker.list}")
    private String brokersList;

    @Value("${group.id}")
    private String groupId;

    @Bean(name = "consumerFactory")
    public ConsumerFactory<String, String> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, brokersList);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        // 可添加其他消费者配置,如auto.offset.reset等
        return new DefaultKafkaConsumerFactory<>(props);
    }

    @Bean(name = "kafkaListenerContainerFactory")
    public ConcurrentKafkaListenerContainerFactory<String, String> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        // 根据业务需求配置并发数、批量消费等
        factory.setConcurrency(3);
        return factory;
    }
}

2. 修改XML配置

删除原XML中消费者工厂和容器工厂的Bean定义,引入Java配置类,并更新消息驱动适配器的容器工厂引用:

<!-- 引入消费者配置类 -->
<context:component-scan base-package="com.your.config.package" />
<!-- 或直接注册:<bean class="com.your.config.package.KafkaConsumerConfig" /> -->

<!-- 保留原有消息驱动适配器,替换listener-container-factory引用 -->
<int-kafka:message-driven-channel-adapter id="kafkaMessageDrivenAdapter"
                                          listener-container-factory="kafkaListenerContainerFactory"
                                          topics="${topic.name}"
                                          output-channel="kafka-input-channel"/>

核心注意事项

  • Bean名称匹配:Java配置中@Bean(name = "...")指定的名称必须和XML中引用的名称完全一致,避免Spring无法定位Bean。
  • 配置注入方式:除了@Value,也可以用@ConfigurationProperties批量绑定配置项,适合复杂配置场景。
  • 扫描范围校验:确保<context:component-scan>的包路径包含你的Java配置类,否则配置类不会被加载。
  • 业务逻辑复用:Spring Integration的通道、转换器等业务组件可继续保留在XML中,仅替换Kafka底层的工厂、模板Bean即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 18:57:08