主頁 > 資料庫 > kafka學習筆記

kafka學習筆記

2022-01-19 17:20:58 資料庫

1.1 初識kafka

     kafka 是一款基于發布與訂閱的訊息系統,

名詞 解釋
broker 訊息系統處理的一個節點,一個kafka服務器被稱為一個broker,多個broker被稱為kafka集群
topic kafka中的訊息歸類,相當于mysql的表,發布到kafka的每條訊息都要有topic
partition 磁區,每個topic由多個partition組成,相當于mysql的磁區表,partition是訊息物理上的集合,topic是邏輯集合,
producer 訊息生產者,發送訊息的客戶端
consumer 訊息消費者,消費訊息的客戶端
offset 偏移量,記錄消費者消費到哪里了,保存在zk上或者kafka上
consumerGroup 消費者組,每個消費者都屬于一個consumerGroup,多個消費者可以組成一個group,同一個組的消費者只能消費一條topic訊息,不同組可以重復消費同一條訊息,比如對于創建訂單訊息,A 系統可消費,B系統也可以消費,互相之前不影響,

 Controller Broker

    在分布式系統中,通常需要有一個協調者,該協調者會在分布式系統發生例外時發揮特殊的作用,在Kafka中該協調者稱之為控制器(Controller),其實該控制器并沒有什么特殊之處,它本身也是一個普通的Broker,只不過需要負責一些額外的作業(追蹤集群中的其他Broker,并在合適的時候處理新加入的和失敗的Broker節點、Rebalance磁區、分配新的leader磁區等),值得注意的是:Kafka集群中始終只有一個Controller Broker,

Controller Broker的主要職責有很多,主要是一些管理行為,主要包括以下幾個方面:

  • 創建、洗掉主題,增加磁區并分配leader磁區
  • 集群Broker管理(新增 Broker、Broker 主動關閉、Broker 故障)
  • partition leader選舉
  • 磁區重分配

 

一個簡單的kafka集群

 

  topic 和 partition

       kafka的訊息是通過topic進行分類的,相當于資料庫的表,一類業務的資料集合,topic可以被分成若干個磁區,訊息以追加的方式寫入磁區,然后以先進先出的方式順序讀取

注意一個topic包含多個磁區,因此無法保持訊息的整體順序性,只能保證單個磁區的有序性,多磁區設計可以提高程式性能,在我們的流量變大的時候,我們可以增加磁區數來提高性能,

消費者在消費磁區里面的訊息的時候,會有一個offset來記錄消費到哪里了,這里有同學會有疑問,最新加入的消費者應該從哪里消費了,有引數可以控制從哪里開始消費

      #earliest
      #當各磁區下有已提交的offset時,從提交的offset開始消費;無提交的offset時,從頭開始消費
      #latest
      #當各磁區下有已提交的offset時,從提交的offset開始消費;無提交的offset時,消費新產生的該磁區下的資料
      auto-offset-reset: earliest

 

1.2 生產者

 1.2.1發送方式

  • 異步發送,我們不關心是否正常到達,可能會丟失一些訊息
    kafkaTemplate.send(topic, message)
  • 同步發送,使用send()方法的時候,會回傳一個future物件,我們使用get()方法,阻塞執行緒等待結果的回傳
kafkaTemplate.send(topic, message).get(10, TimeUnit.SECONDS);
log.info("kafka訊息發送成功:topic:{}, message:{}", topic, message);
  • 異步回呼,kafka api提供回呼鉤子,使我們可以監聽發送成功和失敗的事件,
kafkaTemplate.send(topic, JSON.toJSONString(t)).addCallback(new ListenableFutureCallback() {
    @Override
    public void onFailure(Throwable throwable) {
        log.error("發送topic{}產生例外資訊:{}", topic, throwable.getMessage(), throwable);
    }

    @Override
    public void onSuccess(Object o) {
        log.info("發送成功topic:{}的訊息:{}", topic,JSON.toJSONString(t));
    }
});

