1. 生產者原理
這里分析的是Apache官方客戶端代碼,原始碼位于:https://github.com/apache/kafka/tree/trunk/clients
從git clone 后 需要手動切換到對應分支,
本文所有案例工程在文末
1.1 生產者消費發送流程
訊息發送的整體流程如下:生產者主要由兩個執行緒協調運行,這兩條執行緒分別為main執行緒和sender執行緒(發送執行緒)

我們可以看原始碼跟蹤一下:從Producer 入口進入即可,
Producer<String, String> producer = new KafkaProducer<String, String>(pros);
進入構造方法,我們可以發現在初始化的時候,創建了一個Sender物件,并且啟動了一個IO執行緒,(我這邊是在第188行,如果找不到的話直接搜索即可)
this.sender = this.newSender(logContext, kafkaClient, this.metadata);
String ioThreadName = "kafka-producer-network-thread | " + this.clientId;
this.ioThread = new KafkaThread(ioThreadName, this.sender, true);
this.ioThread.start();
1.1.1 攔截器
接下來是攔截器的執行,在 producer.send 方法中:
public Future<RecordMetadata> send(ProducerRecord<K, V> record, Callback callback) {
ProducerRecord<K, V> interceptedRecord = this.interceptors.onSend(record);
return this.doSend(interceptedRecord, callback);
}
攔截器的作用是實作訊息的定制化(類似于:Spring Interceptor、Mybatis 插件、Quartz的監聽器等)
那這個攔截器是在哪里定義的呢?我們可以自己實作一下:
// 添加攔截器
List<String> interceptors = new ArrayList<>();
interceptors.add("com.demo.interceptor.ChargingInterceptor");
props.put(ProducerConfig.INTERCEPTOR_CLASSES_CONFIG, interceptors);
可以在生產者的屬性中指定多個攔截器,形成攔截器鏈,
舉個栗子,假設發訊息的時候需要扣錢,發一條訊息一分錢,就可以使用攔截器實作,
public class ChargingInterceptor implements ProducerInterceptor<String, String> {
// 發送訊息的時候觸發
@Override
public ProducerRecord<String, String> onSend(ProducerRecord<String, String> record) {
System.out.println("我要開始扣錢啦~");
return record;
}
// 收到服務端的ACK的時候觸發
@Override
public void onAcknowledgement(RecordMetadata metadata, Exception exception) {
System.out.println("訊息被服務端接收啦");
}
@Override
public void close() {
System.out.println("生產者關閉了");
}
// 用鍵值對配置的時候觸發
@Override
public void configure(Map<String, ?> configs) {
System.out.println("configure...");
}
}
我們只需要將該攔截器配置進引數中,當生產者發送訊息的時候就會觸發對應的方法,

1.1.2 序列化
呼叫send方法后,第二步是利用指定的工具對key和value進行序列化:(我這邊是363行)
byte[] serializedKey;
try {
serializedKey = this.keySerializer.serialize(record.topic(), record.headers(), record.key());
} catch (ClassCastException var21) {
throw new SerializationException("Can't convert key of class " + record.key().getClass().getName() + " to class " + this.producerConfig.getClass("key.serializer").getName() + " specified in key.serializer", var21);
}
byte[] serializedValue;
try {
serializedValue = this.valueSerializer.serialize(record.topic(), record.headers(), record.value());
} catch (ClassCastException var20) {
throw new SerializationException("Can't convert value of class " + record.value().getClass().getName() + " to class " + this.producerConfig.getClass("value.serializer").getName() + " specified in value.serializer", var20);
}
Serializer.java 針對不同的資料型別自帶了相應的序列化工具:

除了自帶的序列化工具之外,可以使用如JSON,Protobuf等,或者使用自定義型別的序列化器來實作,實作serialzer介面介面,
我們可以自定義一個序列化介面,然后發送訊息的時候添加相關引數即可,
props.put("value.serializer", "com.demo.serializer.ProtobufSerializer");
1.1.3 路由指定(磁區器)
看過序列化之后,就來到了路由指定,(377)
int partition = this.partition(record, serializedKey, serializedValue, cluster);
一條訊息會發送到那個partition呢?他回傳的是一個磁區的編號,從0開始,
首先我們將磁區分為四種情況:
- 指定了partition;
- 沒有指定partition,自定義了磁區器;
- 沒有指定partition,沒有自定義磁區器,但是key不為空;
- 沒有指定partition,沒有自定義磁區器,key為空;
partition數量可以自行去組態檔中修改,
第一種情況:
指定partition的情況下,直接將指定的值直接作為partition值,
for (int i = 0; i < 10; i++) {
ProducerRecord<String, Integer> producerRecord = new ProducerRecord<String, Integer>(topic, i, null, i);
RecordMetadata metadata = producer.send(producerRecord).get();
System.out.println("Sent to partition: " + metadata.partition() + ", offset: " + metadata.offset());
}

