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

8張圖帶你全面了解Java?kafka的核心機(jī)制

 更新時間:2023年05月17日 10:51:46   作者:JAVA旭陽  
kafka是目前企業(yè)中很常用的消息隊列產(chǎn)品,可以用于削峰、解耦、異步通信,本文就通過幾張圖帶大家全面認(rèn)識一下kafka,現(xiàn)在我們不妨帶入kafka設(shè)計者的角度去思考該如何設(shè)計,它的架構(gòu)是怎么樣的、都有哪些組件組成、如何進(jìn)行擴(kuò)展等等,需要的朋友可以參考下

kafka基礎(chǔ)架構(gòu)

現(xiàn)在假如有100T大小的消息要發(fā)送到kafka中,數(shù)據(jù)量非常大,一臺機(jī)器存儲不下,面對這種情況,你該如何設(shè)計呢?

很簡單,分而治之,一臺不夠,那就多臺,這就形成了一個kafka集群。如下圖所示,一個broker就是一個kafka節(jié)點(diǎn),100T數(shù)據(jù)就有3個節(jié)點(diǎn)分擔(dān),每個節(jié)點(diǎn)約33T,這樣就能解決問題了,還能提高吞吐量。

  • Topic: 可以理解為一個隊列,一個kafka集群中可以定義很多的topic,比如上圖中的topicA
  • Partition: 為了實(shí)現(xiàn)擴(kuò)展性,提高吞吐量,一個非常大的 topic 可以分布到多個 broker(即服務(wù)器)上,一個 topic 可以分為多個 partition,每個 partition 是一個有序的隊列。比如上圖中的topicA被分成了3個partition
  • Replica: 副本,如果數(shù)據(jù)只放在一個broker中,萬一這個broker宕機(jī)了怎么辦?為了實(shí)現(xiàn)高可用,一個 topic 的每個分區(qū)都有若干個副本,一個 Leader 和若干個Follower。比如上圖中的虛線連接的就是它的副本。
  • Leader: 每個分區(qū)多個副本的“主”,生產(chǎn)者發(fā)送數(shù)據(jù)的對象,以及消費(fèi)者消費(fèi)數(shù)據(jù)的對象都是 Leader。
  • Follower: 每個分區(qū)多個副本中的“從”,實(shí)時從 Leader 中同步數(shù)據(jù),保持和Leader 數(shù)據(jù)的同步。Leader 發(fā)生故障時,某個 Follower 會成為新的 Leader。
  • Producer: 消息生產(chǎn)者,就是向 Kafka broker發(fā)消息的客戶端,后面詳細(xì)講解。
  • Consumer: 消息消費(fèi)者,向 Kafka broker 取消息的客戶端,多個Consumer會組成一個消費(fèi)者組,后面詳細(xì)講解。
  • Zookeeper:用來記錄kafka中的一些元數(shù)據(jù),比如kafka集群中的broker,leader是誰等等,但Kafka2.8.0版本以后也支持非zk的方式,大大減少了和zk的交互。

kafka生產(chǎn)者流程

前面通過一張圖片講解了kafka整體的架構(gòu),那現(xiàn)在我們來看看kafka生產(chǎn)者發(fā)送的整個過程,這里面也是大有文章。

