Kafka 至 Iceberg 管道的效能特徵

本頁面說明 Apache Beam 2.75.0 版的效能特徵,適用於從 Apache Kafka 讀取資料並寫入 Apache Iceberg 資料表的 Dataflow 串流工作。這項測試會評估直接寫入 Apache Iceberg 與透過 Managed BigQuery API 寫入的效能差異,並將這些結果與從 Kafka 到 BigQuery 管道的基準比較。由於 Apache Iceberg I/O 的最佳化作業仍在進行中,這些效能指標可能會有所變動。

基準比較適用於三種主要無狀態對應設定 (也就是從來源讀取資料、將訊息轉換為記錄,然後寫入接收器,不會追蹤狀態或套用複雜的商業邏輯;在基準中稱為 map_only 或 mapping):

  1. Kafka 至 BigQuery (map_only) (以Kafka 至 BigQuery 效能為基準)
  2. Kafka to Iceberg Direct (map_only、autosharding=false)
  3. 使用 Managed BigQuery API 將 Kafka 資料匯入 Iceberg (map_only)

此外,本指南還會評估直接 Apache Iceberg 串流模式 (例如使用 groupbykey 的有狀態批次處理),並詳細說明有關檔案大小分布、自動分片行為和讀取端查詢延遲的重要下游考量。

測試方法

基準測試是使用下列資源進行:

  • Managed Service for Apache Kafka 叢集:流量是使用 Dataflow Streaming Data Generator 範本產生。
    • 輸入總處理量:1 GBps
    • 訊息傳送率:每秒約 1,000,000 則訊息
    • 訊息格式:JSON 文字,具有固定結構定義 (每則訊息約 1 KB)
    • 分區:1,000 個 Kafka 分區
  • 目的地接收器:
    • BigQuery:使用 BigQuery Storage Write API 寫入的標準資料表 (未經分割)。
    • Apache Iceberg:由 Cloud Storage 支援的目錄。直接接收器會使用 bucket(id, 64) 分區 (依主鍵分區到 64 個分片),並使用 hash 分配模式。

水平自動調度資源穩定後,每個管道設定都以穩定狀態執行 24 小時。每個管道案例的基準測試都分別執行了 3 次,而所有回報的值都代表這些執行作業的計算平均值,以確保持續可靠的成效指標。

擷取效能:對應工作負載

無狀態對應管道會從來源讀取資料、將訊息格式轉換為記錄,然後寫入接收器,不會追蹤記錄的狀態。以下各節會分析以 1 GBps 執行的參考架構。

工作設定

設定 將 Kafka 資料寫入 BigQuery (map_only) Kafka 至 Iceberg Direct (autosharding=false) 使用 Managed BigQuery API 將 Kafka 資料寫入 Iceberg
Worker 機型 e2-standard-2 e2-standard-4 e2-standard-4
每個工作站的 vCPU 數量 2 4 4
每個工作人員的 RAM 8 GB 16 GB 16 GB
Streaming Engine 已啟用 已啟用 已啟用
水平自動調度 已啟用 已啟用 已啟用
觸發頻率 5 秒 60 秒 60 秒

處理量和資源用量

直接寫入物件儲存空間中的實體 Parquet 檔案,會比 BigQuery 串流擷取產生更高的 I/O 負擔。與直接寫入 Iceberg 相比,透過 Managed BigQuery API 路由寫入作業可提高工作站 CPU 使用率 (約 70% 對約 60%),並稍微減少 Streaming Engine 消耗量 (約 180 SECU/小時對約 200 SECU/小時),但整體工作站運算需求仍相似 (約 440 個 vCPU 對約 450 個 vCPU)。

