Dataflow の並列処理について

Dataflow は、コンピューティング インスタンスのマネージド プールに処理を分散して、大規模なデータ処理パイプラインを実行するように設計されています。Dataflow が処理を並列化する方法を理解すると、効率的なパイプラインを設計し、パフォーマンスのボトルネックを回避し、リソース費用を最適化できます。

このページでは、Dataflow がデータ処理を並列化する方法、実行を管理およびスケーリングする方法、並列処理を制約する一般的な要因、パイプライン スループットを最適化するために使用できる手法について説明します。

並列処理モデル: 水平と垂直

Dataflow は、2 つの補完的な戦略を使用して並列処理を実現します。

  • 水平並列処理: パイプライン データが分割され、複数のワーカー インスタンス(仮想マシン)で同時に処理されます。Dataflow は、水平自動スケーリングにより、ワークロードの需要に基づいてワーカー プールのサイズを自動的に調整できます。デフォルトでは、Dataflow はジョブあたり 4,000 個のワーカーのリソース上限を設定します。この上限は、割り当てリクエストを使用して調整できます。

  • 垂直並列処理: 単一のワーカー インスタンス内の複数の CPU コアとスレッドが、パイプライン データを同時に処理します。各ワーカー VM は、ワーカー プロセスとハーネス スレッドを実行して、使用可能なコンピューティング リソースを使用します。Dynamic Thread Scaling を使用すると、Dataflow は CPU 使用率とメモリ ヘッドルームに基づいて、バッチ パイプラインのワーカーあたりのアクティブ スレッド数を調整できます。Dataflow Prime では、垂直自動スケーリングにより、ワーカーに割り当てられたメモリとコンピューティングが動的にスケーリングされます。

作業単位と実行階層

ワーカーとスレッド間で処理を分散するために、Dataflow は Apache Beam パイプラインを個別の作業単位に分割します。

  • PCollection とパーティション: PCollection は分散データセットを表します。境界のあるデータ(バッチ パイプライン)の場合、Dataflow はデータセットを分割またはシャードに分割します。制限なしデータ(ストリーミング パイプライン)の場合、データは継続的に到着し、メッセージまたはストリーム パーティションとして取り込まれます。
  • バンドル: Dataflow は、DoFn による処理のために、要素を任意のバンドルにグループ化します。バンドルは失敗と再試行の単位です。要素の処理で未処理の例外が発生すると、バンドル全体が再試行されます。メモリ使用量の多いオペレーションは、ワーカーのメモリ負荷を増やし、メモリ不足エラーにつながる可能性があります。
  • ステージとステップの融合: グラフの最適化中に、Dataflow は隣接する変換を融合された実行ステージに結合し、中間データの具体化のオーバーヘッドを排除します。融合ステージ内では、要素は単一のスレッド上のタイトな実行ループで処理され、次のステージまたはシャッフル境界に渡されます。

パイプラインの変換とグラフの生成の詳細については、パイプラインのライフサイクルをご覧ください。

マネージド並列処理と自動スケーリング

デフォルトでは、Dataflow は次の方法で、手動でパーティションをチューニングすることなく、パイプラインの並列処理を自動的に管理します。

  • 水平自動スケーリング:
    • バッチ パイプライン: 合計推定残作業量、ソース バックログ、CPU 使用率を評価して、ワーカープールをスケールアップまたはスケールダウンし、ジョブを迅速かつ費用対効果の高い方法で完了します。
    • ストリーミング パイプライン: システム レイテンシ、バックログ サイズ、CPU 使用率を分析して、スループットの急増時にワーカーをスケールアップし、トラフィックの少ない期間にスケールダウンします。詳細については、ストリーミングの水平自動スケーリングを調整するをご覧ください。
  • 動的作業再調整(DWR): バッチ パイプラインでは、Dataflow は個々のワーカータスクの進行状況をモニタリングします。ワーカーが早期に終了した場合や、データスキュー(ストラグラー)により別のワーカーが遅延した場合、Dataflow は低速ワーカーの未処理の残りの作業を動的に分割し、アイドル状態のワーカーに再割り当てします。詳細については、動的作業再調整をご覧ください。
  • Dynamic Thread Scaling: Portable Runner を使用するバッチ パイプラインで、CPU 使用率とメモリ ヘッドルームに基づいて、ワーカーあたりの同時処理スレッド数を自動的に調整します。詳細については、Dynamic Thread Scalingをご覧ください。
  • 垂直自動スケーリング: Dataflow Prime では、Dataflow はワーカーのメモリとコンピューティング リソースを動的にスケーリングして、メモリエラーを防ぎ、リソース使用率を最適化します。詳細については、垂直自動スケーリングをご覧ください。

