ConcurrentMessageListenerContainer操作及消费者信息获取异常
你手动通过new ConcurrentMessageListenerContainer()创建的容器实例,没有被注册到KafkaListenerEndpointRegistry中,所以调用查询方法拿不到任何结果。KafkaListenerEndpointRegistry本身只负责管理两类容器:
- 由
@KafkaListener注解标注、框架自动扫描生成的监听容器 - 主动调用其注册方法提交到内部管理集合的容器
你创建完容器后直接调用start()启动,整个过程没有和KafkaListenerEndpointRegistry做任何交互,它自然感知不到你手动创建的容器实例,强转类型也解决不了这个问题。
二选一即可:
方案1:自行维护动态创建的容器映射(推荐,无版本兼容问题)
不要依赖KafkaListenerEndpointRegistry管理你动态生成的容器,直接在你自定义的CustomKafkaContainerRegistration中维护容器实例的缓存Map,对外提供查询、获取实例的方法即可。
修正后的代码示例:
CustomKafkaContainerRegistration.java 修正
@Component public class CustomKafkaContainerRegistration { @Value("${concurrent.consumer.kafka}") private int concurrentConsumer; @Autowired public CustomConsumerFactory customConsumerFactory; // 维护所有动态创建的容器缓存 private final ConcurrentMap<String, ConcurrentMessageListenerContainer<String, String>> consumerContainerMap = new ConcurrentHashMap<>(); public void registerCustomKafkaContainer(Request request) { String consumerId = request.getConsumerId(); // 重复创建校验 if (consumerContainerMap.containsKey(consumerId)) { throw new RuntimeException(String.format("Consumer with id %s already exists", consumerId)); } CustomContainerProperties customContainerProperties = new CustomContainerProperties(request.getTopicName(), consumerId); ConcurrentMessageListenerContainer<String, String> container = new ConcurrentMessageListenerContainer<>( customConsumerFactory.getCustomConsumerFactory(), customContainerProperties.getContainerProperties()); container.setConcurrency(concurrentConsumer); container.setAutoStartup(request.getConsumerActivation()); container.setBeanName(consumerId); // 先存入缓存再启动 consumerContainerMap.put(consumerId, container); container.start(); } // 根据consumerId获取容器实例 public ConcurrentMessageListenerContainer<String, String> getContainerById(String consumerId) { return consumerContainerMap.get(consumerId); } // 获取所有已注册的consumerId public Set<String> getAllConsumerIds() { return consumerContainerMap.keySet(); } // 移除已销毁的容器,避免内存泄漏 public void removeContainer(String consumerId) { consumerContainerMap.remove(consumerId); } }
KafkaConsumerRegistryController.java 修正
把原来注入KafkaListenerEndpointRegistry的地方换成注入CustomKafkaContainerRegistration,所有获取容器的操作都从自定义注册类拿即可,示例:
@Slf4j @RestController @RequestMapping(path = "/api/kafka/registry") public class KafkaConsumerRegistryController { @Autowired private CustomKafkaContainerRegistration customKafkaContainerRegistration; @GetMapping public List<KafkaConsumerResponse> getConsumerIds() { return customKafkaContainerRegistration.getAllConsumerIds() .stream() .map(this::createKafkaConsumerResponse) .collect(Collectors.toList()); } @PostMapping(path = "/create") @ResponseStatus(HttpStatus.CREATED) public void createConsumer(@RequestBody Request request) { customKafkaContainerRegistration.registerCustomKafkaContainer(request); } @PostMapping(path = "/activate") @ResponseStatus(HttpStatus.ACCEPTED) public void activateConsumer(@RequestParam String consumerId) { ConcurrentMessageListenerContainer listenerContainer = customKafkaContainerRegistration.getContainerById(consumerId); if (Objects.isNull(listenerContainer)) { throw new RuntimeException(String.format("Consumer with id %s is not found", consumerId)); } else if (listenerContainer.isRunning()) { throw new RuntimeException(String.format("Consumer with id %s is already running", consumerId)); } else { log.info("Running a consumer with id " + consumerId); listenerContainer.start(); } } @PostMapping(path = "/pause") @ResponseStatus(HttpStatus.ACCEPTED) public void pauseConsumer(@RequestParam String consumerId) { MessageListenerContainer listenerContainer = customKafkaContainerRegistration.getContainerById(consumerId); if (Objects.isNull(listenerContainer)) { throw new RuntimeException(String.format("Consumer with id %s is not found", consumerId)); } else if (!listenerContainer.isRunning()) { throw new RuntimeException(String.format("Consumer with id %s is not running", consumerId)); } else if (listenerContainer.isContainerPaused()) { throw new RuntimeException(String.format("Consumer with id %s is already paused", consumerId)); } else if (listenerContainer.isPauseRequested()) { throw new RuntimeException(String.format("Consumer with id %s is already requested to be paused", consumerId)); } else { log.info("Pausing a consumer with id " + consumerId); listenerContainer.pause(); } } @PostMapping(path = "/resume") @ResponseStatus(HttpStatus.ACCEPTED) public void resumeConsumer(@RequestParam String consumerId) { MessageListenerContainer listenerContainer = customKafkaContainerRegistration.getContainerById(consumerId); if (Objects.isNull(listenerContainer)) { throw new RuntimeException(String.format("Consumer with id %s is not found", consumerId)); } else if (!listenerContainer.isRunning()) { throw new RuntimeException(String.format("Consumer with id %s is not running", consumerId)); } else if (!listenerContainer.isContainerPaused()) { throw new RuntimeException(String.format("Consumer with id %s is not paused", consumerId)); } else { log.info("Resuming a consumer with id " + consumerId); listenerContainer.resume(); } } @PostMapping(path = "/deactivate") @ResponseStatus(HttpStatus.ACCEPTED) public void deactivateConsumer(@RequestParam String consumerId) { ConcurrentMessageListenerContainer listenerContainer = customKafkaContainerRegistration.getContainerById(consumerId); if (Objects.isNull(listenerContainer)) { throw new RuntimeException(String.format("Consumer with id %s is not found", consumerId)); } else if (!listenerContainer.isRunning()) { throw new RuntimeException(String.format("Consumer with id %s is already stop", consumerId)); } else { log.info("Stopping a consumer with id " + consumerId); listenerContainer.stop(); // 容器停止后从缓存移除 customKafkaContainerRegistration.removeContainer(consumerId); } } private KafkaConsumerResponse createKafkaConsumerResponse(String consumerId) { MessageListenerContainer listenerContainer = customKafkaContainerRegistration.getContainerById(consumerId); return KafkaConsumerResponse.builder() .consumerId(consumerId) .groupId(listenerContainer.getGroupId()) .listenerId(listenerContainer.getListenerId()) .active(listenerContainer.isRunning()) .assignments(Optional.ofNullable(listenerContainer.getAssignedPartitions()) .map(topicPartitions -> topicPartitions.stream() .map(this::createKafkaConsumerAssignmentResponse) .collect(Collectors.toList())) .orElse(null)) .build(); } private KafkaConsumerAssignmentResponse createKafkaConsumerAssignmentResponse( TopicPartition topicPartition) { return KafkaConsumerAssignmentResponse.builder() .topic(topicPartition.topic()) .partition(topicPartition.partition()) .build(); } }
方案2:将手动创建的容器注册到KafkaListenerEndpointRegistry
如果你坚持要使用KafkaListenerEndpointRegistry管理,需要在创建完容器后主动调用注册方法,注意不同Spring Kafka版本API有差异。
注意:部分低版本Spring Kafka没有提供直接注册已创建容器实例的公开方法,需要反射操作内部私有Map,不推荐生产环境使用。
注册逻辑核心是构造对应KafkaListenerEndpoint实例,绑定你自定义的容器工厂、监听器、topic参数后,调用registry.registerListenerContainer(endpoint, factory, isLazyStart)完成注册,注册后即可通过consumerId从registry中查询到容器实例。
- 你代码中提到的确认模式
MANUAL_IMMDEDIATE存在拼写错误,正确常量为MANUAL_IMMEDIATE,请修正避免配置不生效。 - 动态创建容器时建议增加重复consumerId校验、容器销毁后的缓存清理逻辑,避免内存泄漏、重复消费问题。
内容的提问来源于stack exchange,提问作者None for Nothing

