Kafka
Kafka 基礎概念
Kafka 已被多家不同類型的公司作為多種類型的數據管道和消息系統使用。 行為流數據是幾乎所有站點在對其網站使用情況做報表時都要用到的數據中最常規的部分。
- 包括頁面訪問量 PV(Page View)、頁面曝光(Expose)、頁面點擊(Click) 等行為事件;
- 實時計算中的 Kafka Source,Dataflow Pipeline; 業務的消息系統,通過發布訂閱消息解耦多組微服務,消除峰值;
Kafka 是由 LinkedIn 開發並開源的分佈式消息系統,因其分佈式及高吞吐率而被廣泛使用,現已與Cloudera Hadoop,Apache Storm,Apache Spark集成。
Kafka簡介
Kafka 是一種分佈式的,基於發布/訂閱的消息系統。主要設計目標如下:
- 以時間複雜度為 O(1) 的方式提供消息持久化能力,即使對 TB 級以上數據也能保證常數時間複雜度的訪問性能;
- 高吞吐率。即使在非常廉價的商用機器上也能做到單機支持每秒 100K 條以上消息的傳輸;
- 支持 Kafka Server 間的消息分區,及分佈式消費,同時保證每個 Partition 內的消息順序傳輸;
- 同時支持離線數據處理和實時數據處理;
- Scale out:支持在線水平擴展;

為何使用消息系統
- 解耦 消息系統在處理過程中間插入了一個隱含的、基於數據的接口層,兩邊的處理過程都要實現這一接口。這允許你獨立的擴展或修改兩邊的處理過程,只要確保它們遵守同樣的接口約束。 而基於消息發布訂閱的機制,可以聯動多個業務下游子系統,能夠不侵入的情況下分步編排和開發,來保證數據一致性。
- 冗餘 有些情況下,處理數據的過程會失敗。除非數據被持久化,否則將造成丟失。消息隊列把數據進行持久化直到它們已經被完全處理,通過這一方式規避了數據丟失風險。許多消息隊列所採用的”插入-獲取-刪除”範式中,在把一個消息從隊列中刪除之前,需要你的處理系統明確的指出該消息已經被處理完畢,從而確保你的數據被安全的保存直到你使用完畢。
- 擴展性 因為消息隊列解耦了你的處理過程,所以增大消息入隊和處理的頻率是很容易的,只要另外增加處理過程即可。不需要改變代碼、不需要調節參數。擴展就像調大電力按鈕一樣簡單。
- 靈活性 & 峰值處理能力 在訪問量劇增的情況下,應用仍然需要繼續發揮作用,但是這樣的突發流量並不常見;如果為以能處理這類峰值訪問為標準來投入資源隨時待命無疑是巨大的浪費。使用消息隊列能夠使關鍵組件頂住突發的訪問壓力,而不會因為突發的超負荷的請求而完全崩潰。
- 可恢復性 系統的一部分組件失效時,不會影響到整個系統。消息隊列降低了進程間的耦合度,所以即使一個處理消息的進程掛掉,加入隊列中的消息仍然可以在系統恢復後被處理。
- 順序保證 在大多使用場景下,數據處理的順序都很重要。大部分消息隊列本來就是排序的,並且能保證數據會按照特定的順序來處理。 Kafka 保證一個 Partition 內的消息的有序性。
- 緩衝 在任何重要的系統中,都會有需要不同的處理時間的元素。消息隊列通過一個緩衝層來幫助任務最高效率的執行———寫入隊列的處理會盡可能的快速。該緩衝有助於控制和優化數據流經過系統的速度。
- 異步通訊 很多時候,用戶不想也不需要立即處理消息。消息隊列提供了異步處理機制,允許用戶把一個消息放入隊列,但並不立即處理它。想向隊列中放入多少消息就放多少,然後在需要的時候再去處理它們。
Topic & Partition
Topic 在邏輯上可以被認為是一個 queue,每條消費都必須指定它的 Topic,可以簡單理解為必須指明把這條消息放進哪個queue 裡。我們把一類消息按照主題來分類,有點類似於數據庫中的表。

