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

Java Spring Kafka项目配置多个消费者监听器启动报错如何解决?

Spring Kafka多消费者类型配置报错解决方案

问题根因

  • 两个消费者配置类中的consumerFactory方法未指定Bean名称,默认以方法名作为Bean名,两个Bean同名导致Spring容器覆盖其中一个实例
  • 两个配置类均添加了@EnableKafka注解,重复开启Kafka注解驱动,引发配置冲突
  • Spring Boot默认的Kafka自动配置会尝试加载泛型为ConsumerFactory<Object, Object>的Bean作为默认消费者工厂,但你定义的两个ConsumerFactory都绑定了具体的自定义消息类型,无法匹配自动配置的类型要求,因此抛出找不到合格Bean的异常

修复步骤

  1. 给两个ConsumerFactory指定不同的Bean名称,避免实例覆盖
  2. 全局仅保留一个@EnableKafka注解,建议放在项目启动类上,移除两个配置类上的该注解
  3. 禁用Spring Boot Kafka默认自动配置,因为你已经完全自定义了消费者工厂和监听器容器工厂,不需要自动装配默认Bean

修改后代码示例

KafkaConsumerConfigPipelineMessage 调整

@Configuration
// 移除类上的@EnableKafka注解
public class KafkaConsumerConfigPipelineMessage {
    @Value(value = "${spring.kafka.bootstrap-servers}")
    private String bootstrapAddress;

    @Value(value = "${spring.kafka.producer.group-id}")
    private String groupId;

    @Value(value = "${saslMechanism}")
    private String saslMechanism;

    @Value(value = "${kafkaUser}")
    private String kafkaUser;

    @Value(value = "${kafkaPassword}")
    private String kafkaPassword;

    @Value(value = "${securityProtocol}")
    private String securityProtocol;

    // 指定Bean名称避免冲突
    @Bean("pipelineConsumerFactory")
    public ConsumerFactory<String, PipelineMessage> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, AvroDeserializer.class);
        props.put("sasl.mechanism",saslMechanism);
        props.put("security.protocol",securityProtocol);
        props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username ='"+kafkaUser+"' password = '"+kafkaPassword+"';");
        return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new AvroDeserializer<>(PipelineMessage.class));
    }

    @Bean("kafkaListenerPipelineMessage")
    public ConcurrentKafkaListenerContainerFactory<String, PipelineMessage> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, PipelineMessage> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

KafkaConsumerConfigBridge 调整

@Configuration
// 移除类上的@EnableKafka注解
public class KafkaConsumerConfigBridge {
    @Value(value = "${spring.kafka.bootstrap-servers}")
    private String bootstrapAddress;

    @Value(value = "${spring.kafka.producer.group-id}")
    private String groupId;

    @Value(value = "${saslMechanism}")
    private String saslMechanism;

    @Value(value = "${kafkaUser}")
    private String kafkaUser;

    @Value(value = "${kafkaPassword}")
    private String kafkaPassword;

    @Value(value = "${securityProtocol}")
    private String securityProtocol;

    // 指定Bean名称避免冲突
    @Bean("bridgeConsumerFactory")
    public ConsumerFactory<String, Bridge> consumerFactory() {
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapAddress);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, AvroDeserializer.class);
        props.put("sasl.mechanism",saslMechanism);
        props.put("security.protocol",securityProtocol);
        props.put("sasl.jaas.config", "org.apache.kafka.common.security.plain.PlainLoginModule required username ='"+kafkaUser+"' password = '"+kafkaPassword+"';");
        return new DefaultKafkaConsumerFactory<>(props, new StringDeserializer(), new AvroDeserializer<>(Bridge.class));
    }

    @Bean("kafkaListenerBridge")
    public ConcurrentKafkaListenerContainerFactory<String, Bridge> kafkaListenerContainerFactory() {
        ConcurrentKafkaListenerContainerFactory<String, Bridge> factory = new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }
}

启动类调整

在启动类的@SpringBootApplication注解中排除Kafka自动配置类,同时添加全局唯一的@EnableKafka注解:

@EnableKafka
@SpringBootApplication(exclude = {
    org.springframework.boot.autoconfigure.kafka.KafkaAutoConfiguration.class,
    org.springframework.boot.autoconfigure.kafka.KafkaAnnotationDrivenConfiguration.class
})
public class 你的启动类类名 {
    public static void main(String[] args) {
        SpringApplication.run(你的启动类类名.class, args);
    }
}

验证说明

原有@KafkaListener注解已经正确指定了对应的containerFactory属性,无需调整,重启应用即可正常注册两个监听器,分别消费对应Topic的不同类型消息。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 07:15:04