Spring Boot Kafka闲置消费者容器监听器查询问题求助
动态Kafka消费者容器复用问题
我通过以下代码实现了传入topic名称动态创建对应的Kafka监听器容器:
@Component public class MessageListenerConfigurer { @Autowired ConcurrentKafkaListenerContainerFactory<String, String> factory; public ConcurrentMessageListenerContainer<String,String> createContainerForTopic(String topicName) { ConcurrentMessageListenerContainer<String, String> container = factory.createContainer(topicName); container.getContainerProperties().setMessageListener(new InstantMessageListener()); return container; } }
监听器实现类:
package org.kafka.listener; import lombok.AllArgsConstructor; import lombok.SneakyThrows; import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.kafka.listener.MessageListener; import javax.sql.DataSource; import javax.transaction.Transactional; import java.util.concurrent.CompletableFuture; @Slf4j @AllArgsConstructor public class InstantMessageListener implements MessageListener<String, String> { public InstantMessageListener(){ } @Autowired DataSource dataSource; @Transactional public void onMessage(ConsumerRecord<String,String> record) { log.info("My message listener got a new record: " + record); log.info("message is: "+record.toString()); log.warn("onMessage:: ===== from eventhub topic: {}, partition: {}, offset: {}, message: {}, timestamp: {}", record.topic(), record.partition(), record.offset(), record.value(), record.timestamp()); CompletableFuture.runAsync(this::sleep).join(); log.info("My message listener done processing record: " + record); } @SneakyThrows private void sleep() { Thread.sleep(5000); } }
我希望查询已创建的容器,判断是否存在闲置容器以复用,避免重复创建同一topic的容器,但以下查询代码始终返回空结果:
package org.kafka.controller; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.config.AutowireCapableBeanFactory; import org.springframework.beans.factory.config.SingletonBeanRegistry; import org.springframework.context.ApplicationContext; import org.springframework.http.HttpStatus; import org.springframework.kafka.listener.ConcurrentMessageListenerContainer; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RequestMapping; import org.springframework.web.bind.annotation.ResponseStatus; import org.springframework.web.bind.annotation.RestController; @RestController class ExportController { @Autowired private ApplicationContext applicationContext; @GetMapping("/beans") @ResponseStatus(value = HttpStatus.OK) String[] registeredBeans() { return printBeans(); } private String[] printBeans() { AutowireCapableBeanFactory autowireCapableBeanFactory = applicationContext.getAutowireCapableBeanFactory(); if (autowireCapableBeanFactory instanceof ConcurrentMessageListenerContainer) { String[] singletonNames = ((SingletonBeanRegistry) autowireCapableBeanFactory).getSingletonNames(); for (String singleton : singletonNames) { System.out.println(singleton); } return singletonNames; } return null; } }
问题原因及解决方案
1. 核心问题分析
- 动态创建的容器未注册到Spring上下文:
factory.createContainer()生成的容器实例只是普通对象,没有被Spring管理,所以无法通过上下文的单例列表查询到。 - 类型判断逻辑错误:
AutowireCapableBeanFactory是Spring的Bean工厂接口,不可能是ConcurrentMessageListenerContainer类型,这是查询返回null的直接原因。
2. 修正容器创建逻辑:缓存+上下文注册
修改MessageListenerConfigurer,新增容器缓存并将容器注册到Spring上下文,同时实现复用逻辑:
@Component public class MessageListenerConfigurer { @Autowired private ConcurrentKafkaListenerContainerFactory<String, String> factory; @Autowired private ApplicationContext applicationContext; // 用ConcurrentHashMap缓存已创建的容器,key为topic名称 private final Map<String, ConcurrentMessageListenerContainer<String, String>> containerCache = new ConcurrentHashMap<>(); public ConcurrentMessageListenerContainer<String, String> createContainerForTopic(String topicName) { // 先检查缓存中是否存在该topic的容器 if (containerCache.containsKey(topicName)) { ConcurrentMessageListenerContainer<String, String> container = containerCache.get(topicName); // 若容器已停止则重新启动 if (!container.isRunning()) { container.start(); } return container; } // 创建新容器 ConcurrentMessageListenerContainer<String, String> container = factory.createContainer(topicName); // 修复监听器注入问题:从Spring上下文获取实例,而非new InstantMessageListener listener = applicationContext.getBean(InstantMessageListener.class); container.getContainerProperties().setMessageListener(listener); // 将容器注册到Spring上下文,让Spring管理其生命周期 AutowireCapableBeanFactory beanFactory = applicationContext.getAutowireCapableBeanFactory(); beanFactory.autowireBean(container); beanFactory.initializeBean(container, "kafkaContainer_" + topicName); // 注册为单例Bean(可选,便于Spring统一管理) if (beanFactory instanceof SingletonBeanRegistry) { ((SingletonBeanRegistry) beanFactory).registerSingleton("kafkaContainer_" + topicName, container); } // 启动容器并加入缓存 container.start(); containerCache.put(topicName, container); return container; } // 对外提供查询容器状态的方法 public Map<String, Boolean> getContainerStatuses() { Map<String, Boolean> statusMap = new HashMap<>(); containerCache.forEach((topic, container) -> statusMap.put(topic, container.isRunning())); return statusMap; } // 提供容器停止与清理方法 public void stopAndRemoveContainer(String topicName) { ConcurrentMessageListenerContainer<String, String> container = containerCache.remove(topicName); if (container != null && container.isRunning()) { container.stop(); } } }
3. 修正监听器注入问题
将InstantMessageListener标记为@Component,让Spring管理其Bean实例,确保DataSource能被正确注入:
@Slf4j @AllArgsConstructor @Component // 添加该注解 public class InstantMessageListener implements MessageListener<String, String> { @Autowired DataSource dataSource; // 其余代码不变 }
4. 修正容器查询接口
直接通过MessageListenerConfigurer的缓存查询容器状态,无需操作Spring上下文的单例列表:
@RestController class ExportController { @Autowired private MessageListenerConfigurer configurer; @GetMapping("/containers") @ResponseStatus(value = HttpStatus.OK) Map<String, Boolean> getRegisteredContainers() { return configurer.getContainerStatuses(); } }
内容的提问来源于stack exchange,提问作者Utkarsh Saraf
相关产品推荐
相关产品推荐