為了使得 Kafka 的吞吐率可以線性提高,物理上把 Topic 分成一個或多個 Partition。 對應到系統上就是一個或若干個目錄。

Broker
Broker:Kafka 集群包含一個或多個服務器,每個服務器節點稱為一個 Broker。
Broker 存儲 Topic 的數據。如果某 Topic 有 N 個 Partition,集群有 N 個 Broker,那麼每個 Broker 存儲該 Topic 的一個 Partition。
從 scale out 的性能角度思考,通過 Broker Kafka server 的更多節點,帶更多的存儲,建立更多的 Partition 把 IO 負載到更多的物理節點,提高總吞吐 IOPS。
從 scale up 的角度思考,一個 Node 擁有越多的 Physical Disk,也可以負載更多的 Partition,提升總吞吐 IOPS。

如果某 Topic 有 N 個 Partition,集群有(N+M)個 Broker,那麼其中有 N 個 Broker 存儲該 Topic 的一個 Partition,剩下的 M 個 Broker 不存儲該 Topic 的 Partition 數據。
如果某 Topic 有 N 個 Partition,集群中 Broker 數目少於 N 個,那麼一個 Broker 存儲該 Topic 的一個或多個 Partition。
Topic 只是一個邏輯概念,真正在 Broker間分佈式的 Partition。 每一條消息被發送到 Broker 中,會根據 Partition 規則選擇被存儲到哪一個 Partition。如果 Partition 規則設置的合理,所有消息可以均勻分佈到不同的 Partition中。

Broker & Partition
實驗條件:3個 Broker,1個 Topic,無Replication, 異步模式,3個 Producer,消息 Payload 為100字節:
當 Partition 數量小於 Broker個數時,Partition 數量越大,吞吐率越高,且呈線性提升。
Kafka 會將所有 Partition 均勻分佈到所有Broker 上, 所以當只有2個 Partition 時,會有2個 Broker 為該 Topic 服務。 3個 Partition 時同理會有3個 Broker 為該 Topic 服務。
當 Partition 數量多於 Broker 個數時,總吞吐量並未有所提升,甚至還有所下降。 可能的原因是,當 Partition 數量為4和5時,不同 Broker 上的 Partition 數量不同, 而 Producer 會將數據均勻發送到各 Partition 上,這就造成各Broker 的負載不同, 不能最大化集群吞吐量。

存储原理
Kafka 的消息是存在於文件系統之上的。 Kafka 高度依賴文件系統來存儲和緩存消息,一般的人認為 “磁盤是緩慢的”。
操作系統還會將主內存剩餘的所有空閒內存空間都用作磁盤緩存, 所有的磁盤讀寫操作都會經過統一的磁盤緩存(除了直接 I/O 會繞過磁盤緩存)。
Kafka 正是利用順序 IO,以及 Page Cache 達成的超高吞吐。
任何發佈到 Partition 的消息都會被追加到 Partition 數據文件的尾部,這樣的順序寫磁盤操作讓 Kafka 的效率非常高。

Kafka 集群保留所有發布的 message,不管這個 message 有沒有被消費過, Kafka 提供可配置的保留策略去刪除舊數據(還有一種策略根據分區大小刪除數據)。
例如,如果將保留策略設置為兩天,在 message 寫入後兩天內,它可用於消費,之後它將被丟棄以騰出空間。 Kafka 的性能跟存儲的數據量的大小無關, 所以將數據存儲很長一段時間是沒有問題的。
Offset:偏移量。每條消息都有一個當前 Partition 下唯一的 64 字節的 Offset,它是相當於當前分區第一條消息的偏移量,即第幾條消息。
消費者可以指定消費的位置信息,當消費者掛掉再重新恢復的時候,可以從消費位置繼續消費。

假設我們現在 Kafka 集群只有一個 Broker,我們創建 2 個 Topic 名稱分別為:「Topic1」和「Topic2」,Partition 數量分別為 1、2。 那麼我們的根目錄下就會創建如下三個文件夾:

在 Kafka 的文件存儲中,同一個 Topic 下有多個不同的 Partition,每個 Partition 都為一個目錄。 而每一個目錄又被平均分配成多個大小相等的 Segment File 中,Segment File 又由 index file 和 data file 組成,他們總是成對出現,後綴 ".index" 和 ".log" 分錶表示 Segment 索引文件和數據文件。

Segment 是 Kafka 文件存儲的最小單位。 Segment 文件命名規則:Partition 全局的第一個 Segment 從 0 開始,後續每個 Segment 文件名為上一個 Segment 文件最後一條消息的 offset 值。
其中以索引文件中元數據 <3, 497=""> 為例, 依次在數據文件中表示第 3 個 Message(在全局 Partition 表示第 368769 + 3 = 368772 個 message)以及該消息的物理偏移地址為 497。3,>
注意該 Index 文件並不是從0開始,也不是每次遞增 1 的,這是因為 Kafka 採取稀疏索引存儲的方式,每隔一定字節的數據建立一條索引。
它減少了索引文件大小,使得能夠把 Index 映射到內存,降低了查詢時的磁盤 IO 開銷,同時也並沒有給查詢帶來太多的時間消耗。
因為其文件名為上一個 Segment 最後一條消息的 Offset ,所以當需要查找一個指定 Offset 的 Message 時,通過在所有 Segment 的文件名中進行二分查找就能找到它歸屬的 Segment。
再在其 Index 文件中找到其對應到文件上的物理位置,就能拿出該 Message。

Kafka 是如何準確的知道 Message 的偏移的呢? 這是因為在 Kafka 定義了標準的數據存儲結構,在 Partition 中的每一條 Message 都包含了以下三個屬性:
- Offset:表示 Message 在當前 Partition 中的偏移量,是一個邏輯上的值,唯一確定了 Partition 中的一條 Message,可以簡單的認為是一個 ID。
- MessageSize:表示 Message 內容 Data 的大小。
- Data:Message 的具體內容。

例如讀取 offset=368776的 message,需要通過下面2個步驟查找。
- 第一步查找 segment file 上述圖2為例,其中00000000000000000000.index 表示最開始的文件,起始偏移量(offset)為0。第二個文件00000000000000368769.index 的消息量起始偏移量為368770 = 368769 + 1,其他後續文件依次類推,以起始偏移量命名並排序這些文件,只要根據 offset 二分查找文件列表,就可以快速定位到具體文件。當 offset=368776時定位到00000000000000368769.index | log。
- 第二步通過 segment file 查找 message 通過第一步定位到 segment file,當 offset=368776時,依次定位到00000000000000368769.index 的元數據物理位置和00000000000000368769.log 的物理偏移地址,然後再通過00000000000000368769.log 順序查找直到offset=368776 為止。

segment index file採取稀疏索引存儲方式,它減少索引文件大小,通過mmap可以直接內存操作。 Kafka高效文件存儲設計特點
Kafka把topic中一個parition大文件分成多個小文件段,通過多個小文件段,就容易定期清除或刪除已經消費完文件,減少磁盤佔用。 通過索引信息可以快速定位message和確定response的最大大小。 通過index元數據全部映射到memory,可以避免segment file的IO磁盤操作。 通過索引文件稀疏存儲,可以大幅降低index文件元數據佔用空間大小。
Kafka 從0.10.0.0版本起,為分片日誌文件中新增了一個 .timeindex 的索引文件,可以根據時間戳定位消息。
同樣我們可以通過腳本 kafka-dump-log.sh 查看時間索引的文件內容。
- 首先定位分片,將 1570793423501 與每個分片的最大時間戳進行對比(最大時間戳取時間索引文件的最後一條記錄時間,如果時間為 0 則取該日誌分段的最近修改時間),直到找到大於或等於 1570793423501 的日誌分段,因此會定位到時間索引文件00000000000003257573.timeindex,其最大時間戳為 1570793423505。
- 重複 offset 找到 log 文件的步驟。

