跳至主要內容
論文精煉 · Streaming

Kafka: a Distributed Messaging System for Log Processing

一套面向大量事件資料的分散式 commit log:topic 切分成 partition、只做 append 的 segment 檔、由 consumer 自己保管 offset 的 pull 模式,加上 zero-copy 傳輸。

作者Jay Kreps、Neha Narkhede、Jun Rao,LinkedIn Corp. 發表於NetDB 2011(Workshop on Networking Meets Databases),希臘雅典,2011 年 6 月 年份2011
閱讀原始論文 PDF 所有論文

一口氣講完 — 整篇論文的濃縮

LinkedIn 的活動與維運日誌資料量比所謂「真正的」資料高出好幾個數量級,但企業級訊息系統功能太重、吞吐量太低,而 Scribe、Flume 這類日誌聚合器又只服務離線批次消費者。Kafka 把每個 topic 切成多個 partition 分散到一群 broker 上,每個 partition 就只是一串 append-only 的 segment 檔(大小約 1GB),並直接用訊息在 log 中的邏輯 offset 當作它唯一的識別。Consumer 採 pull 而非被 push:每個請求帶上起始 offset 與可接受的位元組數,broker 透過一份很小的記憶體內 offset 索引找到對應 segment,再用 sendfile 把位元組從檔案直接送進 socket,而 broker 完全不保存任何 per-consumer 的投遞狀態。消費位置與 broker、consumer、partition 擁有權等註冊資訊一起放在 ZooKeeper,因此 consumer group 不需要 master 就能自行重新平衡;訊息採時間型 SLA 保留(通常 7 天)而非讀完即刪,任何 consumer 都可以倒帶重播。實測結果是單一 producer 對單一 broker,不批次時每秒 50,000 則、批次 50 則時每秒 400,000 則,至少是 RabbitMQ 的兩倍,比 ActiveMQ 高出數個數量級,代價是只提供 at-least-once 保證。

在這篇論文之前 — 它所降落的世界

到 2011 年,消費性網路公司產生的日誌資料——瀏覽、點擊、搜尋、按讚,加上服務呼叫延遲與每台機器的 CPU、記憶體、網路、磁碟指標——量體已比交易資料高出好幾個數量級;論文引用中國移動每天蒐集 5–8TB 通話紀錄、Facebook 每天累積將近 6TB 使用者活動事件。真正改變的其實不是量,而是消費端:這些資料正從隔夜分析走進線上功能——搜尋相關性、推薦、廣告投放與報表、垃圾訊息與爬取防護、動態消息彙整——而這些功能需要在幾秒內拿到資料。當時兩套現成工具都不合用。IBM Websphere MQ、各家 JMS 實作、TIBCO EMS、Oracle EMS 這類企業訊息系統提供跨多個 queue 的原子插入與可亂序的逐則確認,但對於「偶爾掉幾筆瀏覽事件也不會怎樣」的日誌資料而言,這些保證是殺雞用牛刀;它們沒有 producer 端的批次 API,每一則訊息都要付一次完整的 TCP/IP 來回,跨機器分割的支援也很弱,而且一旦未消費訊息開始堆積,效能就急遽下滑。專門的日誌聚合器——Facebook 的 Scribe、Yahoo 的 Data Highway、Cloudera 的 Flume——擴展性沒問題,但都是為了把資料倒進 HDFS 或 NFS 供離線使用而設計,還把「minute files」這類實作細節暴露給消費者,且採用 push 模式,會把慢速消費者灌爆,而不是讓它按自己撐得住的速度來拉。

問題 — 當時真正壞掉的地方

  • 企業訊息系統把複雜度預算都花在投遞保證上,例如跨多個 queue 的原子插入與可亂序的逐則確認,但對日誌資料而言這些是過度設計——偶爾掉幾筆瀏覽事件並不是世界末日。
  • JMS 沒有讓 producer 把多則訊息打包進單一請求的 API,因此每則訊息都需要一次完整的 TCP/IP 來回,在日誌處理所需的吞吐量下根本行不通。
  • 既有訊息系統對分散式的支援很弱,沒有簡便的方式把訊息分割並儲存到多台機器上。
  • 訊息系統假設訊息幾乎會被立刻消費,一旦未消費訊息堆積效能就明顯劣化;而離線資料倉儲消費者做的正是週期性大批載入而非持續消費,剛好製造出這種狀況。
  • 既有日誌聚合器是為離線消費而建,還把「minute files」這類實作細節暴露給消費者,無法支撐 LinkedIn 所需、延遲不超過數秒的即時應用。
  • 日誌聚合器採用 broker 主動轉發資料給 consumer 的 push 模式,消費者可能被灌到超過自己能處理的速度,也難以倒帶回頭重讀舊資料。