第二種情況:
自定義磁區器,將使用自定義的磁區器演算法選擇磁區,比如我們自定一個磁區器,然后指定即可,
public class SimplePartitioner implements Partitioner {
public SimplePartitioner() {
}
@Override
public void configure(Map<String, ?> configs) {
}
@Override
public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) {
String k = (String) key;
System.out.println(k);
if (Integer.parseInt(k) % 2 == 0) {
return 0;
} else {
return 1;
}
}
@Override
public void close() {
}
}
指定自定義磁區器:
props.put("partitioner.class", "com.demo.partition.SimplePartitioner");

第三種情況:
沒有指定partition值但是有key的情況下,使用默認磁區DefaultPartitioner,將key的hash值與topic的partition數進行取余得到partition值;
return Utils.toPositive(Utils.murmur2(keyBytes)) % numPartitions;

第四種情況:
既沒有partition值但有沒有key值的情況下,第一次呼叫時隨機生成一個整數(后面每次呼叫在這個數上自增),將這個值與topic可用的partition總數取余得到partition值,也就是常說的輪詢演算法,
public int partition(String topic, Cluster cluster) {
Integer part = (Integer)this.indexCache.get(topic);
return part == null ? this.nextPartition(topic, cluster, -1) : part;
}
1.1.4 訊息累加器
選擇磁區以后并沒有直接發送訊息,而是把訊息放入了訊息累加器(390);
RecordAppendResult result = this.accumulator.append(tp, timestamp, serializedKey, serializedValue, headers, interceptCallback, remainingWaitMs, true, nowMs);
RecordAccumulator 本質上是一個ConcurentMap;
private final ConcurrentMap<TopicPartition, Deque<ProducerBatch>> batches;
一個partition一個batch,Batch滿了之后會喚醒Sender執行緒,發送訊息,(408)
if (result.batchIsFull || result.newBatchCreated) {
this.log.trace("Waking up the sender since topic {} partition {} is either full or getting a new batch", record.topic(), partition);
this.sender.wakeup();
}
1.1.5 總結
我們可以在攔截器里面自定義訊息處理邏輯,也可以選擇自己喜歡的序列化工具,還可以自由選擇磁區,
1.2 資料可靠性保證ACK
1.2.1 服務端相應策略
生產者的訊息是不是發出去就完事了?如果說網路出問題了,或者說Kafka服務端接受的時候出了問題,這個訊息發送失敗了,生產者是不知道的,
所以,Kafka服務端應該要有一種回應客戶端的方式,只有在服務端確認以后,生產者才會發下一輪的訊息,否則重新發送資料,
那么服務端什么時候才算接受成功呢?因為訊息是存盤在不同的partition里面的,所有是寫入到partition之后才會相應生產者,

