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

@KafkaListener 如何使用

 更新時(shí)間:2023年02月20日 09:23:32   作者:DemonHunter211  
這篇文章主要介紹了@KafkaListener 如何使用,本文通過(guò)圖文實(shí)例代碼相結(jié)合給大家詳細(xì)講解,文末給大家介紹了kafka的消費(fèi)者分區(qū)分配策略,需要的朋友可以參考下

@KafkaListener 如何使用

spring-kafka使用基于@KafkaListener注解,@KafkaListener使用方式如下

@KafkaListener(topics = "topic1")
public void? ?kafkaListen(List<ConsumerRecord<xxx, xxx>> records) {
? ? ...
}

在注解內(nèi)指定topic名稱,當(dāng)對(duì)應(yīng)的topic內(nèi)有新的消息時(shí),testListen方法會(huì)被調(diào)用,參數(shù)就是topic內(nèi)新的消息。這個(gè)過(guò)程是異步進(jìn)行的。

@KafkaListener工作流程主要有以下幾步:

解析;解析@KafkaListener注解。
注冊(cè);解析后的數(shù)據(jù)注冊(cè)到spring-kafka。
監(jiān)聽(tīng);開(kāi)始監(jiān)聽(tīng)topic變更。
調(diào)用;調(diào)用注解標(biāo)識(shí)的方法,將監(jiān)聽(tīng)到的數(shù)據(jù)作為參數(shù)傳入。
下面我們一步一步分析

解析

@KafkaListener注解由KafkaListenerAnnotationBeanPostProcessor類解析,后者實(shí)現(xiàn)了BeanPostProcessor接口,這個(gè)接口如下

public interface BeanPostProcessor {

    Object postProcessBeforeInitialization(Object bean, String beanName) throws BeansException;

    Object postProcessAfterInitialization(Object bean, String beanName) throws BeansException;
}

接口內(nèi)部有2個(gè)方法,分別在bean初始化前后被調(diào)用。

KafkaListenerAnnotationBeanPostProcessor內(nèi)會(huì)在postProcessAfterInitialization方法內(nèi)解析@KafkaListener注解。

注冊(cè)
解析步驟里,我們可以獲取到所有含有@KafkaListener注解的類,之后這些類的相關(guān)信息會(huì)被注冊(cè)到 KafkaListenerEndpointRegistry內(nèi),包括注解所在的方法,當(dāng)前的bean等。KafkaListenerEndpointRegistry這個(gè)類內(nèi)部會(huì)維護(hù)多個(gè)Listener Container,每一個(gè)@KafkaListener都會(huì)對(duì)應(yīng)一個(gè)Listener Container。并且每個(gè)Container對(duì)應(yīng)一個(gè)線程。

監(jiān)聽(tīng)
注冊(cè)完成之后,每個(gè)Listener Container會(huì)開(kāi)始工作,會(huì)新啟一個(gè)新的線程,初始化KafkaConsumer,監(jiān)聽(tīng)topic變更等。

調(diào)用
監(jiān)聽(tīng)到數(shù)據(jù)之后,container會(huì)組織消息的格式,隨后調(diào)用解析得到的@KafkaListener注解標(biāo)識(shí)的方法,將組織后的消息作為參數(shù)傳入方法,執(zhí)行用戶邏輯。

@KafkaListener和@KafkaListners

@KafkaListeners是@KafkaListener的Container Annotation,這也是jdk8的新特性之一,注解可以重復(fù)標(biāo)注。

@KafkaListeners({@KafkaListener(topics="topic1"), @KafkaListener(topics="topic2")})
public void listen(ConsumerRecord<Integer, String> msg) {}
 
等同于
 
@KafkaListener(topics="topic1")
@KafkaListener(topics="topic2")
public void listen(ConsumerRecord<Integer, String> msg) {}

擴(kuò)展:kafka的消費(fèi)者分區(qū)分配策略

kafka有三種分區(qū)分配策略

1. RoundRobin

2. Range

3. Sticky

1. RoundRobin

(1)把所有topic的分區(qū)partition放入一個(gè)隊(duì)列中,按照name的hashcode進(jìn)行排序;

(2)把consumer放在一個(gè)循環(huán)隊(duì)列,按照name的hashcode進(jìn)行排序;