核心概念 — 主要貢獻,以及它們為何成立

把 topic 當成切分過的 log

Topic 就是一條具名的訊息串流:producer 往裡發布、broker 負責儲存、consumer 訂閱它。Kafka 的結構性決策是把每個 topic 切成多個 partition 並散布到叢集中的各個 broker 上,讓一個 topic 的吞吐量不再受限於單一機器的磁碟或網卡。每個 partition 是一條有全序的邏輯 log,而 partition——不是訊息、也不是 topic——才是平行化的最小單位:任何時刻,一個 partition 的所有訊息在同一個 consumer group 內只由一個 consumer 消費。光這一個決定就換到了 partition 內的順序保證,也免除了鎖與逐則狀態維護的開銷,讓橫向擴展變成「把 topic 過度切分」這件簡單的事。

用 offset 取代 message id

存在 Kafka 裡的訊息沒有明確的 message id,而是以它在 log 中的邏輯 offset 定址;下一則訊息的 id 就是目前 id 加上目前訊息的長度。因此 id 遞增但不連續,這個設計也徹底省掉了傳統 broker 為了把 message id 映射到實際儲存位置而維護的輔助性、大量隨機 seek 的索引結構。Broker 只在記憶體中保留一份排序過的 offset 清單,包含每個 segment 檔第一則訊息的 offset,這就足以定位到正確的檔案再往後循序讀取。又因為 consumer 一定是循序讀取某個 partition,確認一個 offset 就等同確認它之前的所有訊息,把逐則的投遞狀態壓縮成一個數字。

用 pull,不用 push

Consumer 發出非同步的 pull 請求,帶上起始 offset 與可接受的抓取位元組數,而不是由 broker 把訊息推過來。每個 consumer 因此都能以自己撐得住的最大速率取用,永遠不會被 producer 或正在清積壓的 broker 灌爆——這正是 LinkedIn 捨棄 Scribe、Data Highway、Flume 所用 push 模式的理由。而且因為位置是請求的參數而非伺服器端狀態,consumer 也可以把它往回調;論文把「倒帶」視為必要功能而非附帶好處。應用程式邏輯出錯修好後可以重播訊息,只做週期性 flush 的全文索引器崩潰重啟後也能從最小的未 flush offset 重新消費。

無狀態 broker 與時間型保留策略

和多數訊息系統不同,Kafka 的 broker 不記錄每個 consumer 消費到哪裡,那是 consumer 自己的事。這替 broker 省下大量簿記工作與磁碟寫入,但也讓 broker 無從得知一則訊息何時可以安全刪除。Kafka 的答案刻意粗暴:訊息在 broker 中保留超過設定的期間(通常是 7 天)就自動刪除,不管誰讀過。這之所以可行,一是 Kafka 的效能不會隨 log 變大而劣化,保留一週只是多花磁碟;二是實務上消費者——包括離線的那些——不是即時消費,就是每小時或每天跑完一輪。

不做應用層快取:靠 page cache 與 sendfile

Kafka 刻意不在自己的行程記憶體中快取訊息,而是仰賴底層檔案系統的 page cache。這避免了雙重緩衝,讓快取即使在 broker 行程重啟後依然是熱的,也讓 JVM 幾乎沒有東西要做垃圾回收,這正是用 VM 語言做出高效實作的前提。又因為 producer 是循序附加、consumer 通常只落後一點點,作業系統原本的快取啟發式(write-through 快取與 read-ahead)剛好就是對的策略;論文回報生產與消費的效能在資料量達到數個 TB 時仍與資料大小呈線性關係。讀取路徑上 Kafka 再用 Unix 的 sendfile API 把 segment 檔的 file channel 直接搬進 socket channel,而 Kafka 是多訂閱者系統、同一批位元組會被重複送出,效益因此被放大。

