kafka动态创建消费者(实时更新topic和servers)_@kafkalistener(topics 动态刷新-程序员宅基地

技术标签: java  kafka  

一、疑问描述
spring-kafka通过 @KafkaListener 的方式配置订阅的topic,通过@Configuration 配置创建kafkaListenerContainerFactory。
如下:

@Configuration
@EnableKafka
public class KafkaConfig {
    

    private static final String KAFKA_SERVERS_CONFIG = "10.192.77.202:9092";
    private static final String LOCAL_GROUP_ID = "test";

    @Bean
    ConcurrentKafkaListenerContainerFactory<Integer, String>
    kafkaListenerContainerFactory() {
    
        ConcurrentKafkaListenerContainerFactory<Integer, String> factory =
                new ConcurrentKafkaListenerContainerFactory<>();
        factory.setConsumerFactory(consumerFactory());
        return factory;
    }

    @Bean
    public ConsumerFactory<Integer, String> consumerFactory() {
    
        return new DefaultKafkaConsumerFactory<>(consumerConfigs());
    }

    @Bean
    public Map<String, Object> consumerConfigs() {
    
        Map<String, Object> props = new HashMap<>();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_SERVERS_CONFIG);
        props.put(ConsumerConfig.GROUP_ID_CONFIG, LOCAL_GROUP_ID);
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, IntegerDeserializer.class);
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
        return props;
    }

    @Bean
    public Map<String, Object> producerConfigs() {
    
        Map<String, Object> props = new HashMap<>();
        props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, KAFKA_SERVERS_CONFIG);
        props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, IntegerSerializer.class);
        props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class);
        return props;
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplate() {
    
        return new KafkaTemplate<String, String>(producerFactory());
    }

    @Bean
    public ProducerFactory<String, String> producerFactory() {
    
        return new DefaultKafkaProducerFactory<>(producerConfigs());
    }

    @KafkaListener(topics = "TEST_TOPIC_NEW")
    public void listen(String data) {
    
        System.out.println("kafkaconfig =listen======="+data);
    }
}

但想要动态的创建监听者对象,如通过数据库的方式配置KAFKA_SERVERS_CONFIG 和LOCAL_GROUP_ID ,并且可以不用重启服务,实现热更新。通过spring-kafka提供的接口没有找到好的解决方法。

二、解决方案
所以,考虑通过最基本的手动创建消费者对象。
通过定时任务,每三分钟check一次,从数据库读取相应配置,将已有配置写入缓存,当读取的配置和缓存不一致时,销毁已有消费者,创建新的消费者。
如果有好的方案,谢谢告知~

/**
 * 每三分钟check一次kafka配置
 * @throws Exception
 */
@Scheduled(cron = "1 1/3 * * * ? ")
public void deviceNotifyConfig(){
    
    Map<String, String> kafkaConfigs = systemConfigService.fetchConfigLikeKey("kafka");
    if(kafkaConfigs != null && kafkaConfigs.size() != 0)
    {
    
        String kafkaIp = kafkaConfigs.get("kafkaIp");
        String kafkaPort = kafkaConfigs.get("kafkaPort");
        String kafkaUserName = kafkaConfigs.get("kafkaUserName");
        String kafkaPassword = kafkaConfigs.get("kafkaPassword");
        if(StringUtils.isNotEmpty(KafkaLinkCache.kafkaConfigCache))
        {
    
            if (!KafkaLinkCache.kafkaConfigCache.equals(kafkaIp + "_" + kafkaPort))
            {
    
                //关闭已有消费者对象
                KafkaConsumer<String, String> consumer = KafkaLinkCache.DEVICE_CONSUMER_MAP.get("kafkaComsumer");
                if(consumer != null)
                {
    
                    resourceNotifyConsumer.closeConsumer();
                }
                KafkaLinkCache.DEVICE_CONSUMER_MAP.clear();
                this.handlerConsumer(kafkaIp, kafkaPort);
            }
        }
        else
        {
    
            this.handlerConsumer(kafkaIp, kafkaPort);
        }
    }else
    {
    
        //关闭已有消费者对象
        KafkaConsumer<String, String> consumer = KafkaLinkCache.DEVICE_CONSUMER_MAP.get("kafkaComsumer");
        if(consumer != null)
        {
    
            resourceNotifyConsumer.closeConsumer();
        }
        KafkaLinkCache.DEVICE_CONSUMER_MAP.clear();
    }
}

