Dataflow での柔軟性と効率性を考慮した設計

このドキュメントでは、復元性に優れた Dataflow パイプラインを構築するための一連のアーキテクチャのベスト プラクティスである、柔軟性と効率性のための設計(DFE)について説明します。

厳格なインフラストラクチャの制約から柔軟なリソース定義に移行すると、次のことが可能になります。

  • コンピューティング リソースの取得可能性を最大化します。
  • リージョンの需要が高い期間にシームレスな自動スケーリングを確保します。
  • パイプラインのリリース遅延を防ぎ、容量のボトルネックを解消します。

たとえば、パイプラインを 1 つのゾーンの 1 つの特定のマシンタイプ(us-central1-a の n1-standard-4 ワーカーが必要など)に制限する代わりに、最小リソース要件(4 個の vCPU と 16 GB の RAM など)を設定できます。us-central1-a または N1 マシンシリーズで一時的な容量制約が発生した場合、Dataflow は他のゾーンとマシン ファミリー(E2、N2、N2D など)で互換性のあるワーカー VM を自動的にプロビジョニングできます。この柔軟性により、単一の制約されたハードウェア プールを待つことなく、パイプラインを開始してスケーリングできます。

このドキュメントは、Dataflow ワークロードを管理し、パイプラインの信頼性、スループット、インフラストラクチャの可用性を最適化したいデータ エンジニア、クラウド アーキテクト、プラットフォーム管理者を対象としています。

DFE の概要

Dataflow は、Apache Beam パイプラインを実行するために Compute Engine 仮想マシン(VM)インスタンスを動的にプロビジョニングする、フルマネージドのサーバーレス データ処理サービスです。大規模なバッチ処理とストリーミング パイプラインでは、ワーカープールが数十台から数百台の VM インスタンスにスケールアップされることがよくあります。

インフラストラクチャの制約が厳しいパイプラインは、需要が高い期間にプロビジョニングの遅延が発生しやすくなります。厳格な制約の例を次に示します。

  • n1-standard-4 などの単一のマシンタイプをハードコードする。
  • パイプラインを特定の Compute Engine ゾーンに固定する。

その特定のマシンタイプまたはゾーンで一時的に需要が高くなると、Dataflow はコンピューティング リソースを割り当てることができません。これにより、プロビジョニングの遅延や ZONE_RESOURCE_POOL_EXHAUSTED や RESOURCE_POOL_EXHAUSTED などのエラーが発生する可能性があります。

DFE の原則を使用すると、パイプライン アーキテクチャを、柔軟で要件ベースのリソース定義に移行できます。この柔軟性により、Dataflow は Google Cloudで利用可能な多様なハードウェア プールにコンピューティングを動的に分散できます。これにより、運用オーバーヘッドを最小限に抑えながら、コンピューティングの取得可能性を最大化できます。

DFE のベスト プラクティス

次のベスト プラクティスを採用して、コンピューティングの取得可能性を最大化し、自動スケーリングの応答性を高め、復元力のあるパイプラインを構築します。

VM の自動選択を有効にする

ワーカー マシンタイプ パイプライン オプションを使用して静的マシンタイプをハードコードするのではなく、Apache Beam リソースのヒントで自動 VM 選択を使用します。最小リソース要件(min_ram または cpu_count)を指定すると、Dataflow はインスタンスの柔軟性を自動的に有効にし、互換性のあるマシンタイプのリストからワーカーをプロビジョニングします。

ワークロードのサポート:

  • バッチ パイプライン: リソースヒントを指定すると、Right Fitting と自動 VM 選択が自動的に有効になります。
  • ストリーミング パイプライン: Right Fitting では、--experiments=enable_streaming_rightfitting パイプライン オプションと、水平自動スケーリング(デフォルトで有効)および Streaming Engine(--enable_streaming_engine)を設定する必要があります。

自動 VM 選択を構成するには、コマンドライン オプション、SDK パイプライン オプション、または Flex テンプレート実行パラメータを使用して、パイプライン レベルで最小リソース要件(min_ram または cpu_count)を指定します。Java と Python の詳細な設定手順とコード例については、リソース ヒントを使用するをご覧ください。