兩端都成批:message set

批次在 Kafka 裡是 API 的一等公民,而不是事後補的最佳化。Producer 用單一 publish 請求送出一個 MessageSet;consumer 的 pull 請求同樣一次取回多則訊息、通常數百 KB,儘管客戶端的 iterator 還是一次一則交給應用程式。這攤平了 RPC 與 TCP/IP 來回的成本——論文認為那正是 JMS 逐則發布的致命傷——也讓 sendfile 這條路徑真正划算,因為一次呼叫搬的是一大段連續位元組。實驗中把批次大小從 1 調到 50,producer 吞吐量提升了將近一個數量級。

去中心化協調,沒有 master 節點

Kafka 不選出 master broker 或中央協調者,而是讓 consumer 以去中心化的方式彼此協調,理由講得很直白:多一個 master 就多一種要處理的故障。ZooKeeper 提供底層設施:broker 註冊表、consumer 註冊表、每個 group 的 partition 擁有權註冊表,以及每個 group 的 offset 註冊表;前三者都是 ephemeral 路徑,建立者一消失就自動移除。每個 consumer 在 broker 與 consumer 註冊表上註冊 watcher,成員一有變動所有人都會收到通知,接著各自對同一組排序後的輸入跑同一套決定性的分配演算法。協調因此只發生在重新平衡的時刻——一個不常發生的事件——完全不落在訊息路徑上。

運作方式 — 具體的機制

發布端:選 partition 與 publish 請求

訊息被定義成只有一段位元組 payload,序列化方式由使用者自己選;LinkedIn 在上面疊了 Avro,把 Avro schema 的 id 與序列化後的位元組一起放進 payload,再透過一個輕量的 schema registry 服務把 id 換回 schema。多則訊息聚成一個 MessageSet,用單一個 send 呼叫指定 topic 送出。Producer 選擇目標 partition 的方式有兩種:隨機挑,或是用 partitioning key 搭配 partitioning function 決定——後者正是讓帶有相同 join key 的訊息全部落到同一個 partition、因而落到同一個 consumer 行程的關鍵鉤子。Producer 不等 broker 的確認,broker 收多快就送多快,這是發布吞吐量衝高的原因,但也意味著未經確認的訊息可能悄悄遺失。

Broker 儲存:segment 檔與記憶體內 offset 索引

Topic 的每個 partition 對應一條邏輯 log,實體上實作成一組大小大致相同的 segment 檔,例如 1GB。發布動作就只是往最後一個 segment 檔 append——不 seek、不更新索引、也不寫逐則的中繼資料。Segment 檔要等到發布的訊息數達到設定值、或經過設定的時間之後才 flush 到磁碟,而訊息也只有在 flush 之後才對 consumer 可見,耐久性其實是在這個點上決定的。Broker 在記憶體中保留一份排序過的 offset 清單,包含每個 segment 檔第一則訊息的 offset,因此服務一次抓取就是搜尋這份清單找到所屬 segment,然後往後讀;保留策略則是從 log 前端整個 segment 地刪。

消費端:message stream、pull 請求與 offset 運算

Consumer 對某個 topic 呼叫 createMessageStreams,取得一或多條 message stream,發布到該 topic 的訊息會平均分配到這些子串流上;每條串流提供的 iterator 與一般 iterator 不同——它永不終止,log 讀完時會阻塞,等到有新訊息發布才繼續。底層則是持續發出非同步 pull 請求,每個請求帶著開始消費的 offset 與可接受的抓取位元組數,先把資料緩衝好等應用程式取用。收到訊息後,consumer 把訊息長度加上去算出下一個 offset,用在下一次請求。把多個 consumer 放進同一個 group 就得到 point-to-point 投遞;放進不同 group 就得到 publish/subscribe,而不同 group 之間完全不需要協調。

用 sendfile 做 zero-copy 傳輸