(3)循環(huán)遍歷consumer,從partition隊(duì)列pop出一個(gè)partition,分配給當(dāng)前consumer;以此類推,取下一個(gè)consumer,繼續(xù)從partition隊(duì)列pop出來(lái)分配給當(dāng)前consumer;直到partition隊(duì)列中的元素被分配完;

2. Range

(1)假設(shè)topicA有4個(gè)分區(qū),topicB有5個(gè)分區(qū),topicC有6個(gè)分區(qū);一共有3個(gè)consumer;

(2)遍歷3個(gè)topic的分區(qū)集合,先取topicA的分區(qū)集合,然后準(zhǔn)備依次給3個(gè)consumer分配分區(qū);對(duì)于第1個(gè)consumer,所分配的分區(qū)數(shù)量根據(jù)以下公式:假設(shè)消費(fèi)者數(shù)量為N,當(dāng)前主題剩下的分區(qū)數(shù)量為M,則當(dāng)前消費(fèi)者應(yīng)該分配的分區(qū)數(shù)量 = M%N==0? M/N +1 : M/N ;按照公式,3個(gè)消費(fèi)者應(yīng)該分配的分區(qū)數(shù)量依次為:2/1/1,即topicA-partition-0/1分配給consumer-0,topicA-partition-2分配給consumer-1,topicA-partition-3分配給consumer-2;

(3)按照上述規(guī)則按序把topicB和topicC的分區(qū)分配給3個(gè)consumer;依次為:2/2/1,2/2/2;

3. Sticky

kafka在0.11版本引入了Sticky分區(qū)分配策略,它的兩個(gè)主要目的是:

1. 分區(qū)的分配要盡可能的均勻,分配給消費(fèi)者者的主題分區(qū)數(shù)最多相差一個(gè);

2. 分區(qū)的分配盡可能的與上次分配的保持相同;

當(dāng)兩者發(fā)生沖突時(shí),第一個(gè)目標(biāo)優(yōu)先于第二個(gè)目標(biāo);

粘性分區(qū)是由Kafka從0.11x版本開(kāi)始引入的分配策略,首先會(huì)盡量均衡的分配分區(qū)到消費(fèi)者上面,在出現(xiàn)同一消費(fèi)組內(nèi)消費(fèi)者出現(xiàn)問(wèn)題的時(shí)候,會(huì)盡量保持原來(lái)的分配的分區(qū)不變;

Sticky分區(qū)初始分配分區(qū)的方法與Range相似,但是不同;拿7個(gè)分區(qū)3個(gè)消費(fèi)者為例,消費(fèi)者消費(fèi)的分區(qū)依舊是3/2/2,但是不同與Range的是Range分區(qū)是排好序的,但是Sticky分區(qū)是隨機(jī)的;

到此這篇關(guān)于@KafkaListener 如何使用方式的文章就介紹到這了,更多相關(guān)@KafkaListener 使用內(nèi)容請(qǐng)搜索腳本之家以前的文章或繼續(xù)瀏覽下面的相關(guān)文章希望大家以后多多支持腳本之家!