1.2.2 生產者引數配置

   1.acks 可以保證生產者不丟訊息,

  • acks=0
    意思就是我的KafkaProducer在客戶端,只要把訊息發送出去,不管發送出去的資料有沒有同步完成,都不用管,直接就認為這個訊息發送成功了,意味著producer不等待broker同步完成的確認,繼續發送下一條(批)資訊,提供了最低的延遲,但是持久性最差;當服務器發生故障時,就很可能發生資料丟失,例如leader已經死亡,producer不知情,還會繼續發送訊息broker接收不到資料就會資料丟失,就比如可能你發送出去的訊息還在半路,結果呢,Partition Leader所在的Broker就直接掛了,然后你的客戶端還認為訊息發送成功了,此時就會導致這條訊息就丟失了,
  • acks=1
    意思就是說只要Partition Leader接收到訊息而且寫入本地磁盤了,就認為成功了,不管其他的Follower有沒有同步過去這條訊息了,這種設定其實是kafka默認的設定,大家請注意,劃重點!這是默認的設定,也就是說,默認情況下,你要是不重新設定acks這個引數,只要Partition Leader寫成功導磁盤就算成功,
    但是這里有一個問題,萬一Partition Leader剛剛接收到訊息,Follower還沒來得及同步過去,結果Leader所在的broker宕機了,此時也會導致這條訊息丟失,因為客戶端已經認為發送成功了,
    意味著producer要等待leader成功收到資料并得到確認,才發送下一條message,此選項提供了較好的持久性較低的延遲性,但如果Partition的Leader死亡,follwer尚未復制,資料就會丟失,
  • acks=all
    這個意思就是說,Partition Leader接收到訊息之后,還必須要求ISR串列里跟Leader保持同步的那些Follower都要把訊息同步過去,才能認為這條訊息是寫入成功了,
  • acks=-1
    這個和all是一樣的,只是設定成這樣的時候,引數min.insync.replicas才能生效,min.insync.replicas這個引數設定ISR中的最小副本數是多少,默認值為1,為什么要這個引數,就是ISR中如果只有leader的時候,這時候leader也掛了,資料依然丟失,但設定了最小副本數,就能保證最少有一個副本存在,

 

   2.batch.sizelinger.ms

     kafka發送訊息,并不是單次發往服務器的,是把多個訊息打包成一個批次,一次性發往服務器,該引數指定了同一個批次的大小,如果達到這個閾值,就開始發送到服務器,否則留在記憶體中,同時還有linger.ms這個引數,這個是同一個批次資料停留在記憶體的最長時間,如果時間達到了,也會先把批次資料發送到服務器就算訊息大小沒有達到batch.size閾值,也就是這兩個閾值,誰先達到都會把記憶體的訊息發送到服務器,

 

   3. buffer.memory
      該引數用來設定生產者記憶體緩沖區的大小,生產者用它緩沖要發送到服務器的訊息,如果應用程式發送訊息的速度超過發送到服務器的速度,會導致生產者空間不足,這個時候,send() 方法呼叫要么被阻塞,要么拋出例外,取決于如何設定 block.on.buffer.full引數(在0.9.0.0版本里被替換成了max.block.ns,表示在拋出例外之前可以阻塞一段時間),

 

    4.retries

     retries 引數的值決定了生產者在發送訊息的時候如果遇到錯誤,可以重發訊息的次數,每次重試之間等待100ms,也可以通過retry.backoff.ms來改變,

 

   5.max.in.flight.requests.per.connection
       這個引數指定了生產者在收到服務器回應前可以發送多少個訊息,把它設定成1可以保證訊息是按順序發送的服務器的,即使發生了重試,第一個批次在發送失敗后,第二個批次資料是不能發送到服務器的,會等第一個完成才能進行第二個批次發送,但如果是大于1,第二個批次資料就會發送到服務器,順序就會反了,=1可以保證訊息順序,但這樣設定嚴重影響吞吐,

   

   6.max.request.size

     發送單個訊息的最大值,比如1MB

   7.request.timeout.ms   

     指定了生產者在發送資料時等待服務器回傳回應的時間

 

 1.2.3 生產者磁區策略

      kafka有自己的磁區策略的,如果未指定,就會使用默認的磁區策略,如果沒有指定key,那么會使用輪詢的隨機策略均衡的將訊息發送的topic的各磁區上,如果指定了key,Kafka根據傳遞訊息的key來進行磁區的分配,即hash(key) % numPartitions,如果Key相同的話,那么就會分配到統一磁區,當然也可以自定義磁區策略,只要實作Partitioner這個介面,

 