把本機檔案送到遠端 socket 的傳統路徑是 4 次資料複製與 2 次系統呼叫:儲存媒體到作業系統 page cache、page cache 到應用程式緩衝區、應用程式緩衝區到核心 socket 緩衝區、再從 socket 緩衝區送上網路。Kafka 改用 Unix 的 sendfile API,把位元組從 segment 檔的 file channel 直接送進 socket channel,省掉其中 2 次複製與 1 次系統呼叫。這之所以做得到,是因為磁碟上的表示法與線路上的表示法是同一批位元組:broker 不做任何 per-consumer 的轉換、不重新封裝、也不需要 id 到位置的查表。而 Kafka 是多訂閱者系統,同一則訊息可能被不同應用程式消費很多次,這份節省因此被反覆放大。

ZooKeeper 註冊表與重新平衡演算法

Broker 啟動時把主機名稱、port 以及自己存放的 topic 與 partition 集合寫進 broker 註冊表;consumer 啟動時把所屬 group 與訂閱的 topic 寫進 consumer 註冊表。每個 consumer group 另外擁有一份 ownership 註冊表——每個訂閱的 partition 一條路徑,值是目前正在讀它的 consumer id——以及一份 offset 註冊表,記錄每個 partition 最後消費到的 offset。Broker、consumer 與 ownership 路徑都是 ephemeral,故障時對應項目會自動消失,只有 offset 註冊表是持久的。Watcher 觸發時,consumer Ci 執行 Algorithm 1:先從 ownership 註冊表移除自己擁有的 partition,讀取兩份註冊表,算出 topic T 可用的 partition 集合 PT 與訂閱 T 的 consumer 集合 CT,兩者都排序,令 j 為 Ci 在 CT 中的索引、N = |PT|/|CT|,認領第 j*N 到 (j+1)*N-1 個 partition,把自己寫成擁有者,再為每個認領到的 partition 開一條執行緒,從 offset 註冊表存的位置開始拉資料。

衝突、資料毀損與投遞保證

由於重新平衡的通知抵達各個 consumer 的時間略有先後,某個 consumer 可能會去搶一個仍被別人持有的 partition;發生時它就釋放自己目前擁有的全部 partition,稍等一下再重試,實務上通常重試幾次就穩定下來。全新的 consumer group 在 offset 註冊表裡沒有任何紀錄,此時會依設定從各 partition 可用的最小或最大 offset 開始,使用 broker 為此提供的 API。Kafka 為 log 中的每則訊息存一份 CRC,讓 broker 遇到 I/O 錯誤時能跑復原流程剔除 CRC 不一致的訊息,也讓客戶端能在發布或消費後檢查網路錯誤。投遞保證是 at-least-once:consumer 未正常關閉就崩潰時,最後提交到 ZooKeeper 的 offset 之後的訊息會留著,接手該 partition 的 consumer 可能重複投遞,因此在意重複的應用程式必須自行用回傳的 offset 或訊息中的唯一鍵去重。

LinkedIn 的部署:鏡像、稽核與 Hadoop 載入

每個跑使用者面向服務的資料中心都同址部署一座 Kafka 叢集;前端服務成批把日誌資料發布進去,中間由硬體負載平衡器把發布請求平均分配到各 broker,線上消費者也跑在同一個資料中心內。另有一座分析用資料中心,設在靠近 Hadoop 叢集與資料倉儲設施的位置,其中的 Kafka 叢集跑一組 embedded consumer 從每個線上叢集拉資料,形成一份複本,載入工作、報表與臨時腳本都對它執行。正確性是靠稽核驗證而非假設:每則訊息帶著產生時間戳與伺服器名稱,每個 producer 週期性地把監控事件發到另一個 topic,記錄固定時間窗內自己對每個 topic 送出的訊息數,consumer 再拿收到的計數去對帳。Hadoop 端的匯入使用自製的 Kafka input format,讓 MapReduce 工作直接從 broker 讀取;又因為 offset 由客戶端保管,資料與 offset 都只在工作成功完成時才寫入 HDFS,所以任務失敗重啟既不會重複也不會遺漏。