在消息發(fā)送的過程中,涉及到了兩個線程——main 線程和 Sender 線程。在 main 線程中創(chuàng)建了一個雙端隊列 RecordAccumulator。main 線程將消息發(fā)送給 RecordAccumulatorSender 線程不斷從 RecordAccumulator 中拉取消息發(fā)送到 Kafka Broker。

  • 在主線程中由 kafkaProducer 創(chuàng)建消息,然后通過可能的攔截器、序列化器和分區(qū)器的作用之后緩存到消息累加器(RecordAccumulator, 也稱為消息收集器)中。
  • 攔截器: 可以用來在消息發(fā)送前做一些準(zhǔn)備工作,比如按照某個規(guī)則過濾不符合要求的消息、修改消息的內(nèi)容等,也可以用來在發(fā)送回調(diào)邏輯前做一些定制化的需求,比如統(tǒng)計類工作。
  • 序列化器: 用于在網(wǎng)絡(luò)傳輸中將數(shù)據(jù)序列化為字節(jié)流進(jìn)行傳輸,保證數(shù)據(jù)不會丟失。
  • 分區(qū)器: 用于按照一定的規(guī)則將數(shù)據(jù)分發(fā)到不同的kafka broker節(jié)點(diǎn)中
  • Sender 線程負(fù)責(zé)從 RecordAccumulator 獲取消息并將其發(fā)送到 Kafka 中。
  • RecordAccumulator 主要用來緩存消息以便 Sender 線程可以批量發(fā)送,進(jìn)而減少網(wǎng)絡(luò)傳輸?shù)馁Y源消耗以提升性能。
  • RecordAccumulator 緩存的大小可以通過生產(chǎn)者客戶端參數(shù) buffer.memory 配置,默認(rèn)值為 33554432B ,即 32M
  • 主線程中發(fā)送過來的消息都會被迫加到 RecordAccumulator 的某個雙端隊列( Deque )中,RecordAccumulator 內(nèi)部為每個分區(qū)都維護(hù)了一個雙端隊列,即 Deque<ProducerBatch>, 消息寫入緩存時,追加到雙端隊列的尾部。
  • Sender 讀取消息時,從雙端隊列的頭部讀取。ProducerBatch 是指一個消息批次;與此同時,會將較小的 ProducerBatch 湊成一個較大 ProducerBatch ,也可以減少網(wǎng)絡(luò)請求的次數(shù)以提升整體的吞吐量。ProducerBatch 大小可以通過batch.size 控制,默認(rèn)16kb。
  • Sender 線程會在有數(shù)據(jù)積累到batch.size,默認(rèn)16kb,或者如果數(shù)據(jù)遲遲未達(dá)到batch.sizeSender線程等待linger.ms設(shè)置的時間到了之后就會獲取數(shù)據(jù)。linger.ms單位ms,默認(rèn)值是0ms,表示沒有延遲。
  • SenderRecordAccumulator 獲取緩存的消息之后,會將數(shù)據(jù)封裝成網(wǎng)絡(luò)請求<Node,Request> 的形式,這樣就可以將 Request 請求發(fā)往各個 Node 了。
  • 請求在從 sender 線程發(fā)往 Kafka 之前還會保存到 InFlightRequests 中,它的主要作用是緩存了已經(jīng)發(fā)出去但還沒有收到服務(wù)端響應(yīng)的請求。InFlightRequests默認(rèn)每個分區(qū)下最多緩存5個請求,可以通過配置參數(shù)為max.in.flight.request.per. connection修改。
  • 請求Request通過通道Selector發(fā)送到kafka節(jié)點(diǎn)。
  • 發(fā)送后,需要等待kafka的應(yīng)答機(jī)制,取決于配置項(xiàng)acks.
  • 0:生產(chǎn)者發(fā)送過來的數(shù)據(jù),不需要等待數(shù)據(jù)落盤就應(yīng)答。
  • 1:生產(chǎn)者發(fā)送過來的數(shù)據(jù),Leader 收到數(shù)據(jù)后應(yīng)答。
  • -1(all):生產(chǎn)者發(fā)送過來的數(shù)據(jù),Leader和副本節(jié)點(diǎn)收齊數(shù)據(jù)后應(yīng)答。默認(rèn)值是-1,-1 和all 是等價的。
  • Request請求接受到kafka的響應(yīng)結(jié)果,如果成功的話,從InFlightRequests 清除請求,否則的話需要進(jìn)行重發(fā)操作,可以通過配置項(xiàng)retries決定,當(dāng)消息發(fā)送出現(xiàn)錯誤的時候,系統(tǒng)會重發(fā)消息。retries表示重試次數(shù)。默認(rèn)是 int 最大值,2147483647
  • 清理消息累加器RecordAccumulator 中的數(shù)據(jù)。

