最新国产好看的视频,伊人天堂AV在线,国产Aaaaaa视频,蜜臀视频在线观看一区,人妻av色图,密臀久久久精品影片,青青视频免费观看毛片,久草在线观看视,国产三级精品色情在线

springboot+kafka中@KafkaListener動(dòng)態(tài)指定多個(gè)topic問(wèn)題

 更新時(shí)間:2022年12月27日 11:35:56   作者:Forward233  
這篇文章主要介紹了springboot+kafka中@KafkaListener動(dòng)態(tài)指定多個(gè)topic問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教

說(shuō)明

本項(xiàng)目為springboot+kafak的整合項(xiàng)目,故其用了springboot中對(duì)kafak的消費(fèi)注解@KafkaListener

首先,application.properties中配置用逗號(hào)隔開的多個(gè)topic。

方法:利用Spring的SpEl表達(dá)式,將topics 配置為:@KafkaListener(topics = “#{’${topics}’.split(’,’)}”)

運(yùn)行程序,console打印的效果如下:

因?yàn)橹婚_了一條消費(fèi)者線程,所以所有的topic和分區(qū)都分配給這條線程。

如果你想開多條消費(fèi)者線程去消費(fèi)這些topic,添加@KafkaListener注解的參數(shù)concurrency的值為自己想要的消費(fèi)者個(gè)數(shù)即可(注意,消費(fèi)者數(shù)要小于等于你開的所有topic的分區(qū)數(shù)總和)

運(yùn)行程序,console打印的效果如下:

總結(jié)一下大家問(wèn)的最多的一個(gè)問(wèn)題

如何在程序運(yùn)行的過(guò)程中,改變topic,消費(fèi)者能夠消費(fèi)修改后的topic?

ans: 經(jīng)過(guò)嘗試,使用@KafkaListener注解實(shí)現(xiàn)不了此需求,在程序啟動(dòng)的時(shí)候,程序就會(huì)根據(jù)@KafkaListener的注解信息初始化好消費(fèi)者去消費(fèi)指定好的topic。如果在程序運(yùn)行的過(guò)程中,修改topic,不會(huì)讓此消費(fèi)者修改消費(fèi)者的配置再重新訂閱topic的。

不過(guò)我們可以有個(gè)折中的辦法,就是利用@KafkaListener的topicPattern參數(shù)來(lái)進(jìn)行topic匹配。

具體如何操作的可以看下這篇文章:

http://www.fzitv.net/article/271098.htm

終極方法

思路

不使用@KafkaListener,使用kafka原生客戶端依賴,手動(dòng)初始化消費(fèi)者,開啟消費(fèi)者線程。

在消費(fèi)者線程中,每次循環(huán)都從配置、數(shù)據(jù)庫(kù)或者其他配置源獲取最新的topic信息,與之前的topic比較,如果發(fā)生變化,重新訂閱topic或者初始化消費(fèi)者。

實(shí)現(xiàn)

加入kafka客戶端依賴(本次測(cè)試服務(wù)端kafka版本:2.12-2.4.0)

<dependency>
	<groupId>org.apache.kafka</groupId>
	<artifactId>kafka-clients</artifactId>
	<version>2.3.0</version>
</dependency>

代碼

@Service
@Slf4j
public class KafkaConsumers implements InitializingBean {

    /**
     * 消費(fèi)者
     */
    private static KafkaConsumer<String, String> consumer;
    /**
     * topic
     */
    private List<String> topicList;

    public static String getNewTopic() {
        try {
            return org.apache.commons.io.FileUtils.readLines(new File("D:/topic.txt"), "utf-8").get(0);
        } catch (IOException e) {
            e.printStackTrace();
        }
        return null;
    }

    /**
     * 初始化消費(fèi)者(配置寫死是為了快速測(cè)試,請(qǐng)大家使用配置文件)
     *
     * @param topicList
     * @return
     */
    public KafkaConsumer<String, String> getInitConsumer(List<String> topicList) {
        //配置信息
        Properties props = new Properties();
        //kafka服務(wù)器地址
        props.put("bootstrap.servers", "192.168.9.185:9092");
        //必須指定消費(fèi)者組
        props.put("group.id", "haha");
        //設(shè)置數(shù)據(jù)key和value的序列化處理類
        props.put("key.deserializer", StringDeserializer.class);
        props.put("value.deserializer", StringDeserializer.class);
        //創(chuàng)建消息者實(shí)例
        KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
        //訂閱topic的消息
        consumer.subscribe(topicList);
        return consumer;
    }

    /**
     * 開啟消費(fèi)者線程
     * 異常請(qǐng)自己根據(jù)需求自己處理
     */
    @Override
    public void afterPropertiesSet() {
        // 初始化topic
        topicList = Splitter.on(",").splitToList(Objects.requireNonNull(getNewTopic()));
        if (org.apache.commons.collections.CollectionUtils.isNotEmpty(topicList)) {
            consumer = getInitConsumer(topicList);
            // 開啟一個(gè)消費(fèi)者線程
            new Thread(() -> {
                while (true) {
                    // 模擬從配置源中獲取最新的topic(字符串,逗號(hào)隔開)
                    final List<String> newTopic = Splitter.on(",").splitToList(Objects.requireNonNull(getNewTopic()));
                    // 如果topic發(fā)生變化
                    if (!topicList.equals(newTopic)) {
                        log.info("topic 發(fā)生變化:newTopic:{},oldTopic:{}-------------------------", newTopic, topicList);
                        // method one:重新訂閱topic:
                        topicList = newTopic;
                        consumer.subscribe(newTopic);
                        // method two:關(guān)閉原來(lái)的消費(fèi)者,重新初始化一個(gè)消費(fèi)者
                        //consumer.close();
                        //topicList = newTopic;
                        //consumer = getInitConsumer(newTopic);
                        continue;
                    }
                    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
                    for (ConsumerRecord<String, String> record : records) {
                        System.out.println("key:" + record.key() + "" + ",value:" + record.value());
                    }
                }
            }).start();
        }
    }
}