論文證明了什麼 — 量測數據與證明

  • 在 producer 測試中——producer 與 broker 各一台機器,每台 8 顆 2GHz 核心、16GB 記憶體、6 顆磁碟組 RAID 10,以 1Gb 網路相連——Kafka 發布 1,000 萬則 200 位元組的訊息,批次大小 1 時平均每秒 50,000 則,批次大小 50 時每秒 400,000 則,比 ActiveMQ v5.4 高出數個數量級,也至少是 RabbitMQ v2.4 的兩倍。
  • 批次大小為 50 時,單一個 Kafka producer 幾乎就把 producer 與 broker 之間的 1Gb 連線塞滿;光是批次本身,就透過攤平 RPC 開銷把吞吐量提升將近一個數量級。
  • 在 consumer 測試中,各系統都設定成每次請求預取大致相同的量(最多 1000 則訊息或約 200KB),且資料全在快取中,Kafka 平均每秒消費 22,000 則訊息,是 ActiveMQ 與 RabbitMQ 的四倍以上。
  • Kafka 每則訊息的儲存額外開銷平均 9 位元組,ActiveMQ 則是 144 位元組,等於同樣 1,000 萬則訊息 ActiveMQ 多用了 70% 空間;作者把一部分歸因於 JMS 要求的厚重訊息標頭,並觀察到 ActiveMQ 最忙的執行緒之一大部分時間都在存取 B-Tree 以維護訊息中繼資料與狀態。
  • Consumer 測試期間 Kafka broker 上完全沒有磁碟寫入活動,而 ActiveMQ 有一條執行緒忙著把 KahaDB 頁面寫進磁碟——這正是「broker 維護逐則投遞狀態」代價的直接量測。
  • 在 LinkedIn 的正式環境中,Kafka 在線上與分析資料中心之間每天累積數百 GB、將近十億則訊息,端到端管線的平均延遲約 10 秒,而且沒有做太多調校。

限制與取捨 — 論文自承的,以及後來被發現的

  • 論文自承:Kafka 完全沒有複本機制。broker 一旦當掉,存在上面尚未被消費的訊息就無法取得;若該 broker 的儲存永久損毀,那些訊息就永遠遺失。跨 broker 的內建複寫(同時支援非同步與同步兩種模式)被列為未來工作的第一項。
  • 論文自承:producer 不等待確認,因此無法保證發布出去的訊息真的被 broker 收到。作者明白地接受這一點——對許多類型的日誌資料而言,只要掉的訊息數量相對少,用耐久性換吞吐量是划算的——同時表示未來會為更關鍵的資料處理耐久性問題。
  • 論文自承:投遞保證只有 at-least-once。consumer 未正常關閉就崩潰時,接手者可能重送最後提交到 ZooKeeper 的 offset 之後的訊息,論文把去重推給應用程式,理由是這比兩階段提交更划算。後來的 Kafka 推翻了這個前提:0.11 版以冪等 producer 與交易支援,在訊息路徑上不做兩階段提交也達成了 exactly-once 語意。
  • 論文自承:順序只在單一 partition 內成立,跨 partition 沒有任何保證;而且既然 partition 是平行化的最小單位,一個 consumer group 的有效 consumer 數就永遠無法超過 partition 數。論文開的藥方是把 topic 過度切分,等於把容量決策提前到建立 topic 的時刻。
  • 後續研究揭露:在當時的規模下,由客戶端透過 ZooKeeper 協調看似便宜,卻撐不過成長——每個 consumer 都監看每份註冊表,造成 herd effect 與反覆的重新平衡風暴,逐 partition 的 offset 提交也讓 ZooKeeper 變成寫入熱點。Kafka 0.8.2 把 offset 移進內部的 __consumer_offsets topic,0.9 把 group 成員管理移到 broker 端的協調者,KIP-500 最終把 ZooKeeper 從 Kafka 中完全移除,改用以 Raft 為基礎的內部中繼資料 log。

它後來變成什麼 — 繼承這個想法的系統

