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.size 和 linger.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
標籤:其他
