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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 18:24:45