當然,單個partition(leader)寫入成功,還是不夠可靠,如果有多個副本,follower也要寫入成功才可以,
服務端發送ACK給生產者總體上有兩種思路:
第一種是需要有半數以上的follower節點同步完成,這樣的話客戶端等待的時間就短一些,延遲低,(所以我們通常部署節點的數量都是奇數,如果是偶數,兩邊一樣就很尷尬,)
第二種是需要所有的follower節點全部完成同步,才發送ACK給客戶端,延遲來說相對高一些,但是節點掛掉的可能性比較小,因為所有的節點資料都是完整的,
Kafka會選擇那種方案呢?
Kafka選擇了第二種,部署同樣機器數量的情況下,第二種方案可靠性更高,同時網路延遲對Kafka的影響不是很大,
1.2.2 ISR
如果直接采用第二種思路,不考慮網路延遲,有沒有別的問題呢?
假設leader收到資料,所有follower都開始同步資料,但是有一個follower出了問題,沒有辦法從leader同步資料,按照這個規則,leader就要一直等待,無法發送ACK…
從概率的角度來說,這種問題肯定是會出現的,就是某個follower出問題了,怎么解決這種問題呢?
所以我們的規則就不能那么粗暴了,不能因為一個follower的問題導致無法發送ACK;我們把規則改一下,不是所有的follower都有權力讓leader等待,而是只有那些正常作業的follower同步資料的時候leader才會等待,
我們應該把那些正常和leader保持同步的replica維護起來,放到一個動態list里面,這個就叫做in-sync replica set(ISR),現在只要ISR里面的follower同步完資料之后,leader就給客戶端發送ACK,
如果一個follower長時間不同步資料,就將其從ISR中移除,那么到底多久沒有同步資料才會被剔除呢?這個是由引數replica.lag.time.max.ms決定,默認是30秒,當然了如果follower活過來了,則還能進入ISR中,
如果leader掛了,ISR會重新選擇leader,這部分下文再說,
1.2.3 ACK應答機制
Kafka為客戶端提供了三種可靠性機制,用戶根據對可靠性和延遲的要求自行權衡,選擇相應的配置,
引數配置如下:
pros.put("acks", "1");
舉例:topic的partition0有三個副本,
-
ack = 0
producer不等待broker的ACK,這一操作提供了一個最低的延遲,broker一接收到還沒有寫入磁盤就已經回傳,當broker故障時有可能丟失資料,

-
ack = 1 (默認)
producer 等到 broker 的ack,partition 的 leader 落盤成功后回傳ack,如果在follower同步成功之前leader故障,那么將會丟失資料,

-
ack = -1 (all)
producer 等待broker的ack,partition 的leader 和follower 全部落盤成功后才回傳ack,
這種方案是完美的嗎?會出現問題嗎?
如果在follower同步完成后,broker發送ack之前,leader發生故障,沒有給生產者發送ACK,那么會遭成資料重復,
在這種情況下,把reties 設定成0(不重發),才不會重復,

三種機制,性能依次遞減(producer吞吐量降低),資料健壯性則依次遞增,我們可以根據業務場景選擇合適的引數,
2. Broker 存盤原理
2.1 檔案的存盤結構
路徑設定:config/server.propertise中的logs.dir配置
默認/tmp/kafka/logs
2.1.1 partition 磁區
為了實作橫向擴展,把不同的資料存放在不同的Broker上,同時降低單臺服務器的訪問壓力,我們把一個topic中的資料分割成多個partition,
一個partition中訊息是有序的,順序寫入,但是全域不一定有序,

在服務器上,每個partition都有一個物理目錄,topic名字后面的資料標號則代表磁區,

2.1.2 replica 副本
為了提高磁區的可靠性,Kafka有設計了副本機制,
創建Topic的時候,通過指定 replication-factor 確定副本的數,
注意:副本數必須小于等于節點數,而不能大于Broker的數量,否則會保存,
錯誤示范:
./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 4 --partitions 1 --topic overrep
這樣就可以保證,絕對不會有一個磁區的兩個副本分布在同一個節點上,不然副本機制也失去了備份的意義了,

這些所有的副本分為兩種角色,leader對外提供讀寫服務,follower唯一的任務就是從leader異步拉取資料,
為什么只有leader提供讀寫服務呢?而不是像mysql一樣讀寫分離?
答:這個是設計思想的不同,讀寫都發生在leader節點上,就不存在讀寫分離帶來的一致性問題,這個叫做單調讀一致性,
2.1.3 leader在哪里?
問題來了,如果磁區有多個副本,哪一個節點上的副本是leader呢?
./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 3 --partitions 3 --topic a3part3rep
怎么查看所有的副本中誰是leader?(需要集群環境)
./kafka-topics.sh --topic a3part3rep --describe --zookeeper localhost:2181

解釋:
這個topic有三個磁區三個副本,
第一個磁區的3個副本編號 1 , 2 ,3 (注意副本的編號是從1開始的),同步中的也是 1 ,2 ,3 ,第一個副本是leader,

假設 topic 有 4個磁區2個副本呢?
./kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 2 --partitions 4 --topic a4part2rep
查看:
./kafka-topics.sh --topic a4part2rep --describe --zookeeper localhost:2181

解釋:
這個磁區有4個磁區兩個2副本;第一個磁區的2個副本編號2,3 ;同步中的也是2,3,第三個副本是leader,