指標 將 Kafka 資料寫入 BigQuery (map_only) Kafka 至 Iceberg Direct (autosharding=false) 使用 Managed BigQuery API 將 Kafka 資料寫入 Iceberg
每個工作人員的平均輸入總處理量 ~15 MBps ~9 MB/秒 ~9 MB/秒
平均 CPU 使用率 ~70% 約 60% ~70%
1 GBps 輸入的預估 vCPU 數量 ~126 個 vCPU ~450 個 vCPU ~440 個 vCPU
1 GBps 輸入的預估工作站數量 約 63 名工作人員 約 110 名員工 約 110 名員工
每小時每 1 GBps 的預估 SECU ~58 SECU/小時 約 200 SECU/小時 約 180 SECU/小時

寫入延遲設定檔

由於物件儲存空間中繼資料提交限制,直接 Iceberg 寫入作業會出現嚴重的尾部延遲 (第 99 百分位數)。使用 Managed BigQuery API 可消除尾端延遲尖峰,同時維持低中位數延遲。

端對端寫入延遲 將 Kafka 資料寫入 BigQuery Kafka 至 Iceberg Direct (autosharding=false) 使用 Managed BigQuery API 將 Kafka 資料寫入 Iceberg
P50 (中位數) ~1,200 毫秒 ~1,000 毫秒 ~1,000 毫秒
P95 ~3,000 毫秒 ~7,400 毫秒 ~1,900 毫秒
P99 (尾部) 約 5,400 毫秒 ~14,000 毫秒 ~2,700 毫秒

自動分片注意事項和設計選項

本節將探討寫入 Apache Iceberg 時,自動分片對檔案大小和管道延遲的影響。

為什麼選擇 autosharding=false 做為基準

在初步測試中,啟用自動分片功能會導致檔案大小縮減為極小的區塊,並因區域性執行緒層級負載暴增而觸發動態分片分割,導致檔案大小任意波動,即使在總輸入負載恆定的情況下也是如此。

為維持穩定且可預測的 Parquet 檔案版面配置 (平均約 800 KB),並確保公平的基準,不會過早排清,因此選擇了 autosharding=false 做為直接接收器設定。

停用自動分片與啟用自動分片有何不同?

  • 使用 autosharding=false (基準):與自動分片相比,您可獲得較大的初始檔案大小 (平均約 800 KB)。雖然與理想的 Iceberg 檔案大小 (128 到 512 MB) 相比,這個大小仍偏小,但需要的下游壓縮量明顯較少。不過,由於物件儲存中繼資料瓶頸,代價是寫入尾端延遲時間較長 (P99 達到約 14.0 秒)。
  • 如果啟用自動分片:Dataflow 會動態調度寫入器執行緒,吸收本機輸送量尖峰,進而縮短寫入尾端延遲時間。不過,這會產生大量小型 Parquet 檔案 (約 100 KB 以下),導致儲存層受到影響。這些檔案大小差異極大,且在執行期間會任意波動 (平均約為 39 KB 到 100 KB),因此需要積極進行下游壓縮維護。

分區調整與建議

在評估期間,我們針對目的地資料表實驗了各種固定分割值,以找出最佳平衡點。我們發現,使用 64 個 bucket (例如 bucket(id, 64)) 進行目的地資料表分區,可產生目標檔案大小,同時維持良好的使用率和處理量。這種做法可讓我們享有自動分片帶來的效能優勢,同時避免與完全動態調整相關的任意檔案大小片段問題。

實務人員建議:建議客戶使用目標分割區設定執行類似的初步測試,找出最佳平衡點,在不影響 Parquet 檔案大小的情況下,盡量提高管道平行處理能力。

下游讀取影響:檔案大小和壓縮