相關(guān)文章

  • Java+Selenium實(shí)現(xiàn)控制瀏覽器的啟動(dòng)選項(xiàng)Options

    Java+Selenium實(shí)現(xiàn)控制瀏覽器的啟動(dòng)選項(xiàng)Options

    這篇文章主要為大家詳細(xì)介紹了如何使用java代碼利用selenium控制瀏覽器的啟動(dòng)選項(xiàng)Options的代碼操作,文中的示例代碼講解詳細(xì),感興趣的可以了解一下
    2023-01-01
  • Java中對(duì)于雙屬性枚舉的使用案例

    Java中對(duì)于雙屬性枚舉的使用案例

    今天小編就為大家分享一篇關(guān)于Java中對(duì)于雙屬性枚舉的使用案例,小編覺(jué)得內(nèi)容挺不錯(cuò)的,現(xiàn)在分享給大家,具有很好的參考價(jià)值,需要的朋友一起跟隨小編來(lái)看看吧
    2018-12-12
  • Java使用BigDecimal公式精確計(jì)算及精度丟失問(wèn)題

    Java使用BigDecimal公式精確計(jì)算及精度丟失問(wèn)題

    在工作中經(jīng)常會(huì)遇到數(shù)值精度問(wèn)題,比如說(shuō)使用float或者double的時(shí)候,可能會(huì)有精度丟失問(wèn)題,下面這篇文章主要給大家介紹了關(guān)于Java使用BigDecimal公式精確計(jì)算及精度丟失問(wèn)題的相關(guān)資料,需要的朋友可以參考下
    2023-01-01
  • IDEA中Spring項(xiàng)目的工程構(gòu)建

    IDEA中Spring項(xiàng)目的工程構(gòu)建

    這篇文章主要介紹了IDEA中Spring項(xiàng)目的工程構(gòu)建,Spring框架是輕量級(jí)的JavaEE框架,可以解決企業(yè)應(yīng)用開(kāi)發(fā)的復(fù)雜性,有兩個(gè)核心部分:IOC和Aop,今天來(lái)學(xué)習(xí)如何構(gòu)建spring項(xiàng)目,需要的朋友可以參考下
    2023-05-05
  • Java 線程池ThreadPoolExecutor源碼解析

    Java 線程池ThreadPoolExecutor源碼解析

    這篇文章主要介紹了Java 線程池ThreadPoolExecutor源碼解析
    2022-03-03
  • JPA使用樂(lè)觀鎖應(yīng)對(duì)高并發(fā)方式

    JPA使用樂(lè)觀鎖應(yīng)對(duì)高并發(fā)方式

    這篇文章主要介紹了JPA使用樂(lè)觀鎖應(yīng)對(duì)高并發(fā)方式,具有很好的參考價(jià)值,希望對(duì)大家有所幫助。如有錯(cuò)誤或未考慮完全的地方,望不吝賜教
    2021-10-10
  • java模擬http請(qǐng)求的錯(cuò)誤問(wèn)題整理

    java模擬http請(qǐng)求的錯(cuò)誤問(wèn)題整理

    本文是小編給大家整理的在用java模擬http請(qǐng)求的時(shí)候遇到的錯(cuò)誤問(wèn)題整理,以及相關(guān)分析,有興趣的朋友參考下。
    2018-05-05
  • 詳解Java8新特性如何防止空指針異常

    詳解Java8新特性如何防止空指針異常

    要說(shuō) Java 編程中哪個(gè)異常是你印象最深刻的,那 NullPointerException 空指針可以說(shuō)是臭名昭著的,不要說(shuō)初級(jí)程序員會(huì)碰到, 即使是中級(jí),專家級(jí)程序員稍不留神,就會(huì)掉入這個(gè)坑里,本文就和大家聊聊Java8新特性如何防止空指針異常
    2023-08-08
  • java利用jieba進(jìn)行分詞的實(shí)現(xiàn)

    java利用jieba進(jìn)行分詞的實(shí)現(xiàn)

    本文主要介紹了在Java中使用jieba-analysis庫(kù)進(jìn)行分詞,文中通過(guò)示例代碼介紹的非常詳細(xì),對(duì)大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來(lái)一起學(xué)習(xí)學(xué)習(xí)吧
    2025-03-03
  • SpringBoot統(tǒng)一返回處理出現(xiàn)cannot?be?cast?to?java.lang.String異常解決

    SpringBoot統(tǒng)一返回處理出現(xiàn)cannot?be?cast?to?java.lang.String異常解決

    這篇文章主要給大家介紹了關(guān)于SpringBoot統(tǒng)一返回處理出現(xiàn)cannot?be?cast?to?java.lang.String異常解決的相關(guān)資料,文中通過(guò)圖文介紹的非常詳細(xì),需要的朋友可以參考下
    2023-09-09

最新評(píng)論

睢宁县| 临泽县| 湘潭县| 宜宾县| 运城市| 郧西县| 阿合奇县| 木里| 岗巴县| 乌什县| 壶关县| 奈曼旗| 卓资县| 四会市| 清河县| 明水县| 海安县| 南和县| 淳化县| 湟源县| 金秀| 广昌县| 蓬莱市| 大庆市| 汤阴县| 宜州市| 宁海县| 弥渡县| 定远县| 肇源县| 荥经县| 南汇区| 衡水市| 永胜县| 容城县| 灵宝市| 南川市| 东乡县| 衡阳市| 万山特区| 泉州市|