Kafka から Iceberg へのパイプラインのパフォーマンス特性

このページでは、Apache Beam バージョン 2.75.0 における、Apache Kafka から読み取り、Apache Iceberg テーブルに書き込む Dataflow ストリーミング ジョブのパフォーマンス特性について説明します。Apache Iceberg への直接書き込みと、Managed BigQuery API を介した書き込みのパフォーマンスの違いを評価し、これらの結果を Kafka to BigQuery パイプラインのベースライン ベンチマークと比較します。Apache Iceberg I/O の最適化は継続的に行われているため、これらのパフォーマンス指標は変更される可能性があります。

ベンチマークの比較は、3 つの主要なステートレス マッピング構成(ソースから読み取り、メッセージをレコードに変換し、状態を追跡したり複雑なビジネス ロジックを適用したりせずにシンクに書き込む。ベンチマークでは map_only または mapping と呼ばれます)で行われています。

  1. Kafka to BigQuery (map_only) (baseline from Kafka to BigQuery のパフォーマンス)
  2. Kafka to Iceberg Direct(map_only、autosharding=false)
  3. Managed BigQuery API を使用した Kafka to Iceberg(map_only)

また、このガイドでは、groupbykeyステートフル バッチ処理など、Apache Iceberg の直接ストリーミング パターンを評価し、ファイルサイズの分布、自動シャーディングの動作、および読み取り側のクエリ レイテンシに関する重要なダウンストリームの考慮事項について詳しく説明します。

テスト方法

ベンチマークは次のリソースを使用して実施されました。

  • Managed Service for Apache Kafka クラスタ: トラフィックは、Dataflow Streaming Data Generator テンプレートを使用して生成されました。
    • 入力スループット: 1 GBps
    • メッセージ レート: 1 秒あたり約 1,000,000 メッセージ
    • メッセージ形式: 固定スキーマの JSON テキスト(1 メッセージあたり約 1 KB)
    • パーティション: 1,000 個の Kafka パーティション
  • 宛先シンク:
    • BigQuery: BigQuery Storage Write API を使用して書き込まれた標準テーブル(パーティショニングなし)。
    • Apache Iceberg: Cloud Storage にバックアップされたカタログ。Direct シンク は bucket(id, 64)(主キーで 64 個のシャードにバケット化)を使用してパーティショニングされ、hash 分散モードを使用します。

水平自動スケーリングが安定した後、各パイプライン構成は 24 時間安定した状態で実行されました。各パイプライン ケースのベンチマークは 3 回別々に実行され、報告された値はすべて、持続的で信頼性の高いパフォーマンス指標を確保するために、これらの実行の平均値を表しています。

取り込みパフォーマンス: マッピング ワークロード

ステートレス マッピング パイプラインは、ソースから読み取り、メッセージ形式をレコードに変換し、レコード間で状態を追跡せずにシンクに書き込みます。以降のセクションでは、1 GBps で実行されるリファレンス アーキテクチャについて分析します。

ジョブ構成

設定 Kafka to BigQuery(map_only) Kafka to Iceberg Direct(autosharding=false) Managed BigQuery API を使用した Kafka to Iceberg
ワーカーのマシンタイプ 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 to BigQuery(map_only) Kafka to Iceberg Direct(autosharding=false) Managed BigQuery API を使用した Kafka to Iceberg
ワーカーあたりの平均入力スループット 約 15 MBps 約 9 MBps 約 9 MBps
平均 CPU 使用率 約 70% 約 60% 約 70%
1 GBps 入力に必要な vCPU の推定数 約 126 vCPU 約 450 vCPU 約 440 vCPU
1 GBps 入力に必要なワーカーの推定数 約 63 個のワーカー 約 110 個のワーカー 約 110 個のワーカー
1 GBps あたりの 1 時間あたりの SECU の推定数 約 58 SECU/時間 約 200 SECU/時間 約 180 SECU/時間

書き込みレイテンシ プロファイル

オブジェクト ストレージのメタデータの commit 制約により、Iceberg への直接書き込みではテール レイテンシ(P99)が大きくなります。Managed BigQuery API を使用すると、中央値のレイテンシを低く抑えながら、テール レイテンシの急増を解消できます。