kafka消費(fèi)者流程

原來kafka生產(chǎn)者發(fā)送經(jīng)過了這么多流程,我們現(xiàn)在來看看kafka消費(fèi)者又是如何進(jìn)行的呢?

Kafka 中的消費(fèi)是基于拉取模式的。消息的消費(fèi)一般有兩種模式:推送模式和拉取模式。推模式是服務(wù)端主動將消息推送給消費(fèi)者,而拉模式是消費(fèi)者主動向服務(wù)端發(fā)起請求來拉取消息。

kafka是以消費(fèi)者組進(jìn)行消費(fèi)的,一個消費(fèi)者組,由多個consumer組成。形成一個消費(fèi)者組的條件,是所有消費(fèi)者的groupid相同。

  • 消費(fèi)者組內(nèi)每個消費(fèi)者負(fù)責(zé)消費(fèi)不同分區(qū)的數(shù)據(jù),一個分區(qū)只能由一個組內(nèi)消費(fèi)者消費(fèi)。如果向消費(fèi)組中添加更多的消費(fèi)者,超過主題分區(qū)數(shù)量,則有一部分消費(fèi)者就會閑置,不會接收任何消息。
  • 消費(fèi)者組之間互不影響。所有的消費(fèi)者都屬于某個消費(fèi)者組,即消費(fèi)者組是邏輯上的一個訂閱者。

那么問題來了,kafka是如何指定消費(fèi)者組的每個消費(fèi)者消費(fèi)哪個分區(qū)?每次消費(fèi)的數(shù)量是多少呢?

一、如何制定消費(fèi)方案

  • 消費(fèi)者consumerA,consumerB, consumerC向kafka集群中的協(xié)調(diào)器coordinator發(fā)送JoinGroup的請求。coordinator主要是用來輔助實(shí)現(xiàn)消費(fèi)者組的初始化和分區(qū)的分配。
  • coordinator老大節(jié)點(diǎn)選擇 = groupidhashcode值 % 50( __consumer_offsets內(nèi)置主題位移的分區(qū)數(shù)量)例如: groupid的hashcode值 為1,1% 50 = 1,那么__consumer_offsets 主題的1號分區(qū),在哪個broker上,就選擇這個節(jié)點(diǎn)的coordinator作為這個消費(fèi)者組的老大。消費(fèi)者組下的所有的消費(fèi)者提交offset的時候就往這個分區(qū)去提交offset。
  • 選出一個 consumer作為消費(fèi)中的leader,比如上圖中的ConsumerB。
  • 消費(fèi)者leader制定出消費(fèi)方案,比如誰來消費(fèi)哪個分區(qū)等
  • 把消費(fèi)方案發(fā)給coordinator
  • 最后coordinator就把消費(fèi)方 案下發(fā)給各個consumer, 圖中只畫了一條線,實(shí)際上是有下發(fā)各個consumer。

注意,每個消費(fèi)者都會和coordinator保持心跳(默認(rèn)3s),一旦超時(session.timeout.ms=45s),該消費(fèi)者會被移除,并觸發(fā)再平衡;或者消費(fèi)者處理消息的時間過長(max.poll.interval.ms=5分鐘),也會觸發(fā)再平衡,也就是重新進(jìn)行上面的流程。

二、消費(fèi)者消費(fèi)細(xì)節(jié)