Producer & Consumer
Producer
Producer 發送消息到 Broker 時,會根據 Partition 機制選擇將其存儲到哪一個 Partition。 如果 Partition 機制設置合理,所有消息可以均勻分佈到不同的 Partition裡,這樣就實現了負載均衡。 指明 Partition 的情況下,直接將給定的 Value 作為 Partition 的值。 沒有指明 Partition 但有 Key 的情況下,將 Key 的 Hash 值與分區數取餘得到 Partition 值。 既沒有 Partition 有沒有 Key 的情況下,第一次調用時隨機生成一個整數(後面每次調用都在這個整數上自增),將這個值與可用的分區數取餘,得到 Partition 值,也就是常說的 Round-Robin 輪詢算法。

為保證 Producer 發送的數據,能可靠地發送到指定的 Topic,Topic 的每個 Partition 收到 Producer 發送的數據後,都需要向 Producer 發送 ACK。 如果 Producer 收到 ACK,就會進行下一輪的發送,否則重新發送數據。
- 選擇完分區後,生產者知道了消息所屬的主題和分區,它將這條記錄添加到相同主題和分區的批量消息中,另一個線程負責發送這些批量消息到對應的 Kafka Broker。
- 當 Broker 接收到消息後,如果成功寫入則返回一個包含消息的主題、分區及位移的 RecordMetadata 對象,否則返回異常。
- 生產者接收到結果後,對於異常可能會進行重試。

Producer Exactly Once
0.11 版本的 Kafka,引入了冪等性:Producer 不論向 Server 發送多少重複數據,Server 端都只會持久化一條。
- 要啟用冪等性,只需要將 Producer 的參數中 enable.idompotence 設置為 true 即可。
- 開啟冪等性的 Producer 在初始化時會被分配一個 PID,發往同一 Partition 的消息會附帶 Sequence Number。
- 而 Borker 端會對
做緩存,當具有相同主鍵的消息提交時,Broker 只會持久化一條。 - 但是 PID 重啟後就會變化,同時不同的 Partition 也具有不同主鍵,所以冪等性無法保證跨分區會話的 Exactly Once。
Consumer
假設這麼個場景:我們從 Kafka 中讀取消息,並且進行檢查,最後產生結果數據。
我們可以創建一個消費者實例去做這件事情,但如果生產者寫入消息的速度比消費者讀取的速度快怎麼辦呢? 這樣隨著時間增長,消息堆積越來越嚴重。對於這種場景,我們需要增加多個消費者來進行水平擴展。 Kafka 消費者是消費組的一部分,當多個消費者形成一個消費組來消費主題時,每個消費者會收到不同分區的消息。
假設有一個 T1 主題,該主題有 4 個分區;同時我們有一個消費組 G1,這個消費組只有一個消費者 C1。 那麼消費者 C1 將會收到這 4 個分區的消息。

如果我們增加新的消費者 C2 到消費組 G1,那麼每個消費者將會分別收到兩個分區的消息。 相當於 T1 Topic 內的 Partition 均分給了 G1 消費的所有消費者,在這裡 C1 消費 P0 和 P2,C2 消費 P1 和 P3。

如果增加到 4 個消費者,那麼每個消費者將會分別收到一個分區的消息。 這時候每個消費者都處理其中一個分區,滿負載運行。

但如果我們繼續增加消費者到這個消費組,剩餘的消費者將會空閒,不會收到任何消息。 總而言之,我們可以通過增加消費組的消費者來進行水平擴展提升消費能力。 這也是為什麼建議創建主題時使用比較多的分區數,這樣可以在消費負載高的情況下增加消費者來提升性能。 另外,消費者的數量不應該比分區數多,因為多出來的消費者是空閒的,沒有任何幫助。 如果我們的 C1 處理消息仍然還有瓶頸,我們如何優化和處理?
把 C1 內部的消息進行二次 sharding,開啟多個 goroutine worker 進行消費,為了保障 offset 提交的正確性,需要使用 watermark 機制,保障最小的 offset 保存,才能往 Broker 提交。