エンドツーエンドの書き込みレイテンシ Kafka to BigQuery Kafka to Iceberg Direct(autosharding=false) Managed BigQuery API を使用した Kafka to 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)を維持し、早期フラッシュのない公平なベースラインを確保するため、Direct シンク構成にはautosharding=false が選択されました。

自動シャーディングを無効にした場合と有効にした場合の違い

  • autosharding=false(ベースライン)の場合: 自動シャーディングと比較して、初期ファイルサイズが大きくなります(平均約 800 KB)。これは、理想的な Iceberg ファイルサイズ(128 ~ 512 MB)と比較するとまだ小さいですが、ダウンストリームの圧縮が大幅に少なくなります。ただし、オブジェクト ストレージのメタデータのボトルネックにより、書き込みテール レイテンシ(P99 が約 14.0 秒に達する)が高くなるというトレードオフがあります。
  • 自動シャーディングが有効になっている場合: Dataflow は、ローカル スループットの急増を吸収するためにライター スレッドを動的にスケーリングし、書き込みテール レイテンシを短縮します。ただし、小さな Parquet ファイル(約 100 KB 以下)が大量に生成されるため、ストレージ レイヤが損なわれます。これらのファイルサイズは分散が大きく、実行ごとに任意に変動するため(平均で約 39 KB ~約 100 KB)、積極的なダウンストリーム圧縮メンテナンスの必要性が高まります。

パーティションのチューニングと推奨事項

評価では、最適なバランスを見つけるために、宛先テーブルのさまざまな固定パーティション値を試しました。宛先テーブルのパーティショニングに64 個のバケット (例: bucket(id, 64))を使用すると、適切な使用率とスループットを維持しながら、目標のファイルサイズが得られることがわかりました。このアプローチにより、完全な動的スケーリングに関連する任意のファイルサイズ フラグメンテーションの問題を回避しながら、自動シャーディングのパフォーマンス上のメリットを実現できました。

実践者向けの推奨事項: Parquet ファイルサイズを損なうことなくパイプラインの並列処理を最大化する最適なポイントを見つけるために、ターゲット パーティション設定で同様の予備テストを実施することをおすすめします。

ダウンストリームの読み取りに関する影響: ファイルサイズと圧縮

書き込み側の指標では Iceberg の取り込みに Managed BigQuery API が適していますが、パイプライン全体の効率はダウンストリームの読み取りパフォーマンスに大きく依存します。

  • Managed BigQuery API での小さいファイルの生成: Managed BigQuery API は、書き込みレイテンシを短くするために、データを頻繁にフラッシュします。この動作により、ターゲット Iceberg カタログに大量の小さな Parquet ファイルが書き込まれます。
  • 読み取りクエリ レイテンシの影響: 数百万の小さな Parquet ファイルを含むテーブルを読み取るクエリエンジン(Starburst/Trino、Apache Spark、BigQuery、Dremio など)では、メタデータの解析オーバーヘッドとパーティション スキャン ペナルティが大きくなります。
  • 圧縮要件: Managed BigQuery API を使用する場合(または直接書き込みで自動シャーディングが有効になっている場合)に読み取りパフォーマンスの低下を防ぐには、定期的に Iceberg 圧縮メンテナンス ジョブ(REWRITE DATA FILES など)を実行します。圧縮のコンピューティング オーバーヘッドは、アーキテクチャ全体の設計に考慮する必要があります。
  • 直接書き込み(autosharding=false)のファイル分布: 固定シャーディングを使用した Iceberg への直接書き込みでは、平均 Parquet ファイルサイズが大きくなります(約 800 KB)。これにより、即時圧縮を必要とせずに、すぐにクエリにアクセスできるフラグメンテーションの少ないレイアウトが実現します(理想的な範囲を下回っています)。

ステートフルな直接 Iceberg パイプライン(groupbykey)

手動バッチ処理戦略を評価するために、ステートフル キー グループ化(groupbykey)をベースラインの Kafka to Iceberg Direct(map_only、autosharding=false) パイプラインに対してテストしました。どちらの構成でも、Parquet ファイルはオブジェクト ストレージに直接書き込まれます。

ベンチマークの比較