雖然寫入端指標偏好使用 Managed BigQuery API 擷取 Iceberg,但整體管道效率很大程度取決於下游讀取效能:

  • 透過 Managed BigQuery API 產生小型檔案:Managed BigQuery API 會頻繁排清資料,確保寫入延遲時間較短。這項行為會導致大量小型 Parquet 檔案寫入目標 Iceberg 目錄。
  • 讀取查詢延遲影響:查詢引擎 (例如 Starburst/Trino、Apache Spark、BigQuery、Dremio) 讀取含有數百萬個小型 Parquet 檔案的資料表時,會產生大量的中繼資料剖析負擔和分割區掃描懲罰。
  • 壓縮需求:為避免使用 Managed BigQuery API 時讀取效能下降 (或直接寫入時啟用自動分片),請定期執行 Iceberg 壓縮維護工作 (例如 REWRITE DATA FILES)。壓縮的運算負擔應納入整體架構設計。
  • 直接寫入 (autosharding=false) 檔案發布:直接 Iceberg 寫入搭配固定分片會產生較大的平均 Parquet 檔案 (~800 KB),因此可提供較不分散的版面配置,方便立即查詢存取,不必立即壓縮 (但仍低於理想範圍)。

有狀態的直接 Iceberg 管道 (groupbykey)

為評估手動批次處理策略,我們針對基準 Kafka to Iceberg Direct (map_only、autosharding=false) 管道測試了具狀態的鍵分組 (groupbykey)。這兩種設定都會直接將 Parquet 檔案寫入物件儲存空間。

基準比較

指標 / 功能 直接接收器基準 (autosharding=false) 有狀態直接接收器 (groupbykey) 對效能的影響
每 1 GBps 的預估 vCPU 數 ~450 個 vCPU ~520 個 vCPU ~+16% 運算資源需求
平均 CPU 使用率 ~60% 約 50% ~-17% 的工作人員效率
每小時每 1 GBps 的預估 SECU ~200 SECU/小時 ~300 SECU/小時 Streaming Engine 負載增加約 50%
平均檔案大小 ~800 KB ~100 KB 產生較小的檔案批次
P50 延遲時間 ~1,000 毫秒 ~1,200 毫秒 中位數速度慢約 20%
第 95 個百分位數的延遲時間 ~7,400 毫秒 ~5,500 毫秒 延遲時間縮短約 26%
P99 延遲時間 ~14,000 毫秒 ~13,000 毫秒 尾端延遲變化幅度不大

取捨分析

  1. Streaming Engine 負荷:新增有狀態的 groupbykey 步驟時,Beam 必須在視窗界線之間儲存中間狀態。這會使 Streaming Engine 運算單元消耗量增加 ~50% (從 ~200 SECU/小時增加至 ~300 SECU/小時)。
  2. 緩衝延遲時間:手動金鑰匯總會強制執行視窗緩衝,導致中位數寫入延遲時間 (P50) 增加至 ~1,200 毫秒,P95 延遲時間增加至 ~5.5 秒。

反向管道:從 Iceberg 串流至 Kafka

為評估雙向 lakehouse 功能,我們也針對反向流動的串流資料進行基準測試,也就是從 Apache Iceberg 資料表讀取僅附加串流,然後發布回 Apache Kafka。

工作設定和效率

與必須處理大量物件儲存空間檔案寫入或中繼資料提交瓶頸的擷取管道不同,從 Iceberg 讀取及串流輸出變更的效率很高:

指標 Iceberg 至 Kafka (僅限附加、完全一次)
Worker 機型 e2-standard-4
1 GBps 輸入的預估 vCPU 數量 約 30 個 vCPU
1 GBps 輸入的預估工作站數量 約 7 名工作人員
每小時每 1 GBps 的預估 SECU 約 0.2 SECU/小時

反向管道的重點摘要

  • 大幅降低運算負荷:從 Iceberg 讀取及投射 CDC 串流時,需要的運算資源大幅減少 (約 30 個 vCPU,直接寫入則約 450 個 vCPU),因為這項作業可避免將大量 Parquet 檔案分割、編碼及提交至物件儲存空間的繁重工作。
  • 資源效率:與傳入擷取路徑相比,從 Lakehouse 格式向下游事件驅動消耗或複製回串流層的效率極高。

架構建議摘要

