在运行正常的Spring Boot Kafka消费者服务中添加生产者配置引发Bean依赖异常的求助
首先,我们拆解下报错的核心原因:
你的消费者配置类KafkaConsumerConfig中的kafkaListenerContainerFactory方法,声明需要注入一个KafkaTemplate<String, CarrierStopEBO>类型的Bean,但当前Spring容器里只有你自定义的KafkaTemplate<String, CarrierModelDTO>(Bean名称kafkaTemplatetest),没有符合类型要求的Bean。同时,因为你自定义了KafkaTemplate,Spring自动配置的默认KafkaTemplate被@ConditionalOnMissingBean规则阻止创建,导致没有兜底的Bean可以匹配。
更关键的是,看你的代码发现这个注入的kafkaTemplate1参数根本没在方法里用到——这大概率是一个不必要的依赖,所以我们有两种解决思路:
解决方案1:移除不必要的依赖(最直接)
既然kafkaListenerContainerFactory方法并没有使用注入的KafkaTemplate,直接把这个参数删掉即可,完全不影响原有功能:
修改KafkaConsumerConfig中的kafkaListenerContainerFactory方法:
@Bean public ConcurrentKafkaListenerContainerFactory<String, CarrierStopEBO> kafkaListenerContainerFactory() { ConcurrentKafkaListenerContainerFactory<String, CarrierStopEBO> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.getContainerProperties().setAckOnError(false); factory.getContainerProperties().setSyncCommits(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setAckDiscarded(true); factory.getContainerProperties().setAuthorizationExceptionRetryInterval(Duration.ofMillis(30000)); return factory; }
这样Spring就不会再去寻找不存在的KafkaTemplate<String, CarrierStopEBO> Bean,消费者配置可以正常初始化,同时你的生产者配置也能正常工作。
解决方案2:如果确实需要该KafkaTemplate(比如后续业务要用到)
如果你的业务逻辑后续需要在消费者容器工厂中使用KafkaTemplate<String, CarrierStopEBO>,那么需要显式创建这个类型的Bean:
步骤1:添加CarrierStopEBO对应的ProducerFactory和KafkaTemplate
在KafkaProducerConfig中新增以下配置:
// 为CarrierStopEBO创建独立的ProducerFactory @Bean public ProducerFactory<String, CarrierStopEBO> carrierStopEboProducerFactory() { Map<String, Object> configProperties = new HashMap<>(); configProperties.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); configProperties.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class); configProperties.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, JsonSerializer.class); configProperties.put(ProducerConfig.PARTITIONER_CLASS_CONFIG, DefaultPartitioner.class); configProperties.put(CommonClientConfigs.SECURITY_PROTOCOL_CONFIG, securityProtocol); // 补充SASL相关配置 configProperties.put(SaslConfigs.SASL_MECHANISM, saslMechanism); configProperties.put(SaslConfigs.SASL_JAAS_CONFIG, jaasConfig); return new DefaultKafkaProducerFactory<>(configProperties); } // 创建对应的KafkaTemplate,指定Bean名称方便后续注入匹配 @Bean("kafkaTemplate1") public KafkaTemplate<String, CarrierStopEBO> kafkaTemplate1() { return new KafkaTemplate<>(carrierStopEboProducerFactory()); }
步骤2:在消费者容器工厂中指定注入的Bean
修改KafkaConsumerConfig中的方法,用@Qualifier明确指定要注入的Bean:
@Bean public ConcurrentKafkaListenerContainerFactory<String, CarrierStopEBO> kafkaListenerContainerFactory( @Qualifier("kafkaTemplate1") KafkaTemplate<String, CarrierStopEBO> kafkaTemplate1) { // 原有逻辑保持不变 ConcurrentKafkaListenerContainerFactory<String, CarrierStopEBO> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory()); factory.getContainerProperties().setAckOnError(false); factory.getContainerProperties().setSyncCommits(true); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); factory.setAckDiscarded(true); factory.getContainerProperties().setAuthorizationExceptionRetryInterval(Duration.ofMillis(30000)); return factory; }
这样Spring就能找到匹配类型和名称的Bean,完成注入,同时你的两个不同类型的KafkaTemplate都能正常工作。
内容的提问来源于stack exchange,提问作者Vikram Srinivasan