1.3 消費者

    想知道如何從kafka消費訊息(kafka是拉訊息的模式,有些mq是推,),需要先了解消費者和消費者組的概念,

    1.3.1 消費者和消費者組

    kafka 消費者從屬于消費者群組,一個群組里的消費者訂閱同一個topic,組里的每個消費者訊息topic一部分訊息,從而達到高吞吐,在生產者速率大于消費者的時候,我們可以增加消費者來提高消費速率,不同的群組可以消費同一個topic,比如訂單創建訊息,客服系統可用消費,CRM系統也可以消費,不同群組之間消費互不影響,

下圖就展示了消費者組合消費者的關系,一般情況下磁區數=消費者數,一個消費者負責消費一個磁區訊息,如果消費者數小于磁區數,就會出現一個消費者消費兩個磁區的訊息,注意如果消費者數大于磁區數,多余的消費者將不會消費任何訊息(一個磁區只能一個消費者消費)

    1.3.2 如何保證順序消費

  • 我們知道同一個磁區訊息是有序的,那我們可以把訊息固定發送到同一個磁區,但這樣的話就失去分布式的特性了,
  • max.in.flight.requests.per.connection 設定成1,這個引數指定了生產者在收到服務器回應前可以發送多少個訊息,把它設定成1可以保證訊息是按順序發送的服務器的,即使發生了重試

     

    1.3.3 偏移量的提交

      我們每次poll()訊息的時候,總是回傳沒有被消費的訊息,那kafka是如何做到的,它內部有一個偏移量,記錄了當前消費者消費到磁區的什么位置,當我們消費完訊息更新磁區位置的操作叫偏移量提交,

     (1)自動提交 enable.auto.commit=true,每隔5s,消費者會自動吧從poll()方法手動的最大偏移量提交上去,這種方式有個問題,如果在最近一次提交后3s發生了rebalance,那么這時候偏移量沒有提交上去,這3s內的訊息會被重復消費,

     (2)手動提交(同步or異步)enable.auto.commit=false,使用commitSync()同步提交這個方法會重試到提交成功,使用commitAsync()異步提交,異步的不會重試,因為可能另外一個更大的執行緒提交了偏移量,如果我們重試之前比較低的偏移量就會導致重復消費訊息,

 

     1.3.4 rebalance

       kafka 怎么均勻地分配某個 topic 下的所有 partition 到各個消費者,從而使得訊息的消費速度達到最快,這就是平衡(balance),而 rebalance(重平衡)其實就是重新進行 partition 的分配,從而使得 partition 的分配重新達到平衡狀態,以下幾種情況會發送rebalance,

  • 訂閱 Topic 的磁區數發生變化,
  • 訂閱的 Topic 個數發生變化,
  • 消費組內成員個數發生變化,例如有新的 consumer 實體加入該消費組或者離開組,

     特別針對第三種情況說明下,消費者組成員變化也分以下幾種情況

  • 新成員加入
  • 組成員主動離開
  • 組成員崩潰,收不到心跳,或者在一定時間內沒有完成消費,這篇文章寫的挺詳細的,可以看下 https://www.cnblogs.com/chanshuyi/p/kafka_rebalance_quick_guide.html

     幾個重要的消費者引數都會導致rebalance

       session.timeout.ms 表示 consumer 向 broker 發送心跳的超時時間,例如 session.timeout.ms = 180000 表示在最長 180 秒內 broker 沒收到 consumer 的心跳,那么 broker 就認為該 consumer 死亡了,會啟動 rebalance,

       heartbeat.interval.ms 表示 consumer 每次向 broker 發送心跳的時間間隔,heartbeat.interval.ms = 60000 表示 consumer 每 60 秒向 broker 發送一次心跳,一般來說,session.timeout.ms 的值是 heartbeat.interval.ms 值          的 3 倍以上,

       max.poll.interval.ms 表示 consumer 每兩次 poll 訊息的時間間隔,簡單地說,其實就是 consumer 每次消費訊息的時長,如果訊息處理的邏輯很重,那么市場就要相應延長,否則如果時間到了 consumer 還么消費完,broker           會默認認為 consumer 死了,發起 rebalance,

       max.poll.records 表示每次消費的時候,獲取多少條訊息,獲取的訊息條數越多,需要處理的時間越長,所以每次拉取的訊息數不能太多,需要保證在 max.poll.interval.ms 設定的時間內能消費完,否則會發生 rebalance,

 