Consumer Group
Kafka 一個很重要的特性就是,只需寫入一次消息,可以支持任意多的應用讀取這個消息。 換句話說,每個應用都可以讀到全量的消息。為了使得每個應用都能讀到全量消息,應用需要有不同的消費組。 對於上面的例子,假如我們新增了一個新的消費組 G2,而這個消費組有兩個消費者如圖。 在這個場景中,消費組 G1 和消費組 G2 都能收到 T1 主題的全量消息,在邏輯意義上來說它們屬於不同的應用。 最後,總結起來就是:如果應用需要讀取全量消息,那麼請為該應用設置一個消費組;如果該應用消費能力不足,那麼可以考慮在這個消費組裡增加消費者。

可以看到,當新的消費者加入消費組,它會消費一個或多個分區,而這些分區之前是由其他消費者負責的。 另外,當消費者離開消費組(比如重啟、當機等)時,它所消費的分區會分配給其他分區。這種現象稱為重平衡(Rebalance)。 重平衡是 Kafka 一個很重要的性質,這個性質保證了高可用和水平擴展。不過也需要注意到,在重平衡期間,所有消費者都不能消費消息,因此會造成整個消費組短暫的不可用。 而且,將分區進行重平衡也會導致原來的消費者狀態過期,從而導致消費者需要重新更新狀態,這段期間也會降低消費性能。 消費者通過定期發送心跳(Hearbeat)到一個作為組協調者(Group Coordinator)的 Broker 來保持在消費組內存活。這個 Broker 不是固定的,每個消費組都可能不同。 當消費者拉取消息或者提交時,便會發送心跳。如果消費者超過一定時間沒有發送心跳,那麼它的會話(Session)就會過期,組協調者會認為該消費者已經當機,然後觸發重平衡。
可以看到,從消費者當機到會話過期是有一定時間的,這段時間內該消費者的分區都不能進行消息消費。 通常情況下,我們可以進行優雅關閉,這樣消費者會發送離開的消息到組協調者,這樣組協調者可以立即進行重平衡而不需要等待會話過期。 在 0.10.1 版本,Kafka 對心跳機制進行了修改,將發送心跳與拉取消息進行分離,這樣使得發送心跳的頻率不受拉取的頻率影響。 另外更高版本的 Kafka 支持配置一個消費者多長時間不拉取消息但仍然保持存活,這個配置可以避免活鎖(livelock)。活鎖,是指應用沒有故障但是由於某些原因不能進一步消費。 但是活鎖也很容易導致連鎖故障,當消費端下游的組件性能退化,那麼消息消費會變的很慢,會很容易出發 livelock 的重新均衡機制,反而影響力吞吐。
Partition 會為每個 Consumer Group 保存一個偏移量,記錄 Group 消費到的位置。
Kafka 0.9開始將消費端的位移信息保存在集群的內部主題(__consumer_offsets)中,該主題默認為50個分區,每條日誌項的格式都是:

分組協調者(Group Coordinator)是一個服務,kafka集群中的每個節點在啟動時都會啟動這樣一個服務,該服務主要是用來存儲消費分組相關的元數據信息,每個消費組均會選擇一個協調者來負責組內各個分區的消費位移信息存儲,選擇的主要步驟如下:
首選確定消費組的位移信息存入哪個分區:前面提到默認的__consumer_offsets主題分區數為50, 通過以下算法可以計算出對應消費組的位移信息應該存入哪個分區 partition = Math.abs(groupId.hashCode() % groupMetadataTopicPartitionCount) 其中 groupId 為消費組的id,這個由消費端指定, groupMetadataTopicPartitionCount 為主題分區數。 根據partition尋找該分區的leader所對應的節點broker,該broker的Coordinator即為該消費組的Coordinator。
Consumer Commit Offset
消費端可以通過設置參數 enable.auto.commit 來控制是自動提交還是手動, 如果值為 true 則表示自動提交,在消費端的後台會定時的提交消費位移信息,時間間隔由 auto.commit.interval.ms(默認為5秒)。
- 可能存在重複的位移數據提交到消費位移主題中,因為每隔5秒會往主題中寫入一條消息,不管是否有新的消費記錄,這樣就會產生大量的同 key 消息,其實只需要一條,因此需要依賴前面提到日誌壓縮策略來清理數據。
- 重複消費,假設位移提交的時間間隔為5秒,那麼在5秒內如果發生了 rebalance,則所有的消費者會從上一次提交的位移處開始消費,那麼期間消費的數據則會再次被消費。
我們來看看集中 Delivery Guarantee:
- 讀完消息先 commit 再處理消息。這種模式下,如果 Consumer 在 commit 後還沒來得及處理消息就 crash 了,下次重新開始工作後就無法讀到剛剛已提交而未處理的消息,這就對應於 At most once。
- [讀完消息先處理再 commit。這種模式下,如果在處理完消息之後 commit 之前 Consumer crash 了,下次重新開始工作時還會處理剛剛未 commit 的消息,實際上該消息已經被處理過了。這就對應於At least once。
在很多使用場景下,消息都有一個主鍵,所以消息的處理往往具有冪等性,即多次處理這一條消息跟只處理一次是等效的,那就可以認為是Exactly once。 (筆者認為這種說法比較牽強,畢竟它不是Kafka本身提供的機制,主鍵本身也並不能完全保證操作的冪等性。而且實際上我們說delivery guarantee 語義是討論被處理多少次,而非處理結果怎樣,因為處理方式多種多樣,我們不應該把處理過程的特性——如是否冪等性,當成Kafka本身的Feature)
Consumer Exactly Once
Flink 提供的 checkpoint 機制,結合 Source/Sink 端配合支持 Exactly Once 語義,以 Hive 為例:
- 從 Kafka 消費數據,寫入到臨時目錄
- ck snapshot 階段,將 Offset 存儲到 State 中,Sink 端關閉寫入的文件句柄,以及保存 ckid 到 State 中
- ck complete 階段,commit kafka offset,將臨時目錄中的數據移到正式目錄
- ck recover 階段,恢復 state 信息,reset kafka offset;恢復 last ckid,將臨時目錄的數據移動到正式目錄
如果一定要做到Exactly once,就需要協調offset和實際操作的輸出。經典的做法是引入兩階段提交。如果能讓offset和操作輸入存在同一個地方,會更簡潔和通用。
Push vs Pull
作為一個消息系統,Kafka遵循了傳統的方式,選擇由 Producer 向 Broker push 消息並由 Consumer 從 Broker pull 消息。一些 logging-centric system,比如 Facebook 的 Scribe 和 Cloudera 的 Flume,採用 push 模式。事實上,push 模式 和 pull 模式各有優劣。
push 模式很難適應消費速率不同的消費者,因為消息發送速率是由 Broker 決定的。 push 模式的目標是盡可能以最快速度傳遞消息,但是這樣很容易造成 Consumer 來不及處理消息,典型的表現就是拒絕服務以及網絡擁塞。而 pull 模式則可以根據Consumer 的消費能力以適當的速率消費消息。
對於 Kafka 而言,pull 模式更合適。 pull 模式可簡化 Broker 的設計,Consumer 可自主控制消費消息的速率,同時 Consumer 可以自己控制消費方式——即可批量消費也可逐條消費,同時還能選擇不同的提交方式從而實現不同的傳輸語義。
而 Pull 模式則可以根據 Consumer 的消費能力以適當的速率消費消息。 Pull 模式不足之處是,如果 Kafka 沒有數據,消費者可能會陷入循環中,一直返回空數據。
因為消費者從 Broker 主動拉取數據,需要維護一個長輪詢,針對這一點, Kafka 的消費者在消費數據時會傳入一個時長參數 timeout。如果當前沒有數據可供消費,Consumer 會等待一段時間之後再返回,這段時長即為 timeout。
Leader & Follower
Replication
Kafka 在0.8以前的版本中,並不提供 HA 機制,一旦一個或多個 Broker 當機,則當機期間其上所有 Partition 都無法繼續提供服務。若該 Broker 永遠不能再恢復,亦或磁盤故障,則其上數據將丟失。
在 Kafka 在0.8以前的版本中,是沒有 Replication 的,一旦某一個 Broker 當機,則其上所有的 Partition 數據都不可被消費,這與 Kafka 數據持久性及 Delivery Guarantee 的設計目標相悖。同時 Produce r都不能再將數據存於這些 Partition 中。
- 如果 Producer 使用同步模式則 Producer 會在嘗試重新發送 message.send.max.retries(默認值為3)次後拋出 Exception,用戶可以選擇停止發送後續數據也可選擇繼續選擇發送。而前者會造成數據的阻塞,後者會造成本應發往該 Broker 的數據的丟失。
- 如果 Producer 使用異步模式,則 Producer 會嘗試重新發送 message.send.max.retries(默認值為3)次後記錄該異常並繼續發送後續數據,這會造成數據丟失並且用戶只能通過日誌發現該問題。
由此可見,在沒有 Replication 的情況下,一旦某機器當機或者某個 Broker 停止工作則會造成整個系統的可用性降低。隨著集群規模的增加,整個集群中出現該類異常的機率大大增加,因此對於生產系統而言 Replication 機制的引入非常重要。
Leader
引入 Replication 之後,同一個 Partition 可能會有多個 Replica,而這時需要在這些Replication 之間選出一個 Leader,Producer 和 Consumer 只與這個 Leader 交互,其它 Replica 作為 Follower 從 Leader 中復制數據。
因為需要保證同一個 Partition 的多個 Replica 之間的數據一致性(其中一個宕機後其它 Replica 必須要能繼續服務並且即不能造成數據重複也不能造成數據丟失)。
如果沒有一個 Leader,所有 Replica 都可同時讀/寫數據,那就需要保證多個 Replica 之間互相(N×N條通路)同步數據,數據的一致性和有序性非常難保證,大大增加了 Replication 實現的複雜性,同時也增加了出現異常的機率。而引入 Leader 後,只有 Leader 負責數據讀寫,Follower 只向 Leader 順序 Fetch 數據(N條通路),系統更加簡單且高效。
和大部分分佈式系統一樣,Kafka 處理失敗需要明確定義一個 Broker 是否“活著”。對於 Kafka 而言,Kafka 存活包含兩個條件:
- 副本所在節點需要與 ZooKeeper 維持 session (這個通過 ZK 的 Heartbeat 機制來實現)。
- 從副本的最後一條消息的 offset 需要與主副本的最後一條消息 offset 差值不超過設定閾值(replica.lag.max.messages)或者副本的 LEO 落後於主副本的 LEO 時長不大於設定閾值(replica.lag.time.max.ms),官方推薦使用後者判斷,並在新版本 kafka0.10.0 移除了replica.lag.max.messages 參數。
Leader 會跟踪與其保持同步的 Replica 列表,該列表稱為 ISR(即in-sync Replica)。如果一個 Follower 宕機,或者落後太多,Leader 將把它從 ISR 中移除。當其再次滿足以上條件之後又會被重新加入集合中。 ISR 的引入主要是解決同步副本與異步複製兩種方案各自的缺陷:
- 同步副本中如果有個副本宕機或者超時就會拖慢該副本組的整體性能。
- 如果僅僅使用異步副本,當所有的副本消息均遠落後於主副本時,一旦主副本宕機重新選舉,那麼就會存在消息丟失情況。
replicated log 是分佈式日誌系統,主要保證:
- commit log 不會丟失
- commit log 在不同機器上是一致的
羅列幾個常見的基於主從復制的 replicated log 實現:
- raft:基於多數節點的 ack,節點一般稱為 leader/follower,kafka 將要使用
- pacificA:基於所有節點的 ack,節點一般稱為 primary/secondary,kafka 正在使用
- bookkeeper:基於法定個數節點的 ack,節點一般稱為 writer/bookie,pulsar 正在使用
Kafka 在 Zookeeper 中動態維護了一個 ISR(in-sync replicas),這個 ISR 裡的所有 Replica都跟上了 leader,只有 ISR 裡的成員才有被選為 Leader 的可能。在這種模式下,對於 f+1 個Replica,一個 Partition 能在保證不丟失已經 commit的消息的前提下容忍 f 個 Replica 的失敗。在大多數使用場景中,這種模式是非常有利的。事實上,為了容忍 f 個 Replica 的失敗,Majority Vote 和 ISR 在 commit 前需要等待的 Replica 數量是一樣的,但是 ISR 需要的總的Replica 的個數幾乎是 Majority Vote 的一半。
而對於Producer而言,它可以選擇是否等待消息commit,這可以通過request.required.acks來設置。這種機制確保了只要ISR有一個或以上的Follower,一條被commit的消息就不會丟失。
Preferred Replica
Kafka 集群 Partition Replication 默認自動分配。
在 Kafka 集群中,每個 Broker 都有均等分配Partition 的 Leader 機會。
- 上述圖 Broker Partition 中,箭頭指向為副本,以Partition-0 為例:Broker1 中 parition-0 為 Leader,Broker2 中 Partition-0 為副本。
- 上述圖種每個 Broker (按照 BrokerId 有序)依次分配主 Partition,下一個 Broker 為副本,如此循環迭代分配,多副本都遵循此規則。
副本分配算法如下:
- 將所有 N Broker 和待分配的 i 個 Partition 排序。
- 將第 i 個 Partition 分配到第(i mod n)個 Broker 上。
- 將第 i 個 Partition 的第 j 個副本分配到第((i + j) mod n)個 Broker 上。
創建1個 Topic 包含4個 Partition,2 Replication:

當集群中新增2節點,Partition 增加到6個:

ACK
對於某些不太重要的數據,對數據的可靠性要求不是很高,能夠容忍數據的少量丟失,所以沒必要等 ISR 中的 Follower 全部接受成功。 只有被 ISR 中所有 Replica 同步的消息才被 Commit,但Producer 發布數據時,Leader 並不需要 ISR 中的所有 Replica 同步該數據才確認收到數據。
- 0:Producer 不等待 Broker 的 ACK,這提供了最低延遲,Broker 一收到數據還沒有寫入磁盤就已經返回,當 Broker 故障時有可能丟失數據。
- 1:Producer 等待 Broker 的 ACK,Partition 的 Leader 落盤成功後返回 ACK,如果在 Follower 同步成功之前 Leader 故障,那麼將會丟失數據。
- -1(all):Producer 等待 Broker 的 ACK,Partition 的 Leader 和 Follower 全部落盤成功後才返回 ACK。但是在 Broker 發送 ACK 時,Leader 發生故障,則會造成數據重複。

數據可靠性
數據一致性
可用性
性能優化
Reference
- 我以为我对Kafka很了解,直到我看了这篇文章
- 全面解析kafka架构与原理
- 我用kafka两年踩过的一些非比寻常的坑
- Kafka设计解析(一) Kafka背景及架构介绍
- Kafka设计解析(二)- Kafka High Availability (上)
- Kafka设计解析(三)- Kafka High Availability (下
- Kafka设计解析(四)- Kafka Consumer设计解析
- Kafka设计解析(五)- Kafka性能测试方法及Benchmark报告
- Kafka设计解析(六)- Kafka高性能架构之道
- 两万字深入剖析Kafka,你学会了吗?
- kafka集群partition分布原理分析
- Kafka文件存储机制那些事
- 一文彻底搞清 Kafka 的副本复制机制
- 如何理解Kafka的消息可靠性策略?
- 端到端一致性,流系统Spark/Flink/Kafka/DataFlow对比总结(压箱宝具呕血之作)
- Kafka集群突破百万partition 的技术探索
- 美团数据平台Apache Kafka系统实践
- 基于SSD的Kafka应用层缓存架构设计与实现
- 千亿级数据量的 Kafka 深度实践
- LinkedIn的Kafka:我是如何做到1秒发布450万+条消息!
- 万亿级消息队列Kafka在滴滴的实践
- 快手万亿级别Kafka集群应用实践与技术演进之路
- kafka的leader选举过程