リージョン ワーカーの配置を使用する(ゾーンの固定を回避する)

選択したリージョン内の正常なゾーン全体でワーカー VM を動的にスケジュールするように Dataflow を構成します。

--region パイプライン オプションを指定し、--zone と --worker_zone を省略します。次に例を示します。

--region=us-central1

マネージド サービスを使用して状態とシャッフルを分離する

マネージド バックエンド サービスを使用しないパイプラインは、シャッフル データ オペレーションとストリーミング状態の保存をワーカー VM のディスクとメモリで直接実行します。この密結合では、ワーカー ディスクのサイズが大きくなり、ワークロードの存続が特定の VM インスタンスにバインドされるため、容量制約時にワーカーの置き換えが難しくなります。

  • バッチジョブの場合 - Dataflow Shuffle を使用する: Dataflow Shuffle は、サポートされているワーカー マシンタイプで実行されるバッチ パイプラインでデフォルトで有効になっており、シャッフル オペレーションをワーカー VM から専用の Google 管理のバックエンド サービスにオフロードします。
  • ストリーミング ジョブの場合 - Streaming Engine を使用します。Streaming Engine は、ウィンドウ状態ストレージとタイマー管理をワーカー VM から特殊な高応答性のバックエンド インフラストラクチャにオフロードします。Apache Beam SDK 2.30.0 以降を使用するパイプラインの場合、Streaming Engine はデフォルトで有効になっています。明示的に有効にするには、--enable_streaming_engine パイプライン オプションを渡します。

バッチ パイプラインに柔軟なリソース スケジューリング(FlexRS)を使用する

夜間の ETL、データレイクの取り込み、日次ロールアップなど、時間の制約が厳しくないバッチ ワークロードには、Flexible Resource Scheduling(FlexRS)を使用します。

FlexRS を有効にするには、flexRS の目標パイプライン オプションを設定します。

  • Python パイプラインの場合: --flexrs_goal=COST_OPTIMIZED
  • Java パイプラインの場合: --flexRSGoal=COST_OPTIMIZED

Flex テンプレートの柔軟なランチャー VM タイプを構成する

Flex テンプレートを使用してパイプラインを起動する場合、パイプライン ランチャー VM のデフォルトは e2-standard-2 です。ほとんどの場合、デフォルトの VM で問題ありませんが、容量の制約が発生した場合は、gcloud dataflow flex-template run コマンドの実行時に --launcher-machine-type オプションを使用して構成をカスタマイズできます。

gcloud dataflow flex-template run my-job \
    --template-file-gcs-location="gs://my-bucket/template.json" \
    --region="us-central1" \
    --launcher-machine-type="n2-standard-2"

運用上の考慮事項とトレードオフ

DFE のベスト プラクティスを採用すると、コンピューティングの取得可能性、自動スケーリングの応答性、運用の信頼性が大幅に向上しますが、アーキテクチャを設計する際は、次の運用上の要因とトレードオフを考慮してください。

