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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 13:41:12