為什么第一個磁區兩個副本選擇在2,3broker;第二個磁區兩個副本選擇1,3;第三個磁區的副本選擇1,2broker呢?
2.1.4 副本在Broker中的分布,
副本在Broker的分布有什么規則嗎?
a4part2rep這個topic,4個磁區2個副本,一共8個副本,怎么分布到3臺機器?
結果剛剛我們已經看到過了,大家也可以自行前往tmp/kafka-logs中查看對應的資料資訊,
實際上,這種分配策略是由AdminUtils.scala的assignReplicasToBrokers函式決定的,
規則如下:
-
first of all
副本因子不能大于Broker的個數
-
第一個磁區(編號為0的磁區)的第一個副本位置是隨機從brokerList選擇的;
-
其他磁區的第一個副本放置位置相對于第0個磁區依次往后移,
也就是說:如果我們有5個Broker,5個磁區,假設第一個磁區的第一個副本放在第四個Broker上,那么第2個磁區的第一個副本將會放在第5個Broker上;第三個磁區的第一個副本將會放在第一個Broker上;第四個磁區的第一個副本將會放在第二個Broker上,以此類推;
-
每個磁區剩余副本數相對于第一個副本位置其實是由nextReplicaShift決定的,這個數也是隨機生成的,

為什么這么設計呢?
在每個磁區的第一個副本錯開之后,一般第一個磁區的第一個副本(按Broker編號排序)都是leader,leader是錯開的,以至于broker掛了之后影響太大,
bin目錄下的kafka-reassign-partition.sh可以根據Broker數量變化情況重新分配磁區,
一個磁區是不是只有一個檔案呢?也就是說,訊息日志檔案會不會無限變大?
2.1.5 segment
為了防止log不斷追加導致檔案過大,導致檢索訊息效率變低,一個partition又被劃分成多個segment來組織資料,
在磁盤上,每個segment由一個log檔案和2個index檔案組成,

這三個檔案是成套出現的,(其他檔案先忽略)
leader-epoch-checkpoint 中保存了每一任leader開始寫入訊息時的offset,
log日志檔案
在一個新的segment檔案里面,日志時被追加寫入的,如果滿足一定條件,就會切分日志檔案,產生一個新的segment,什么時候會觸segment的切分呢?
第一種是根據日志檔案大小,當一個segment寫滿以后,會創建一個新的segment,用最新的offset作為名稱,這個例子可以通過往一個topic發送大量訊息產生,
segment的默認大小是1G,通過以下引數控制:
log.segment.bytes
第二種是根據訊息的最大時間戳和當前系統時間戳的差值,
有一個默認的引數:168小時 (一周)
log.roll.hours=168
意味著,如果服務器上次寫入時間是一周之前,舊的segment就不寫了,重新創建一個segment;
還可以從更加精細的時間單位進行控制,如果配置了毫秒級別的日志切分時間間隔會優先使用這個單位,否則就用小時的,
log.roll.ms
第三種情況就是offset索引檔案或者timestamp索引檔案達到了一定的大小,默認是10M,如果要減少日志檔案的切分,可以把這個值調大一點,
log.index.size.max.bytes
意思就是:索引檔案寫滿了,資料檔案也要跟著拆分,不然這一套東西對不上,
2.1.6 索引
由于一個segment的檔案里面可能存放很多訊息,如果要根據offset獲取訊息,必須要有一種快速檢索訊息的機制,這個就是索引,
在Kafka中設計了兩種索引,
偏移量索引檔案記錄的是offset和訊息物理地址(在log檔案中的位置)的映射關系,時間戳索引檔案記錄是時間戳和offset的關系,
當然,內容是二進制的 檔案,不能以純文本的形式查看,bin目錄下有dumplog工具,
./kafka-dump-log.sh --files /tmp/kafka-logs/mytopic-0/000000.index|head -n 10