現(xiàn)在已經(jīng)初始化消費(fèi)者組信息,知道哪個消費(fèi)者消費(fèi)哪個分區(qū),接著我們來看看消費(fèi)者細(xì)節(jié)。

  • 消費(fèi)者創(chuàng)建一個網(wǎng)絡(luò)連接客戶端ConsumerNetworkClient, 發(fā)送消費(fèi)請求,可以進(jìn)行如下配置:
  • fetch.min.bytes: 每批次最小抓取大小,默認(rèn)1字節(jié)
  • fetch.max.bytes: 每批次最大抓取大小,默認(rèn)50M
  • fetch.max.wait.ms:最大超時時間,默認(rèn)500ms
  • 發(fā)送請求到kafka集群
  • 成功的回調(diào),會將數(shù)據(jù)保存到completedFetches隊列中
  • 消費(fèi)者從隊列中抓取數(shù)據(jù),根據(jù)配置max.poll.records一次拉取數(shù)據(jù)返回消息的最大條數(shù),默認(rèn)500條。
  • 獲取到數(shù)據(jù)后,需要經(jīng)過反序列化器、攔截器等。

kafka的存儲機(jī)制

我們都知道消息發(fā)送到kafka,最終是存儲到磁盤中的,我們看下kafka是如何存儲的。

一個topic分為多個partition,每個partition對應(yīng)于一個log文件,為防止log文件過大導(dǎo)致數(shù)據(jù)定位效率低下,Kafka采取了分片和索引機(jī)制,每個partition分為多個segment。每個segment包括:“.index”文件、“.log”文件和.timeindex等文件,Producer生產(chǎn)的數(shù)據(jù)會被不斷追加到該log文件末端。

上圖中t1即為一個topic的名稱,而“t1-0/t1-1”則表明這個目錄是t1這個topic的哪個partition。

kafka中的索引文件以稀疏索引(sparseindex)的方式構(gòu)造消息的索引,如下圖所示:

1.根據(jù)目標(biāo)offset定位segment文件

2.找到小于等于目標(biāo)offset的最大offset對應(yīng)的索引項(xiàng)

3.定位到log文件

4.向下遍歷找到目標(biāo)Record

注意:index為稀疏索引,大約每往log文件寫入4kb數(shù)據(jù),會往index文件寫入一條索引。通過參數(shù)log.index.interval.bytes控制,默認(rèn)4kb

那kafka中磁盤文件保存多久呢?

kafka 中默認(rèn)的日志保存時間為 7 天,可以通過調(diào)整如下參數(shù)修改保存時間。

  • log.retention.hours,最低優(yōu)先級小時,默認(rèn) 7 天。
  • log.retention.minutes,分鐘。
  • log.retention.ms,最高優(yōu)先級毫秒。
  • log.retention.check.interval.ms,負(fù)責(zé)設(shè)置檢查周期,默認(rèn) 5 分鐘。

總結(jié)

其實(shí)kafka中的細(xì)節(jié)十分多,本文也只是對kafka的一些核心機(jī)制從理論層面做了一個總結(jié),更多的細(xì)節(jié)還是需要自行去實(shí)踐,去學(xué)習(xí)。

以上就是8張圖帶你全面了解Java kafka的核心機(jī)制的詳細(xì)內(nèi)容,更多關(guān)于Java kafka核心機(jī)制的資料請關(guān)注腳本之家其它相關(guān)文章!