private void handlerConsumer(String kafkaIp, String kafkaPort) {
    
    Properties props = new Properties();
    props.setProperty("bootstrap.servers", kafkaIp + ":" + kafkaPort);
    // key反序列化
    props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
    // value反序列化
    props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer");
    // 每个消费者都必须属于某一个消费组,所以必须指定group.id
    props.put("group.id", "test");

    // 构造消费者对象
    deviceNoifyThreadExecutor.execute(()->{
    
        KafkaConsumer<String, String> consumerObj = null;
        // 指定多主题:
        List<String> topics = CbdmOptUtil.stringToStringList(PropertiesUtil.getProperty("kafka.subscribe.topics"), ConstParamErrorCode.DEFAULT_SPLIT_KEY, false);
        try {
    
            consumerObj = new KafkaConsumer<>(props);
            if(consumerObj != null) {
    
                consumerObj.subscribe(topics);
                resourceNotifyConsumer.setConsumer(consumerObj);
                KafkaLinkCache.DEVICE_CONSUMER_MAP.put("kafkaComsumer", consumerObj);
                resourceNotifyConsumer.onMessage();
            }
        } catch(Exception e) {
    
            LogUtils.logError(RunTimeLogUtil.toErrorLog(ConstParamErrorCode.SYSTEM_CODE_FAIL + "", LogObjectTypeEnum.SYSTEM,"consume",
                    "resolve data platform notify error"),e);
        }finally {
    
            // 关闭
            consumerObj.close();
        }
    });

    //保存配置
    KafkaLinkCache.kafkaConfigCache = kafkaIp + "_" + kafkaPort;
}

@Component(value = "resourceNotifyConsumer")
public class ResourceNotifyConsumer {
    

    private Logger logger = LoggerFactory.getLogger(ResourceNotifyConsumer.class);

    @Resource
    IAccessDeviceService resourceService;

    private KafkaConsumer<String, String> consumer = null;

    public KafkaConsumer<String, String> getConsumer() {
    
        return consumer;
    }

    public void setConsumer(KafkaConsumer<String, String> consumer) {
    
        this.consumer = consumer;
    }

    public void closeConsumer()
    {
    
        //consumer非线程安全,依靠gc回收
        consumer = null;
    }

    public void onMessage(){
    
        try{
    
            logger.info(RunTimeLogUtil.toLog(LogObjectTypeEnum.SYSTEM,"consume","Get resource Notify start",null,null));

            while (true) {
    
                if(consumer != null)
                {
    
                    // timeout 阻塞时间,从kafka中取出100毫秒的数据,有可能一次取出0到N条
                    List<Map<String,Object>> datas = new ArrayList<>();
                    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                    // 遍历
                    for (ConsumerRecord<String, String> record : records) {
    
                        Map<String,Object> notifyDto = ( Map<String,Object> ) JsonUtils.jsonToMap(record.value());
                        datas.add(notifyDto);
                    }
                    // 拿出结果
                    if(CollectionUtils.isNotEmpty(datas)){
    
                        logger.info(RunTimeLogUtil.toLog(LogObjectTypeEnum.SYSTEM,"consume","Get resource Notify",null,null, "record"),JsonUtils.object2Json(datas));
                        // 起线程处理 资源变更通知
                        resourceHandle(datas);
                    }
                } else {
    
                    break;
                }
            }
        }catch (Throwable e){
    
            logger.error(RunTimeLogUtil.toErrorLog(ConstParamErrorCode.SYSTEM_CODE_FAIL + "",LogObjectTypeEnum.SYSTEM,"consume",
                   "resolve resource notify error"),e);
        }
    }

    /**
     *
     * @param datas
     */
    private void resourceHandle(List<Map<String,Object>> datas){
    
        if(CollectionUtils.isNotEmpty(datas)){
    
            try {
    
                new Thread(() -> resourceService.dealResource(datas)).start();
            }catch (Throwable e){
    
                logger.error(RunTimeLogUtil.toErrorLog(ConstParamErrorCode.SYSTEM_CODE_FAIL + "",LogObjectTypeEnum.SYSTEM,"consume",
                    "resourceHandle error"),e);
            }
        }else{
    
            logger.info(RunTimeLogUtil.toLog(LogObjectTypeEnum.SYSTEM,"consume","resource notify data is empty!",null,null));
        }
    }
}
版权声明:本文为博主原创文章,遵循 CC 4.0 BY-SA 版权协议,转载请附上原文出处链接和本声明。
本文链接:https://blog.csdn.net/weixin_41422086/article/details/104849127

智能推荐

JWT(Json Web Token)实现无状态登录_无状态token登录-程序员宅基地

文章浏览阅读685次。1.1.什么是有状态?有状态服务,即服务端需要记录每次会话的客户端信息,从而识别客户端身份,根据用户身份进行请求的处理,典型的设计如tomcat中的session。例如登录:用户登录后,我们把登录者的信息保存在服务端session中,并且给用户一个cookie值,记录对应的session。然后下次请求,用户携带cookie值来,我们就能识别到对应session,从而找到用户的信息。缺点是什么?服务端保存大量数据,增加服务端压力 服务端保存用户状态,无法进行水平扩展 客户端请求依赖服务.._无状态token登录

SDUT OJ逆置正整数-程序员宅基地