注意Kafka的索引并不是每一條訊息都會建立索引,二是一種稀疏索引,
稀疏索引的稀疏程度是根據訊息的大小來控制的,默認是4KB;
log.index.interval.bytes=4096
只要寫入的資訊超過了4KB,偏移量索引檔案和時間戳索引檔案就會增加一條記錄,
這個值設定的越小,索引則越密集;值越大則索引越稀疏,
相對來說越稠密的索引資料檢索更快,但是會消耗更多的空間存盤,
稀疏消費的空間少,但是插入和洗掉時開銷比較大,
第二種所以是時間戳索引,
為什么會有時間戳索引檔案呢?光有offset索引還不夠嗎?會根據時間戳來查找訊息嗎?
首先訊息是必須要記錄時間戳的,客戶端封裝的ProducerRecord就有timestamp屬性,
為什么需要呢?
- 如果要基于時間切分日志檔案,必須要有時間戳;
- 如果要基于時間清理訊息,必須要有時間戳,
既然都記錄時間戳了,那干脆就可以直接設計一個時間戳索引,可以根據時間戳查詢,
注意創建時間戳有兩種:一種是創建訊息的時間戳,一種是消費在Broker追加寫入的時間,我們改用那個引數?這個也可以通過引數控制:
log.message.timestamp.type=CreateTime
默認是創建時間,如果要改成日志追加時間,則修改為LogAppendTime;
快速檢索
Kafka是如何基于索引快速檢索呢?比如我要檢索偏移量是10002673的訊息,
- 消費的時候是能夠確定磁區的,所以第一步是找到在那個segment中,segment檔案使用base offset命名的,所以可以用二分法很快確定(找到名字不小于10002673的segment);
- 這個segment有對應的索引檔案,他們是成套出現的,所以現在要在索引檔案中根據offset找position,
- 得到position后,到對應的log檔案開始查找offset,和訊息的offset進行比較,直到找到訊息,
為什么不用B+樹?
因為Kafka是寫多查少,如果用B+樹,首先會出現大量的B+樹,大量的插入會非常消耗性能,
2.2 訊息保留(清理)機制
都知道Kafka是將資料保存在磁盤中的,那么很多的舊資料我們該怎么辦?
2.2.1 開關與策略
清理策略默認是開啟的:
log.cleaner.enable=true
kafka里面提供了兩種方式,一種是直接洗掉,一種是對日志進行壓縮,默認是直接洗掉,
log.cleanup.policy=delete
2.2.2 洗掉策略
如果是洗掉,什么時候洗掉呢?日志洗掉是通過定時任務實作的,默認5分鐘執行一次,
log.retention.check.interval.ms=300000
那么從哪里洗掉呢?當然是從老資料開始,那么什么才是老資料呢?
通過以下引數進行控制:
log.retention.hours
默認值是168小時(一周),也就是時間戳超過一周的資料才會洗掉,
Kafka也提供了另外粒度更細的配置:分鐘和毫秒,
這里還有一種情況,假設Kafka產生訊息的速度是不均勻的,有的時候一周幾百萬條資料,有的時候一周幾千條資料,那這個按照時間洗掉就不太合理了,所以第二種洗掉策略就是根據日志檔案大小洗掉,先刪舊的資料,一直刪到不超過這個大小為止,
log.retention.bytes
默認值是-1,代表不限制大小,想寫多少就寫多少,它指的是所有檔案的大小,我們也可以對單個segment檔案大小進行限制,
log.segment.bytes
默認是1G.
2.2.3 壓縮策略
第二種策略是不洗掉,對日志資料進行壓縮,
問題:如果同一個key重復寫入多次,會存盤多次還是更新?
比如用來存盤唯一的這個特殊topic:_consumer_offsets,存盤的是消費者ID和partition的offset關系,消費者不斷消費資訊commit的時候是直接更新原來的offset,還是不斷的寫入呢?
答案是肯定存盤多次,不然我們怎么實作順序寫呢,
當有了這些key相同value不同的訊息的時候,存盤空間就被浪費了,壓縮就是將相同的Key合并為最后一個value.
2.3 高可用架構
2.3.1 Controller 選舉
當創建添加一個的磁區或者磁區增加了副本的時候,都要從所有副本中選舉一個新的leader出來,
那么我們怎么進行選舉呢?
通過ZK實作嗎?通過ZK的watch機制來實作嗎?
這種方法雖然簡單,但是存在一定的弊端,如果磁區和副本數量過多,所有的副本都直接進行選舉的話,一旦某個出現節點的增,就會遭成大量的watch事件被觸發,ZK的負載就會過重,
Kafka早期版本就是這么實作的,后來換了一種實作方式,
不是所有的replica都參與leader選舉,而是由其中的一個Broker統一來指揮,這個Broker的角色就叫做Controller,類似redis集群中的哨兵機制,
所有的Broker會嘗試在ZK中創建臨時節點,只有一個能創建成功(先到先得),
如果Controller掛掉了或者網路出現了問題,ZK上的臨時節點會訊息,其他的Broker通過watch監聽到了Controller下線的訊息后,開始競選新的controller,方法跟之前還是一樣的,誰先在ZK里面寫入一個controller節點,誰就成為新的controller,
一個節點成為controller,它肩上的責任也比別人重了幾分,
- 監聽Broker變化
- 監聽Topic變化
- 監聽partition變化
- 獲取和管理Broker,Topic,Partition的資訊
- 管理Partition的主從訊息
2.3.2 磁區副本Leader選舉
Controller確定以后,就可以開始做磁區選主的事情了,顯然每個replica都想推薦自己,但是所有的replica都有競選資格嗎?
并不是,這里有幾個概念,
一個磁區所有的副本,叫做Assigned-Replicas(AR),
這些所有的副本中,跟leader資料保持一定程度同步的,叫做In-Sync Replicas(ISR),
跟leader同步滯后過多的副本,叫做Out-Sync-Replicas(OSR),
AR = ISR + OSR ,正常情況下OSR是空的,大家都同步,AR = ISR,
誰能夠參加選舉呢?肯定不是AR ,也不是OSR,而是ISR,而且這個ISR不是固定不動的,還是一個動態串列,
前面我們說過,如果同步延遲超過30秒,就剔除ISR,進入OSR,如果趕上了就加入ISR,
默認情況下,當leader副本發生故障時,只有在ISR集合中的副本才有資格被選舉成新的leader,
如果ISR為空呢?在這種情況下,可以讓ISR之外的副本參與選舉,允許ISR之外的副本參與選舉,叫做unclean leader election.
unable.leader.election.enable=false;
把這個引數改成true(一般情況不建議開啟,會造成資料丟失,)
選舉規則
分布式系統中常見的選舉協議有哪些?
ZAB(ZK),Raft(Redis)思想歸納起來都是:先到先得,少數服從多數,
但是Kafka沒有用到這些辦法,而是用了一種自己實作的演算法,
為什么呢?比如zab協議,可能會出現腦裂(節點不能互通的時候,出現多個leader,)驚群效應(大量watch事件被觸發)
Kafka的選舉實作類似與微軟的PacificA演算法,
在這種演算法中,默認是讓ISR中第一個replica變成leader,比如ISR是1,5,8.優先讓1成為leader,
2.3.3 主從同步
leader確定以后,客戶端的讀寫操作只能操作leader節點,follower需要向leader同步資料,
不同的replica的offset是不一樣的,到底怎么同步呢?
又要看幾個概念了,,,,