並列処理を制約する要因

次のデータ特性またはパイプライン グラフの設計により、パイプラインで期待される並列処理が実現されないことがあります。

分割できない入力ソース

入力ソースを独立した範囲に分割できない場合、Dataflow は単一のワーカー スレッドでソースを順次読み取る必要があります。

  • 分割できないファイルの圧縮: .gz(gzip)や .bzip2(インデックスなし)などの形式は、任意のバイト オフセットから並列で読み取ることができません。1 つの大きな圧縮ファイルを読み取ると、データが解凍されて再配布されるまで、取り込みステージが単一スレッドに制限されます。
  • 解決策: 分割可能なファイル形式(Parquet、Avro、Snappy 圧縮形式など)でデータを保存するか、Cloud Storage で入力データを複数の小さなファイルに分割します。

ステップ融合と高いファンアウト

ステップ融合は、シリアル化のオーバーヘッドを減らすことでパフォーマンスを向上させますが、並列処理の低いステップで多数の出力要素(「ハイ ファンアウト」オペレーション)が生成されると、並列処理が意図せず制限され、メモリ負荷が増加する可能性があります。

  • 例: ソースが 5 つのファイルを読み取り、1,000,000 個の出力要素を生成する FlatMap 変換と融合されます。FlatMap 変換がダウンストリーム変換と融合されると、1,000,000 個の要素すべてが最大 5 つのワーカー スレッドで実行され続け、ダウンストリーム スループットが大幅に制限されます。また、中間変換がコミット前にメモリ内で大幅に拡張されると、大きなバンドルによって使用可能なワーカーメモリが使い果たされる可能性があります。
  • 解決策: 高いファンアウト ステップとダウンストリーム変換の間に Redistribute(または従来の Reshuffle)変換を挿入して、フュージョンを中断し、ワーカープール全体に作業を再分配します。メモリ関連の問題のデバッグについては、メモリ不足エラーのトラブルシューティングをご覧ください。

キーのスキューとホットキー

集約オペレーション(GroupByKey、CoGroupByKey、Combine.PerKey)は、関連付けられたキーで要素をグループ化します。

  • ホットキーのボトルネック: Dataflow は、同じキーを持つすべての要素を集計のために単一のワーカー スレッドに転送します。1 つのキーにデータセット全体の大部分が含まれている場合、そのワーカーは遅延ワーカーになり、アップストリーム ワーカーでバックプレッシャーが発生する可能性があります。たとえば、デフォルトの null キーや非常に人気のあるカテゴリキーなどです。
  • 解決策:
    1. 可能な場合は GroupByKey ではなく Combiner(CombineFn または Combine.PerKey)を使用します。これにより、Dataflow はシャッフル前に部分的なローカル結合を実行できます。
    2. ホットキーにランダムな整数接頭辞または接尾辞を追加して(キーのソルト化)、キー空間をワーカーに分散し、2 段階目の集計でソルト化された結果を統合します。

ダウンストリーム シンクのスロットリング

パイプライン出力をデータベースやサードパーティ API などの外部サービスに書き込む場合、並列処理のレベルが高いと、宛先システムが飽和する可能性があります。

  • スロットリング: 同時書き込み呼び出しを発行する数百のワーカー スレッドは、レート制限エラー、接続タイムアウト、データベースのパフォーマンス低下につながる可能性があります。
  • 解決策:
    • GroupByKey で要素をグループ化するか、並列処理を制御したバッチ シンクを使用して、書き込みの並列処理を制限します。
    • シンク DoFn 実装にクライアントサイドの指数バックオフと再試行ロジックを実装します。

最適化戦略

Dataflow ジョブの並列処理を最適化するには、次の方法を検討してください。

次のステップ