Java Spring Kafka项目配置多个消费者监听器启动报错如何解决?
Spring Kafka多消费者类型配置报错解决方案
问题根因
- 两个消费者配置类中的
consumerFactory方法未指定Bean名称,默认以方法名作为Bean名,两个Bean同名导致Spring容器覆盖其中一个实例 - 两个配置类均添加了
@EnableKafka注解,重复开启Kafka注解驱动,引发配置冲突 - Spring Boot默认的Kafka自动配置会尝试加载泛型为
ConsumerFactory<Object, Object>的Bean作为默认消费者工厂,但你定义的两个ConsumerFactory都绑定了具体的自定义消息类型,无法匹配自动配置的类型要求,因此抛出找不到合格Bean的异常
修复步骤
- 给两个ConsumerFactory指定不同的Bean名称,避免实例覆盖
- 全局仅保留一个
@EnableKafka注解,建议放在项目启动类上,移除两个配置类上的该注解 - 禁用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
相关产品推荐
相关产品推荐