LEO(Log End Offset):下一條等待寫入的訊息的offset(最新的offset+1);途中分別是9.8.6;
HW(Hign Watermark) :ISR中最小的LEO.,leader會改管理所有ISR中最小的LEO作為HW,目前是6.
**consumer最多只能消費到HW之前的位置(消費到offset5的訊息),**也就是說:其他的副本沒有同步過去的訊息,是不能被消費的,
為什么要這么設計呢?如果在同步成功之前就被消費了,consumer group的offset會偏大,如果leader崩潰,中間會缺失訊息,
有了這兩個offset之后,我們再來看看訊息怎么同步,
follower1同步了1條訊息,follower2同步了兩條訊息,此時HW推進了2,變成了8,

follower1同步了0條訊息,follower2同步了1條訊息,此時HW推進了1,變成了9,LEO和HW重疊了,所有的訊息都可以消費了,

這里我們關注以下,從節點怎么跟主節點保持同步?
- follower節點會向leader發送要給fetch請求,leader向follower發送資料后,需要更新follower的LEO,
- follower接收到資料回應后,依次寫入訊息并更新LEO,
- leader更新HW(ISR最小的LEO)
這中獨特的ISR復制,可以在保障資料一致性情況下又可以提供高吞吐量,
2.3.4 replica 故障處理
follower故障
首先follower發生故障,會被踢出ISR,
follower恢復之后,從哪里開始同步資料呢?假設第一個replica宕機,(中間這個)