Kafka 成了整個產業預設的事件骨幹,而論文所描述的東西幾乎原封不動地留了下來:topic、partition、offset、consumer group、segment 檔與時間型保留策略。當年被列為未來工作的項目後來變成它最關鍵的能力——0.8 版帶來以 in-sync replica set 為基礎的複寫,0.11 版以冪等 producer 與交易達成 exactly-once 語意,而論文期待的「一套好用的串流工具函式庫」則長成了 Kafka Streams,以及更廣生態中的 Samza、Storm、Spark Streaming 與 Apache Flink,這些系統全都把 Kafka partition 當成真相來源。無狀態 broker 加上客戶端保管 offset 的設計,事後看正是串流處理得以成立的關鍵:因為 log 可重播、位置又是 consumer 自己持有的一個數字,「重跑歷史」與「讓一個全新消費者從頭建起」變成同一個操作,這也是 Kappa 架構、以及 Kafka Connect 與 Debezium 這類 log-based change data capture 的基礎。把分散式 log 當成一種基本元件而非佇列,這個觀點被 AWS Kinesis、Apache Pulsar、Azure Event Hubs、Redpanda 與 NATS JetStream 直接繼承,也間接影響了那些以複寫 log 作為系統真相來源的資料庫設計。效率配方——append-only segment、不做應用層快取、sendfile、成批的 message set、磁碟格式即線路格式——則遠遠超出訊息系統的範圍,成為高吞吐量儲存系統的標準做法。Kafka 自己的軌跡最後把這個圈畫完:KRaft 用內部的 Raft 中繼資料 log 取代了 ZooKeeper,拿掉這篇論文當年唯一倚賴的外部相依。

論文原文 — 逐字引用

“與典型的訊息系統不同,存在 Kafka 中的訊息沒有明確的 message id。相反地,每則訊息是以它在 log 中的邏輯 offset 來定址。”

§3.1

“與多數其他訊息系統不同,在 Kafka 中,每個 consumer 已消費多少的資訊不是由 broker 維護,而是由 consumer 自己維護。”

§3.1 Stateless broker

“一般而言,Kafka 只保證 at-least-once 投遞。exactly-once 投遞通常需要兩階段提交,對我們的應用而言並非必要。”

§3.3

術語 — 依本篇論文的用法

Topic(主題)
某一類訊息組成的串流:producer 往 topic 發布、consumer 訂閱一或多個 topic。一個 topic 會被切成多個 partition 分散到叢集中的各個 broker 上。
Partition(分割)
Topic 的一片,在 broker 上以一條具全序的邏輯 log 儲存。它是 Kafka 平行化的最小單位——任何時刻,一個 partition 的所有訊息在每個 consumer group 內只由一個 consumer 消費。
Broker
儲存已發布訊息的伺服器;一個 Kafka 叢集由多個 broker 組成,每個 broker 持有一或多個 partition。在這篇論文中,broker 不保存任何 per-consumer 狀態,也不持有其他 broker 資料的複本。
Offset(位移)
訊息在其所屬 partition log 中的邏輯位置,用來取代明確的 message id。Offset 遞增但不連續:下一個 offset 等於目前的 offset 加上目前訊息的長度。
Segment file(分段檔)
Partition log 的實體單位,大小大致固定,例如 1GB。發布是往最後一個 segment append,保留策略則從 log 前端整個 segment 刪除,broker 的記憶體索引記錄每個 segment 的第一個 offset。
Message set(訊息集)
一次 publish 請求送出、或一次 pull 回應帶回的一批訊息,consumer 端通常是數百 KB。成批可攤平 RPC 與 TCP/IP 來回的成本,而 JMS 式的逐則發布正是躲不掉這筆開銷。
Consumer group(消費者群組)
一或多個 consumer 共同消費一組訂閱的 topic,每則訊息只投遞給群組中的一員。不同 group 各自獨立消費完整串流、彼此不需要協調,同一個 topic 因此能同時支援 point-to-point 與 publish/subscribe。
Rebalance(重新平衡)
把 partition 重新分配給 consumer 的去中心化流程,由 ZooKeeper watcher 在 broker 或 consumer 加入、離開時觸發。每個 consumer 各自對 partition 集合與 consumer 集合排序、認領一段連續區間,再從 offset 註冊表中的位置繼續消費。
At-least-once delivery(至少一次投遞)
Kafka 唯一明說的保證:每則訊息至少會送達每個 consumer group 一次,但 consumer 非正常崩潰時,最後提交到 ZooKeeper 的 offset 之後可能出現重複。在意的應用程式必須用 offset 或訊息中的唯一鍵去重,作者認為這比兩階段提交更划算。

在時間軸上的位置 — 這篇論文在整段故事中的座標

在時間軸上查看