架構模式 P99 寫入延遲時間 檔案版面配置 下游讀取注意事項
Kafka 至 BigQuery (map_only) 最多 5.4 秒 不適用 最佳化 (BigQuery 代管儲存空間引擎)
使用 Managed BigQuery API 將 Kafka 資料寫入 Iceberg ~2.7 秒 任意小型檔案 需要定期壓縮,才能進行大量讀取
Kafka 至 Iceberg 直接 (autosharding=false) ~14.0 秒 ~800 KB 良好 (初始檔案較大,壓縮需求較低)
Kafka 至 Iceberg 直接 (groupbykey) ~13.0 秒 ~100 KB 中等 (運算和狀態負擔較高)

估算費用

如要使用 Google Cloud 價格計算工具,根據資源用量計費估算類似管道的基準費用,請按照下列步驟操作:

  1. 開啟價格計算機。
  2. 按一下「新增至估算值」。
  3. 選取「Dataflow」。
  4. 在「服務類型」部分,選取「Dataflow Classic」。
  5. 選取「進階設定」即可查看所有選項。
  6. 選擇執行工作的位置。
  7. 在「Job type」(工作類型) 區段選取「Streaming」(串流)。
  8. 選取「啟用 Streaming Engine」。
  9. 輸入工作執行時數、工作站節點、工作站機器和永久磁碟儲存空間的相關資訊。
  10. 輸入預估的 Streaming Engine 運算單元數量。

資源用量和費用大致會隨著輸入處理量線性擴展,但如果是只有少數 worker 的小型工作,總費用會以固定費用為主。您可以根據基準測試結果,推斷工作站節點數量和資源消耗量。

舉例來說,假設您使用 Kafka 至 Iceberg Direct (autosharding=false) 架構執行管道,輸入資料速率為 100 MBps。根據 1 GBps 管道的基準測試結果,您可以估算資源需求,如下所示:

  • 縮放比例:(100 MBps) / (1024 MBps) = 約 0.1
  • 預估工作站節點數:110 名工作者 × 0.1 = 約 11 名工作者
  • 每小時預計使用的 Streaming Engine 運算單元數:200 × 0.1 = 每小時約 20 個單元

這個值僅供初步估算。實際處理量和費用可能會因機型、訊息大小分配、使用者程式碼、彙整類型、鍵值平行處理和視窗大小等因素而有顯著差異。詳情請參閱「Dataflow 成本最佳化的最佳做法」。

執行測試管道

如要使用 Dataflow Flex 範本部署 Apache Iceberg 串流工作,請使用 gcloud dataflow flex-template run 指令。

gcloud dataflow flex-template run JOB_NAME \
  --project=PROJECT_ID \
  --region=REGION \
  --template-file-gcs-location=gs://dataflow-templates-us-central1/latest/flex/Kafka_To_Iceberg_Yaml \
  --enable-streaming-engine \
  --parameters ^@^bootstrapServers="KAFKA_BOOTSTRAP_ADDRESS"\
@topic="KAFKA_TOPIC"\
@table="ICEBERG_TABLE_IDENTIFIER"\
@catalogName="CATALOG_NAME"\
@catalogProperties='{"type":"CATALOG_TYPE","warehouse":"gs://BUCKET_NAME/warehouse/"}'\
@triggeringFrequencySeconds=60\
@schema='SCHEMA_DEFINITION'

更改下列內容:

  • JOB_NAME:Dataflow 工作名稱
  • PROJECT_ID:您的 Google Cloud 專案 ID
  • REGION:工作執行的 Google Cloud 區域 (例如 us-central1)
  • KAFKA_BOOTSTRAP_ADDRESS:Apache Kafka 叢集的啟動位址
  • KAFKA_TOPIC:Kafka 主題的名稱
  • ICEBERG_TABLE_IDENTIFIER:目標 Iceberg 資料表的 ID
  • CATALOG_NAME:Iceberg 目錄名稱
  • CATALOG_TYPE:要使用的目錄類型 (例如 hadoop 或 bigquery)
  • BUCKET_NAME:倉庫位置的 Cloud Storage bucket 名稱
  • SCHEMA_DEFINITION:Kafka 主題資料的結構定義 (例如 {"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]})