恢復以后,首先根據之前記錄的HW(6),把高于HW的訊息截掉(6,7).然后向leader同步訊息,追上leader之后,重新加入ISR,
leader故障
假設圖中leader發生故障,
首先選一個leader,因為replica1優先(中間這個),所以它成為leader,
為了保證資料一致性,其他的follower需要把高于HW的訊息截取掉(這里沒有訊息需要截取,)
然后replica2開始同步資料,
注意:這種機制只能保證副本之間的資料一致性,并不能保證資料不丟失或者不重復,
3. 消費者原理
3.1 offset 的維護
我們首先研究以下消費者怎么消費,
3.1.1 offset的存盤
我們知道在partition中,訊息是不會洗掉的,所以在可以追加寫入,寫入的訊息是連續存在的,
這種特性決定了Kafka是可以消費歷史訊息的而且按照訊息的順序消費指定訊息,而不是只能消費對頭的訊息,
正常情況下,我們希望消費沒有被消費過的資料,而且是從最先發送的開始消費,(這樣才是有序和公平的);
對于一個partition,消費者組怎么才能做到接著上次消費的位置(offset)繼續消費呢?我們肯定需要把這個對應關系保存起來,下次消費的時候直接查找就可以了,
(對應上文的索引檔案)

這個對應關系到底保存在哪里呢?首先肯定不是保存在消費者這端的,為什么?因為所有的消費者都可以使用這個consumer group id,放在本地是做不到統一維護的,肯定需要放到服務端,
Kafka早期的版本把消費者和partition的offset直接維護在ZK中,但是讀寫的性能消耗太大了,后來就放在一個特殊的topic中,名字叫_consumer_offsets,默認有50個磁區(offsets.topic.num.partitions默認是50),每個磁區默認一個replication,
這樣的一個特殊的Topic怎么存盤消費者組對于磁區的偏移量呢?
Topic里面是可以存放物件型別的value的(經過序列化和反序列化),這個topic里面主要存盤兩種物件:
-
GroupMetadata
保存了消費者組中各個消費者的資訊(每個消費者有編號)
-
OffsetAndMetadata
保存了消費者組和各個partition的offset位移資訊元資料,
大致結構如下:

我們怎么知道offset會放在那個磁區呢?
System.out.println(Math.abs("test".hashCode() % 50));
3.1.2 如果找不到offset
當然,這個是Broker有記錄offset的情況,如果說新增了一個消費者組去消費一個topic的某個partition,沒有offset的記錄,這個時候我們應該從哪里開始消費?
什么情況下會找不到offset?就是沒有消費過,沒有把當前的offset上報給Broker,
消費者的代碼中有一個引數用來控制如果找不到偏移量的時候從哪里開始消費,
auto.offset.reset
-
latest 默認值
也就是從最新的訊息開始消費(最后發送的訊息),歷史訊息是不能消費的,
-
earliest
代表從最早的訊息開始消費(最先發送的訊息),可以消費到歷史訊息,
-
none
如果消費組在服務端找不到offset會報錯,
3.1.3 offset的更新
前面我們講了,消費者組的offset是保存在broker的,但是是由消費者上報給broker的,并不是消費者組消費了訊息,offset就會更新,消費者必須要有一個commit(提交)的動作,就跟RabbitMQ中消費者的ACK一樣,
消費者可以自動提交或者手動提交,通過以下引數控制:
enable.auto.commit
默認為true,true代表消費者消費訊息以后自動提交,此時Broker會更新消費者組的offset,
另外還可以通過配置引數來控制自動提交的頻率:
auto.commit.interval.ms
默認是5秒,
如果我們要在消費完訊息做完業務邏輯處理之后才commit,就要把這個值改成false,
如果是false,消費者就必須要呼叫一個方法讓Broker更新offset,
有兩種方式:
- consumer.conmmitSync()的手動同步提交
- consumer.conmmitAsync()的手動異步提交
演示代碼如下:
// 手動提交
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
System.out.printf("offset = %d ,key =%s, value= %s, partition= %s%n", record.offset(), record.key(), record.value(), record.partition());
buffer.add(record);
}
if (buffer.size() >= minBatchSize) {
// 同步提交
consumer.commitSync();
buffer.clear();
}
}
如果不提交或者提交失敗,Broker的offset不會更新,消費者下次消費的時候會消費到重復的訊息,
3.2 消費者消費策略(消費者與磁區關系)
3.2.1 消費策略
前面我們講過,一個消費者里面的一個消費,只能消費Topic的一個磁區,
如果磁區數量跟消費者的數量一樣,那就一人消費一個,如果是消費者比磁區多,或者消費者比磁區小,這個時候消費者跟磁區的關系是怎么樣的呢?
例如:2個消費者消費5個磁區,怎么分配呢?
我們首先創建一個5個磁區的topic,
sh bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 5 --topic ass5part
啟動兩個消費者消費(同一個消費者,目標是消費同一個partition,不同的clientid)
// 兩個消費者消費5個磁區
KafkaConsumer<String, String> consumer1 = new KafkaConsumer<String, String>(props);
KafkaConsumer<String, String> consumer2 = new KafkaConsumer<String, String>(props);
// 訂閱佇列
consumer1.subscribe(Arrays.asList("ass5part"));
consumer2.subscribe(Arrays.asList("ass5part"));
給5個磁區分別發送一條訊息:
producer.send(new ProducerRecord<String, String>("ass5part", 0, "0", "0"));
producer.send(new ProducerRecord<String, String>("ass5part", 1, "1", "1"));
producer.send(new ProducerRecord<String, String>("ass5part", 2, "2", "2"));
producer.send(new ProducerRecord<String, String>("ass5part", 3, "3", "3"));
producer.send(new ProducerRecord<String, String>("ass5part", 4, "4", "4"));
結果(列印順序不一定)