文章浏览阅读293次。SDUT OnlineJudge#include<iostream>using namespace std;int main(){int a,b,c,d;cin>>a;b=a%10;c=a/10%10;d=a/100%10;int key[3];key[0]=b;key[1]=c;key[2]=d;for(int i = 0;i<3;i++){ if(key[i]!=0) { cout<<key[i.

年终奖盲区_年终奖盲区表-程序员宅基地

文章浏览阅读2.2k次。年终奖采用的平均每月的收入来评定缴税级数的,速算扣除数也按照月份计算出来,但是最终减去的也是一个月的速算扣除数。为什么这么做呢,这样的收的税更多啊,年终也是一个月的收入,凭什么减去12*速算扣除数了?这个霸道(不要脸)的说法,我们只能合理避免的这些跨级的区域了,那具体是那些区域呢?可以参考下面的表格:年终奖一列标红的一对便是盲区的上下线,发放年终奖的数额一定一定要避免这个区域,不然公司多花了钱..._年终奖盲区表

matlab 提取struct结构体中某个字段所有变量的值_matlab读取struct类型数据中的值-程序员宅基地

文章浏览阅读7.5k次,点赞5次,收藏19次。matlab结构体struct字段变量值提取_matlab读取struct类型数据中的值

Android fragment的用法_android reader fragment-程序员宅基地

文章浏览阅读4.8k次。1,什么情况下使用fragment通常用来作为一个activity的用户界面的一部分例如, 一个新闻应用可以在屏幕左侧使用一个fragment来展示一个文章的列表,然后在屏幕右侧使用另一个fragment来展示一篇文章 – 2个fragment并排显示在相同的一个activity中,并且每一个fragment拥有它自己的一套生命周期回调方法,并且处理它们自己的用户输_android reader fragment

FFT of waveIn audio signals-程序员宅基地

文章浏览阅读2.8k次。FFT of waveIn audio signalsBy Aqiruse An article on using the Fast Fourier Transform on audio signals. IntroductionThe Fast Fourier Transform (FFT) allows users to view the spectrum content of _fft of wavein audio signals

随便推点

Awesome Mac:收集的非常全面好用的Mac应用程序、软件以及工具_awesomemac-程序员宅基地

文章浏览阅读5.9k次。https://jaywcjlove.github.io/awesome-mac/ 这个仓库主要是收集非常好用的Mac应用程序、软件以及工具,主要面向开发者和设计师。有这个想法是因为我最近发了一篇较为火爆的涨粉儿微信公众号文章《工具武装的前端开发工程师》,于是建了这么一个仓库,持续更新作为补充,搜集更多好用的软件工具。请Star、Pull Request或者使劲搓它 issu_awesomemac

java前端技术---jquery基础详解_简介java中jquery技术-程序员宅基地

文章浏览阅读616次。一.jquery简介 jQuery是一个快速的,简洁的javaScript库,使用户能更方便地处理HTML documents、events、实现动画效果,并且方便地为网站提供AJAX交互 jQuery 的功能概括1、html 的元素选取2、html的元素操作3、html dom遍历和修改4、js特效和动画效果5、css操作6、html事件操作7、ajax_简介java中jquery技术

Ant Design Table换滚动条的样式_ant design ::-webkit-scrollbar-corner-程序员宅基地

文章浏览阅读1.6w次,点赞5次,收藏19次。我修改的是表格的固定列滚动而产生的滚动条引用Table的组件的css文件中加入下面的样式:.ant-table-body{ &amp;amp;::-webkit-scrollbar { height: 5px; } &amp;amp;::-webkit-scrollbar-thumb { border-radius: 5px; -webkit-box..._ant design ::-webkit-scrollbar-corner

javaWeb毕设分享 健身俱乐部会员管理系统【源码+论文】-程序员宅基地

文章浏览阅读269次。基于JSP的健身俱乐部会员管理系统项目分享:见文末!

论文开题报告怎么写?_开题报告研究难点-程序员宅基地

文章浏览阅读1.8k次,点赞2次,收藏15次。同学们,是不是又到了一年一度写开题报告的时候呀?是不是还在为不知道论文的开题报告怎么写而苦恼?Take it easy!我带着倾尽我所有开题报告写作经验总结出来的最强保姆级开题报告解说来啦,一定让你脱胎换骨,顺利拿下开题报告这个高塔,你确定还不赶快点赞收藏学起来吗?_开题报告研究难点

原生JS 与 VUE获取父级、子级、兄弟节点的方法 及一些DOM对象的获取_获取子节点的路径 vue-程序员宅基地

文章浏览阅读6k次,点赞4次,收藏17次。原生先获取对象var a = document.getElementById("dom");vue先添加ref <div class="" ref="divBox">获取对象let a = this.$refs.divBox获取父、子、兄弟节点方法var b = a.childNodes; 获取a的全部子节点 var c = a.parentNode; 获取a的父节点var d = a.nextSbiling; 获取a的下一个兄弟节点 var e = a.previ_获取子节点的路径 vue