在這篇論文之前 — 它所降落的世界
到 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 模式,會把慢速消費者灌爆,而不是讓它按自己撐得住的速度來拉。
術語 — 依本篇論文的用法
- 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 或訊息中的唯一鍵去重,作者認為這比兩階段提交更划算。