最后說下kafka高性能的原因,

1.從上面partition可以看出,每個磁區都是順序寫盤,這比隨機寫盤效率高,

2.kafka零拷貝,

3.分布式,多磁區提高了寫入和消費性能,性能不行的時候可以加broker或者增加磁區數,

   

轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/415327.html

標籤:其他

上一篇:MSSQL·WHERE過濾含空格字串的坑

下一篇:AzureDevopsAPI:通過Sha1前綴、分支名稱或標記名稱查找Git提交?

標籤雲
其他(157675) Python(38076) JavaScript(25376) Java(17977) C(15215) 區塊鏈(8255) C#(7972) AI(7469) 爪哇(7425) MySQL(7132) html(6777) 基礎類(6313) sql(6102) 熊猫(6058) PHP(5869) 数组(5741) R(5409) Linux(5327) 反应(5209) 腳本語言(PerlPython)(5129) 非技術區(4971) Android(4554) 数据框(4311) css(4259) 节点.js(4032) C語言(3288) json(3245) 列表(3129) 扑(3119) C++語言(3117) 安卓(2998) 打字稿(2995) VBA(2789) Java相關(2746) 疑難問題(2699) 细绳(2522) 單片機工控(2479) iOS(2429) ASP.NET(2402) MongoDB(2323) 麻木的(2285) 正则表达式(2254) 字典(2211) 循环(2198) 迅速(2185) 擅长(2169) 镖(2155) 功能(1967) .NET技术(1958) Web開發(1951) python-3.x(1918) HtmlCss(1915) 弹簧靴(1913) C++(1909) xml(1889) PostgreSQL(1872) .NETCore(1853) 谷歌表格(1846) Unity3D(1843) for循环(1842)