他是按照范圍連續分配的,你一部分我一部分,

他實際上采用了默認的策略:RangeAssignor,

我們也可以通過配置引數使用其他的消費策略,
props.put("partition.assignment.strategy","org.apache.kafka.clients.consumer.RoundRobinAssignor");
另外兩種策略的查詢及如果如下:
- 輪詢策略

-
StickyAssignor:這種策略比較復雜,但是相對來說均勻一點(每次的結果可能不一樣),
原則:
- 磁區的分配盡可能均勻
- 磁區的分配盡可能和上次分配保持相同
consumer可以指定topic的某個磁區消費嗎?比如我就喜歡坐在講臺旁的寶座,我可以去做嗎?
這個時候我們需要使用打assgin而不是subscribe介面,subscribe會自動分配磁區,而assign是由我們自己指定磁區消費,相當于comsumer group id失效了,
// 訂閱topic,消費指定parptition
TopicPartition tp = new TopicPartition("ass5part", 0);
之前已經說過,在第一次消費的時候,一個組的消費者和磁區的消費就已經確定了,如果分配策略沒動,關系是不會改變的,那什么時候才會重新分配呢?
3.2.2 rebalance 磁區重分配
有兩種情況需要重新分配磁區和消費者的關系:
- 消費者組的消費者數量發生了變化,比如新增了消費者或者消費者關閉連接,
- topic的磁區數量發生了變化,新增或者減少了,
為了讓磁區分配盡量均衡,這個時候會觸發rebalance機制,
大致分為以下幾步:

- 首先找一個話事人出來,他起到監督和保證公平的作用,每個Broker上都有一個用來管理offset,消費者組的實體,叫做GroupCoordinator,第一步就是要從所有GroupCoordinator中找一個話事人出來,
- 第二步:清點人數,所有的消費者連接到GroupCoordiantor報數,這個叫做join group請求,
- 第三步:選組長,GroupCoordinator從所有消費者里面選一個leader,這個消費者會根據消費者的情況和設定的策略,確定一個方案,leader把方案上報給GroupCoordinator,然后GroupCoordinatar會通知所有消費者,
4. Kafka為什么這么快?
總結起來,主要是四點:磁盤順序io,索引機制,批量操作和壓緊,零拷貝,
-
磁盤IO
隨機I/O讀寫的資料在磁盤上分散的,尋址會狠耗時,
順序I/O讀寫的資料在磁盤上是集中的,不需要尋址的程序,
所以順序IO是被隨機IO快的多的,
-
索引在上文已經講過了,
-
批量讀寫
Kafka將所有的訊息變成一個批量檔案,減少網路IO損耗,
-
零拷貝
我覺得從下面那個圖應該可以看出零拷貝大大提升了檔案傳輸的性能,
傳統IO模型:

零拷貝模式:

5. 確保Kafka訊息不丟失的配置
- producer端使用帶有回呼方法的send方法,根據回呼,一旦出現訊息提交失敗的情況則進行處理,
- 設定ack = all,所有Broker接受到訊息后,才算提交完成,
- 設定retries為一個較大的值,使其能自動重試發生訊息,避免訊息丟失,
- 需要三個以上的副本
- 設定 unclean.leader.election.enable=false;不從ISR中選舉leader,
- 確保訊息消費完在提交,最好將自動提交關閉,改成我們手動提交,
6. 專案地址
kafka-demo
轉載請註明出處,本文鏈接:https://www.uj5u.com/qita/382823.html
標籤:其他
上一篇:C語言 圣誕樹(程式員的浪漫)