說(shuō)一下第72行代碼:

ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));

上面這行代碼表示:在100ms內(nèi)等待Kafka的broker返回?cái)?shù)據(jù).超市參數(shù)指定poll在多久之后可以返回,不管有沒(méi)有可用的數(shù)據(jù)都要返回。

在修改topic后,必須等到此次poll拉取的消息處理完,while(true)循環(huán)的時(shí)候檢測(cè)topic發(fā)生變化,才能重新訂閱topic.

poll()方法一次拉取得消息數(shù)默認(rèn)為:500,如下圖,kafka客戶端源碼中設(shè)置的。

如果想自定義此配置,可在初始化消費(fèi)者時(shí)加入

運(yùn)行結(jié)果(測(cè)試的topic中都無(wú)數(shù)據(jù))

注意:KafkaConsumer是線程不安全的,不要用一個(gè)KafkaConsumer實(shí)例開啟多個(gè)消費(fèi)者,要開啟多個(gè)消費(fèi)者,需要new 多個(gè)KafkaConsumer實(shí)例。

總結(jié)

以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。

相關(guān)文章

  • 通過(guò)實(shí)例了解cookie機(jī)制特性及使用方法

    通過(guò)實(shí)例了解cookie機(jī)制特性及使用方法

    這篇文章主要介紹了通過(guò)實(shí)例了解cookie機(jī)制特性及使用方法,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-09-09
  • java自定義序列化的具體使用

    java自定義序列化的具體使用

    本文主要介紹了java自定義序列化的具體使用,文中通過(guò)示例代碼介紹的非常詳細(xì),具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2021-11-11
  • 解決springboot application.yml變灰色的問(wèn)題

    解決springboot application.yml變灰色的問(wèn)題

    這篇文章主要介紹了解決springboot application.yml變灰色的問(wèn)題,具有很好的參考價(jià)值,希望對(duì)大家有所幫助,如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2024-07-07
  • 使用jenkins+maven+git發(fā)布jar包過(guò)程詳解

    使用jenkins+maven+git發(fā)布jar包過(guò)程詳解

    這篇文章主要介紹了使用jenkins+maven+git發(fā)布jar包過(guò)程詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-07-07
  • java開發(fā)之鬧鐘的實(shí)現(xiàn)代碼

    java開發(fā)之鬧鐘的實(shí)現(xiàn)代碼

    本篇文章介紹了,在java中鬧鐘的實(shí)現(xiàn)代碼。需要的朋友參考下
    2013-05-05
  • Java Lambda表達(dá)式原理及多線程實(shí)現(xiàn)

    Java Lambda表達(dá)式原理及多線程實(shí)現(xiàn)

    這篇文章主要介紹了Java Lambda表達(dá)式原理及多線程實(shí)現(xiàn),文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2020-07-07
  • java web中使用cookie記住用戶的賬號(hào)和密碼

    java web中使用cookie記住用戶的賬號(hào)和密碼

    這篇文章主要介紹了java web中使用cookie記住用戶的賬號(hào)和密碼的相關(guān)資料,需要的朋友可以參考下
    2017-01-01
  • java 多線程交通信號(hào)燈模擬過(guò)程詳解

    java 多線程交通信號(hào)燈模擬過(guò)程詳解

    這篇文章主要介紹了java 多線程交通信號(hào)燈模擬過(guò)程詳解,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友可以參考下
    2019-07-07
  • javassist使用指南

    javassist使用指南

    這篇文章主要介紹了javassist的使用方法,文中講解非常細(xì)致,代碼幫助大家更好的理解和學(xué)習(xí),感興趣的朋友可以了解下
    2020-07-07
  • 使用C3P0改造JDBC對(duì)數(shù)據(jù)庫(kù)的連接

    使用C3P0改造JDBC對(duì)數(shù)據(jù)庫(kù)的連接

    這篇文章主要為大家詳細(xì)介紹了使用C3P0改造JDBC對(duì)數(shù)據(jù)庫(kù)的連接,具有一定的參考價(jià)值,感興趣的小伙伴們可以參考一下
    2019-08-08

最新評(píng)論

澄城县| 伊春市| 南开区| 武鸣县| 琼结县| 新营市| 元江| 囊谦县| 阿克苏市| 建平县| 武清区| 白朗县| 大足县| 桂林市| 兰西县| 怀宁县| 朔州市| 宝兴县| 天全县| 新巴尔虎左旗| 长岭县| 滨海县| 竹北市| 凤庆县| 舞钢市| 巴里| 富宁县| 福鼎市| 崇礼县| 汝城县| 额济纳旗| 满洲里市| 邻水| 土默特右旗| 泾川县| 雷州市| 湟中县| 买车| 旌德县| 兴隆县| 辉南县|