相關(guān)文章

  • Spring ApplicationListener的使用詳解

    Spring ApplicationListener的使用詳解

    這篇文章主要介紹了Spring ApplicationListener的使用,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧
    2020-06-06
  • SpringBoot超詳細(xì)深入講解底層原理

    SpringBoot超詳細(xì)深入講解底層原理

    我們知道springboot內(nèi)部是通過spring框架內(nèi)嵌Tomcat實(shí)現(xiàn)的,當(dāng)然也可以內(nèi)嵌jetty,undertow等等web框架;另外springboot還有一個特別重要的功能就是自動裝配,這又是如何實(shí)現(xiàn)的呢
    2022-07-07
  • SpringBoot攔截器使用精講

    SpringBoot攔截器使用精講

    攔截器可以根據(jù) URL 對請求進(jìn)行攔截,主要應(yīng)用于登陸校驗(yàn)、權(quán)限驗(yàn)證、亂碼解決、性能監(jiān)控和異常處理等功能上。SpringBoot同樣提供了攔截器功能。 本文將為大家詳細(xì)介紹一下
    2021-12-12
  • Spring AOP事務(wù)管理的示例詳解

    Spring AOP事務(wù)管理的示例詳解

    這篇文章將通過轉(zhuǎn)賬案例為大家詳細(xì)介紹一下Spring AOP是如何進(jìn)行事務(wù)管理的,文中的示例代碼講解詳細(xì),感興趣的小伙伴可以了解一下
    2022-06-06
  • 詳解使用spring cloud config來統(tǒng)一管理配置文件

    詳解使用spring cloud config來統(tǒng)一管理配置文件

    這篇文章主要介紹了詳解使用spring cloud config來統(tǒng)一管理配置文件,小編覺得挺不錯的,現(xiàn)在分享給大家,也給大家做個參考。一起跟隨小編過來看看吧
    2018-12-12
  • Java中的靜態(tài)代碼塊與構(gòu)造代碼塊詳解

    Java中的靜態(tài)代碼塊與構(gòu)造代碼塊詳解

    靜態(tài)代碼塊在JVM加載類時執(zhí)行,僅執(zhí)行一次,用于初始化靜態(tài)變量或執(zhí)行一次性設(shè)置,構(gòu)造代碼塊在每次創(chuàng)建對象時執(zhí)行,用于實(shí)例變量的初始化或執(zhí)行創(chuàng)建對象時的公共邏輯,靜態(tài)代碼塊按定義順序執(zhí)行,構(gòu)造代碼塊在構(gòu)造方法前執(zhí)行
    2024-11-11
  • Redis中String字符串和sdshdr結(jié)構(gòu)體超詳細(xì)講解

    Redis中String字符串和sdshdr結(jié)構(gòu)體超詳細(xì)講解

    這篇文章主要介紹了Redis中String字符串和sdshdr結(jié)構(gòu)體,文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧
    2023-04-04
  • java通過Idea遠(yuǎn)程一鍵部署springboot到Docker詳解

    java通過Idea遠(yuǎn)程一鍵部署springboot到Docker詳解

    這篇文章主要介紹了java通過Idea遠(yuǎn)程一鍵部署springboot到Docker詳解,Idea是Java開發(fā)利器,springboot是Java生態(tài)中最流行的微服務(wù)框架,docker是時下最火的容器技術(shù),那么它們結(jié)合在一起會產(chǎn)生什么化學(xué)反應(yīng)呢?的相關(guān)資料
    2019-06-06
  • Java 8新特性方法引用詳細(xì)介紹

    Java 8新特性方法引用詳細(xì)介紹

    這篇文章主要介紹了Java 8新特性方法引用詳細(xì)介紹的相關(guān)資料,這里對新特性 方法引用做的資料整理,具有參考價值,需要的朋友可以參考下
    2016-12-12
  • hutool實(shí)戰(zhàn):IoUtil 流操作工具類(將內(nèi)容寫到流中)

    hutool實(shí)戰(zhàn):IoUtil 流操作工具類(將內(nèi)容寫到流中)

    這篇文章主要介紹了Go語言的io.ioutil標(biāo)準(zhǔn)庫使用,是Golang入門學(xué)習(xí)中的基礎(chǔ)知識,需要的朋友可以參考下,如果能給你帶來幫助,請多多關(guān)注腳本之家的其他內(nèi)容
    2021-06-06

最新評論

乳源| 安龙县| 丹巴县| 鄂伦春自治旗| 枣庄市| 县级市| 潮州市| 南平市| 永春县| 汽车| 东宁县| 黄梅县| 五莲县| 柘城县| 乌拉特中旗| 新郑市| 柳林县| 大田县| 临颍县| 德格县| 左贡县| 高阳县| 仙居县| 潞西市| 瓮安县| 普格县| 辰溪县| 琼结县| 金沙县| 通城县| 磐安县| 兴隆县| 疏勒县| 广宁县| 阳新县| 中山市| 白山市| 体育| 鄂州市| 中超| 元朗区|