熱門瀏覽
  • GPU虛擬機創建時間深度優化

    **?桔妹導讀:**GPU虛擬機實體創建速度慢是公有云面臨的普遍問題,由于通常情況下創建虛擬機屬于低頻操作而未引起業界的重視,實際生產中還是存在對GPU實體創建時間有苛刻要求的業務場景。本文將介紹滴滴云在解決該問題時的思路、方法、并展示最終的優化成果。 從公有云服務商那里購買過虛擬主機的資深用戶,一 ......

    uj5u.com 2020-09-10 06:09:13 more
  • 可編程網卡芯片在滴滴云網路的應用實踐

    **?桔妹導讀:**隨著云規模不斷擴大以及業務層面對延遲、帶寬的要求越來越高,采用DPDK 加速網路報文處理的方式在橫向縱向擴展都出現了局限性。可編程芯片成為業界熱點。本文主要講述了可編程網卡芯片在滴滴云網路中的應用實踐,遇到的問題、帶來的收益以及開源社區貢獻。 #1. 資料中心面臨的問題 隨著滴滴 ......

    uj5u.com 2020-09-10 06:10:21 more
  • 滴滴資料通道服務演進之路

    **?桔妹導讀:**滴滴資料通道引擎承載著全公司的資料同步,為下游實時和離線場景提供了必不可少的源資料。隨著任務量的不斷增加,資料通道的整體架構也隨之發生改變。本文介紹了滴滴資料通道的發展歷程,遇到的問題以及今后的規劃。 #1. 背景 資料,對于任何一家互聯網公司來說都是非常重要的資產,公司的大資料 ......

    uj5u.com 2020-09-10 06:11:05 more
  • 滴滴AI Labs斬獲國際機器翻譯大賽中譯英方向世界第三

    **桔妹導讀:**深耕人工智能領域,致力于探索AI讓出行更美好的滴滴AI Labs再次斬獲國際大獎,這次獲獎的專案是什么呢?一起來看看詳細報道吧! 近日,由國際計算語言學協會ACL(The Association for Computational Linguistics)舉辦的世界最具影響力的機器 ......

    uj5u.com 2020-09-10 06:11:29 more
  • MPP (Massively Parallel Processing)大規模并行處理

    1、什么是mpp? MPP (Massively Parallel Processing),即大規模并行處理,在資料庫非共享集群中,每個節點都有獨立的磁盤存盤系統和記憶體系統,業務資料根據資料庫模型和應用特點劃分到各個節點上,每臺資料節點通過專用網路或者商業通用網路互相連接,彼此協同計算,作為整體提供 ......

    uj5u.com 2020-09-10 06:11:41 more
  • 滴滴資料倉庫指標體系建設實踐

    **桔妹導讀:**指標體系是什么?如何使用OSM模型和AARRR模型搭建指標體系?如何統一流程、規范化、工具化管理指標體系?本文會對建設的方法論結合滴滴資料指標體系建設實踐進行解答分析。 #1. 什么是指標體系 ##1.1 指標體系定義 指標體系是將零散單點的具有相互聯系的指標,系統化的組織起來,通 ......

    uj5u.com 2020-09-10 06:12:52 more
  • 單表千萬行資料庫 LIKE 搜索優化手記

    我們經常在資料庫中使用 LIKE 運算子來完成對資料的模糊搜索,LIKE 運算子用于在 WHERE 子句中搜索列中的指定模式。 如果需要查找客戶表中所有姓氏是“張”的資料,可以使用下面的 SQL 陳述句: SELECT * FROM Customer WHERE Name LIKE '張%' 如果需要 ......

    uj5u.com 2020-09-10 06:13:25 more
  • 滴滴Ceph分布式存盤系統優化之鎖優化

    **桔妹導讀:**Ceph是國際知名的開源分布式存盤系統,在工業界和學術界都有著重要的影響。Ceph的架構和演算法設計發表在國際系統領域頂級會議OSDI、SOSP、SC等上。Ceph社區得到Red Hat、SUSE、Intel等大公司的大力支持。Ceph是國際云計算領域應用最廣泛的開源分布式存盤系統, ......

    uj5u.com 2020-09-10 06:14:51 more
  • es~通過ElasticsearchTemplate進行聚合~嵌套聚合

    之前寫過《es~通過ElasticsearchTemplate進行聚合操作》的文章,這一次主要寫一個嵌套的聚合,例如先對sex集合,再對desc聚合,最后再對age求和,共三層嵌套。 Aggregations的部分特性類似于SQL語言中的group by,avg,sum等函式,Aggregation ......

    uj5u.com 2020-09-10 06:14:59 more
  • 爬蟲日志監控 -- Elastc Stack(ELK)部署

    傻瓜式部署,只需替換IP與用戶 導讀: 現ELK四大組件分別為:Elasticsearch(核心)、logstash(處理)、filebeat(采集)、kibana(可視化) 下載均在https://www.elastic.co/cn/downloads/下tar包,各組件版本最好一致,配合fdm會 ......

    uj5u.com 2020-09-10 06:15:05 more