VM の自動選択に関する考慮事項

  • 信頼性とピーク パフォーマンス: 自動 VM 選択では、ピーク時の実行パフォーマンスよりもジョブの起動の信頼性とコンピューティングの取得可能性が優先されます。Dataflow は複数の候補マシン ファミリー(E2、N2、N4、N2D など)からプロビジョニングするため、プロビジョニングされるマシン ファミリーに応じて、ランタイム パフォーマンスとスループットが若干異なる場合があります。厳格な実行 SLA を伴うコンピューティング負荷の高いワークロードの場合は、パイプラインを広範囲にデプロイする前に、自動 VM 選択でパイプラインをテストして、パフォーマンスのベースラインを確立します。ワークロードに特定のハードウェア プラットフォームまたはクロック速度が必要で、容量の制約を許容できる場合は、特定のマシンタイプを引き続き設定できます。
  • 候補ファミリー全体の Compute Engine 割り当て: 自動 VM 選択では、複数の候補マシン ファミリーからワーカーをプロビジョニングできるため、 Google Cloud プロジェクトに、ターゲット リージョンの各候補ファミリーに十分な Compute Engine vCPU とメモリの割り当てがあることを確認してください。プライマリ ファミリーで容量不足が発生し、フォールバック ファミリーの割り当てがプロジェクトにない場合、ワーカーのプロビジョニングは QUOTA_EXCEEDED エラーで失敗します。
  • ストリーミング パイプラインの前提条件: ストリーミング パイプラインでは、Right Fitting と Auto VM Selection はデフォルトで有効になっていません。--experiments=enable_streaming_rightfitting を明示的に指定し、Streaming Engine(--enable_streaming_engine)と水平自動スケーリングの両方が有効になっていることを確認する必要があります。
  • 構成の除外: 次の表の機能またはオプションを構成すると、自動 VM 選択は自動的にバイパスされるか、サポートされません。

    機能 フラグまたは構成オプション メモ
    明示的なマシンタイプ --worker_machine_type または --machine_type(Python)
    --workerMachineType(Java)
    指定されたマシンタイプが優先され、自動 VM 選択はバイパスされます。
    カスタム ディスクタイプ、プロビジョニングされた IOPS、またはスループット --disk_type、--disk_provisioned_iops、または --disk_provisioned_throughput_mibps VM の自動選択がバイパスされます。--disk_size_gb を使用したカスタム ディスクサイズの設定がサポートされています。
    最小 CPU プラットフォーム --min_cpu_platform(Python)
    --minCpuPlatform(Java)
    最小 CPU プラットフォームを設定すると、VM の自動選択がバイパスされます。
    Confidential VMs --experiments=enable_confidential_compute Confidential VM インスタンスは、VM の自動選択ではサポートされていません。
    GPU または TPU アクセラレータ --dataflow_service_options=worker_accelerator=... または accelerator リソースヒント VM の自動選択は、アクセラレータのないワークロードにのみ適用されます。
    Dataflow Prime --dataflow_service_options=enable_prime Dataflow Prime は、自動 VM 選択ではなく、垂直自動スケーリングと動的ライト フィッティングを使用します。
    柔軟なリソース スケジューリング(FlexRS) --flexrs_goal=COST_OPTIMIZED(Python)
    --flexRSGoal=COST_OPTIMIZED(Java)
    FlexRS は独自のワーカープールとスケジューリング バッファを管理します。

Flexible Resource Scheduling(FlexRS)のトレードオフ

  • スケジューリング遅延ウィンドウ: FlexRS では、ジョブの実行が開始される前に最大 6 時間のスケジューリング バッファを導入できます。完了までの時間に関する厳格な SLA や、下流の依存関係が厳しいパイプラインには FlexRS を使用しないでください。

リージョン プレースメントとデータ局所性

  • マネージド サービスの前提条件: リージョン ワーカーの配置は、バッチに Dataflow Shuffle を使用するジョブ、またはストリーミングに Streaming Engine を使用するジョブでのみサポートされます。これらのマネージド バックエンド サービスを使用しないジョブは、リージョン内の最適な単一ゾーンを選択する自動ゾーン配置を使用します。
  • データの局所性とクロスリージョン下り(外向き): リージョン配置では、選択したリージョン内の使用可能なゾーンにワーカーが分散されます。ネットワークのレイテンシを最小限に抑え、リージョン間のネットワーク下り(外向き)料金を回避するには、すべてのデータソースとシンク(Cloud Storage バケット、BigQuery データセット、Pub/Sub トピックなど)が Dataflow ジョブと同じリージョンに存在することを確認します。

Compute Engine の予約

  • 予約アフィニティ: オンデマンド Dataflow ジョブは、ANY 予約アフィニティを使用する一致する Compute Engine 予約を自動的に使用します。ただし、自動 VM 選択では、特定の名前付き予約からインスタンスを使用することはできません。
  • 一時的なワークロードへの適合性: Compute Engine の予約は、一般的に、急増しやすいワークロード、自動スケーリング ワークロード、短期間のバッチ ワークロードにはおすすめできません。また、ゾーンの容量不足がアクティブなときに新しい予約を作成すると、オンデマンド VM の作成と同じ容量制約で失敗します。

次のステップ