指標 / 機能 Direct シンクのベースライン(autosharding=false) ステートフル Direct シンク(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% 遅い
P95 レイテンシ 約 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 へのストリーミング

双方向のレイクハウス機能を評価するために、逆方向に流れるストリーミング データ(Apache Iceberg テーブルから追加専用ストリームを読み取り、Apache Kafka にパブリッシュする)のベンチマークも実施しました。

ジョブ構成と効率

オブジェクト ストアへの大量のファイル書き込みやメタデータの commit ボトルネックに対処する必要がある取り込みパイプラインとは異なり、Iceberg からの変更の読み取りとストリーミングは非常に効率的に動作します。

指標 Iceberg to Kafka(追加専用、1 回限り)
ワーカーのマシンタイプ e2-standard-4
1 GBps 入力に必要な vCPU の推定数 約 30 vCPU
1 GBps 入力に必要なワーカーの推定数 約 7 個のワーカー
1 GBps あたりの 1 時間あたりの SECU の推定数 約 0.2 SECU/時間

リバース パイプラインの重要なポイント

  • コンピューティング オーバーヘッドの大幅な削減: Iceberg から CDC ストリームを読み取って投影するには、パーティショニング、エンコード、大量の Parquet ファイルのオブジェクト ストレージへの commit という負荷の高い処理を回避できるため、必要なコンピューティング リソースが大幅に少なくなります(直接書き込みの場合は約 450 vCPU に対して約 30 vCPU)。
  • リソース効率: レイクハウス形式からストリーミング レイヤへのダウンストリームのイベント ドリブン消費またはレプリケーションは、インバウンドの取り込みパスと比較して非常に効率的です。

アーキテクチャに関する推奨事項の概要

アーキテクチャ パターン P99 書き込みレイテンシ ファイル レイアウト ダウンストリームの読み取りに関する考慮事項
Kafka to BigQuery(map_only) 約 5.4 秒 なし 最適 (Managed BigQuery Storage Engine)
Managed BigQuery API を使用した Kafka to Iceberg 約 2.7 秒 任意の小さいファイル 大量の読み取りには定期的な圧縮が必要
Kafka to Iceberg Direct(autosharding=false) 約 14.0 秒 約 800 KB 良好 (初期ファイルサイズが大きい、圧縮の必要性が低い)
Kafka to Iceberg Direct(groupbykey) 約 13.0 秒 約 100 KB 中程度 (コンピューティングと状態のオーバーヘッドが大きい)

費用を見積もる

Google Cloud 料金計算ツールを使用して、リソースベースの課金で、独自の同等のパイプラインのベースライン費用を見積もることができます。手順は次のとおりです。

  1. 料金計算ツールを開きます。
  2. [Add To Estimate] をクリックします。
  3. Dataflow を選択します。
  4. [サービスタイプ] で [Dataflow Classic] を選択します。
  5. [詳細設定] を選択して、オプションの完全なセットを表示します。
  6. ジョブを実行する場所を選択します。
  7. [サービスの種類] で [Streaming] を選択します。
  8. [Streaming Engine を有効にする] を選択します。
  9. ジョブの実行時間、ワーカーノード、ワーカーマシン、Persistent Disk ストレージの情報を入力します。
  10. Streaming Engine コンピューティング単位数の推定数を入力します。

リソースの使用量と費用は、入力スループットにほぼ比例して増加しますが、ワーカー数が少ない小規模なジョブでは、総費用は固定費が大部分を占めます。まず、ベンチマーク結果からワーカーノードの数とリソース消費量を推定できます。

たとえば、Kafka to Iceberg Direct(autosharding=false) アーキテクチャを使用して、入力データレートが 100 MBps のパイプラインを実行するとします。1 GBps パイプラインのベンチマーク結果に基づいて、リソース要件を次のように見積もることができます。

  • スケーリング ファクタ:(100 MBps)/(1024 MBps)= 約 0.1
  • 予測されるワーカーノード: 110 個のワーカー × 0.1 = 約 11 個のワーカー
  • 1 時間あたりの Streaming Engine コンピューティング単位数の予測数: 200 × 0.1 = 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 テーブルの識別子
  • CATALOG_NAME: Iceberg カタログの名前
  • CATALOG_TYPE: 使用するカタログのタイプ(hadoop、bigquery など)
  • BUCKET_NAME: ウェアハウスのロケーションの Cloud Storage バケットの名前
  • SCHEMA_DEFINITION: Kafka トピックデータのスキーマ定義({"type": "record", "name": "Record", "fields": [{"name": "id", "type": "string"}]} など)