最新发布
  • day02-2-商鋪查詢快取

    功能02-商鋪查詢快取 3.商鋪詳情快取查詢 3.1什么是快取? 快取就是資料交換的緩沖區(稱作Cache),是存盤資料的臨時地方,一般讀寫性能較高。 快取的作用: 降低后端負載 提高讀寫效率,降低回應時間 快取的成本: 資料一致性成本 代碼維護成本 運維成本 3.2需求說明 如下,當我們點擊商店詳 ......

    uj5u.com 2023-04-20 08:33:24 more
  • MySQL中binlog備份腳本分享

    關于MySQL的二進制日志(binlog),我們都知道二進制日志(binlog)非常重要,尤其當你需要point to point災難恢復的時侯,所以我們要對其進行備份。關于二進制日志(binlog)的備份,可以基于flush logs方式先切換binlog,然后拷貝&壓縮到到遠程服務器或本地服務器 ......

    uj5u.com 2023-04-20 08:28:06 more
  • day02-短信登錄

    功能實作02 2.功能01-短信登錄 2.1基于Session實作登錄 2.1.1思路分析 2.1.2代碼實作 2.1.2.1發送短信驗證碼 發送短信驗證碼: 發送驗證碼的介面為:http://127.0.0.1:8080/api/user/code?phone=xxxxx<手機號> 請求方式:PO ......

    uj5u.com 2023-04-20 08:27:27 more
  • 快取與資料庫雙寫一致性幾種策略分析

    本文將對幾種快取與資料庫保證資料一致性的使用方式進行分析。為保證高并發性能,以下分析場景不考慮執行的原子性及加鎖等強一致性要求的場景,僅追求最終一致性。 ......

    uj5u.com 2023-04-20 08:26:48 more
  • sql陳述句優化

    問題查找及措施 問題查找 需要找到具體的代碼,對其進行一對一優化,而非一直把關注點放在服務器和sql平臺 降低簡化每個事務中處理的問題,盡量不要讓一個事務拖太長的時間 例如檔案上傳時,應將檔案上傳這一步放在事務外面 微軟建議 4.啟動sql定時執行計劃 怎么啟動sqlserver代理服務-百度經驗 ......

    uj5u.com 2023-04-20 08:26:35 more
  • 云時代,MySQL到ClickHouse資料同步產品對比推薦

    ClickHouse 在執行分析查詢時的速度優勢很好的彌補了MySQL的不足,但是對于很多開發者和DBA來說,如何將MySQL穩定、高效、簡單的同步到 ClickHouse 卻很困難。本文對比了 NineData、MaterializeMySQL(ClickHouse自帶)、Bifrost 三款產品... ......

    uj5u.com 2023-04-20 08:26:29 more
  • sql陳述句優化

    問題查找及措施 問題查找 需要找到具體的代碼,對其進行一對一優化,而非一直把關注點放在服務器和sql平臺 降低簡化每個事務中處理的問題,盡量不要讓一個事務拖太長的時間 例如檔案上傳時,應將檔案上傳這一步放在事務外面 微軟建議 4.啟動sql定時執行計劃 怎么啟動sqlserver代理服務-百度經驗 ......

    uj5u.com 2023-04-20 08:25:13 more
  • Redis 報”OutOfDirectMemoryError“(堆外記憶體溢位)

    Redis 報錯“OutOfDirectMemoryError(堆外記憶體溢位) ”問題如下: 一、報錯資訊: 使用 Redis 的業務介面 ,產生 OutOfDirectMemoryError(堆外記憶體溢位),如圖: 格式化后的報錯資訊: { "timestamp": "2023-04-17 22: ......

    uj5u.com 2023-04-20 08:24:54 more
  • day02-2-商鋪查詢快取

    功能02-商鋪查詢快取 3.商鋪詳情快取查詢 3.1什么是快取? 快取就是資料交換的緩沖區(稱作Cache),是存盤資料的臨時地方,一般讀寫性能較高。 快取的作用: 降低后端負載 提高讀寫效率,降低回應時間 快取的成本: 資料一致性成本 代碼維護成本 運維成本 3.2需求說明 如下,當我們點擊商店詳 ......

    uj5u.com 2023-04-20 08:24:03 more
  • day02-短信登錄

    功能實作02 2.功能01-短信登錄 2.1基于Session實作登錄 2.1.1思路分析 2.1.2代碼實作 2.1.2.1發送短信驗證碼 發送短信驗證碼: 發送驗證碼的介面為:http://127.0.0.1:8080/api/user/code?phone=xxxxx<手機號> 請求方式:PO ......

    uj5u.com 2023-04-20 08:23:11 more