Flink Adaptive Batch Processing Capability Evolution

Apache Flink はストリーミングおよびバッチコンピューティングフレームワークであり、初期は主にストリーミングコンピューティングのシナリオで使用されていました。近年、ストリームとバッチの統一という概念の推進に伴い、ますます多くの企業が Flink をバッチビジネスの処理に活用し始めています。

Flink はフレームワークレベルでバッチ処理をネイティブにサポートしていますが、実際の運用では依然として課題が存在します。そのため、近年のバージョンではコミュニティが Flink のバッチ処理能力の改善を継続的に進めており、API、実装、運用メンテナンスの 3 つの側面で反映されています。

API レベルでは、SQL の改善と構文の拡張を進め、Hive SQL との互換性を高めています。また、DataStream API も改善し、バッチジョブ開発をより適切にサポートしています。運用メンテナンスレベルでは、Flink バッチを本番環境でより使いやすくするため、History Server を改善して、実行中および完了後のジョブの状態をより明確に把握できるようにしました。同時に、Hive 互換の SQL Gateway も導入しています。実行レベルでは、オペレータの性能と安定性、実行計画、スケジューリング、ネットワーク層の改善を行いました。

核となるアイデアは、データ量、データパターン、実行時間、利用可能なリソースなどのランタイム情報に基づいて、ジョブの実行を適応的に最適化することです。具体的には、データ量に応じてジョブノードに適切な同時実行数を自動設定する、実行予測により遅いノードの影響を特定して軽減する、リソース使用率と処理性能を向上させる適応的データ転送を導入する、動的パーティションプルーニングを複数パーティションテーブルに適用して処理効率を向上させるといった取り組みが含まれます。

これらの改善により、Flink バッチ処理が使いやすくなり、バッチ処理ジョブの安定性が保証され、ジョブの実行性能が向上しています。一部の改善は複数の効果を同時に実現しています。

1. Adaptive Batch Scheduler

以前は、ジョブを本番環境に投入する前に同時実行数の調整が必要でした。バッチジョブのユーザーが直面する状況は以下のとおりです。

バッチ処理ジョブは非常に多いことが多く、ジョブの同時実行数の調整は大きなワークロードとなり、時間と労力を消費します。

日次データ量は変動することが多く、特にプロモーション期間中はデータが数倍から数十倍に増加するため、正確なデータ量の見積もりが困難で、チューニングも難しくなります。同時に、キャンペーン前後で同時実行数の設定を変更する場合も追加の工数がかかります。

Flink は複数のコンピューティングノードが直列に接続された実行トポロジーで構成されています。オペレータの複雑さとデータ自体の特性により、データ量と中間ノードの適切な同時実行数を正確に予測することは困難です。グローバルな統一同時実行数ではリソースの無駄遣いにつながり、追加のスケジューリング、デプロイメント、ネットワーク転送コストも発生します。

さらに、SQL ジョブではソースとシンク以外ではグローバルな統一同時実行数しか設定できず、細粒度な同時実行数を設定できないため、同様にリソースの無駄遣いと追加のオーバーヘッドが発生します。

この課題に対処するため、Flink はアダプティブバッチスケジューラを導入しました。ユーザーは各同時実行インスタンスが処理したいデータ量を設定できます。Flink は実行時に各論理ノードの実際のデータ量に基づいて、各論理ノードの実際の同時実行数を自動決定し、各同時実行インスタンスの処理データ量がユーザーの期待におおむね沿うようにします。

この設定方法の特徴は、データ量に依存しない汎用的な設定であることです。同じ設定を多くのジョブに適用でき、各ジョブを個別にチューニングする必要がありません。次に、自動設定された同時実行数は日次のデータ量の変動に対応できます。同時に、各ノードが実際に処理するデータ量を実行時に収集できるため、ノード粒度の同時実行数を設定し、より良い結果を得られます。

処理フローは上図のようになります。上流の論理ノード A のすべての実行ノードが実行を完了してデータを生成した後、総出力データ量を収集できます。これはノード B が消費すべきデータ量です。次に、この情報を同時実行数計算戦略コンポーネントに渡して適切な同時実行数を算出し、実行ノードトポロジーを動的に生成してスケジューリングとデプロイメントを行います。

従来の Flink 実行では、実行トポロジーは静的であり、ジョブの送信時にすべてのノードの同時実行数が既知でした。そのため、上流は実行中に下流の各実行ノードに対して個別のデータサブパーティションを分割でき、下流の起動時には対応するデータサブパーティションのみを読み取ればデータを取得できました。しかし、動的同時実行数の場合、上流の実行時には下流の同時実行数がまだ確定していないため、解決すべき主な課題は上流ノードの実行を下流ノードの同時実行数から切り離すことです。

動的実行トポロジーをサポートするため、以下の改善を行いました。上流ノードが生成するデータパーティション数は下流の同時実行数ではなく、下流の最大同時実行数によって決定されます。

上図の右側に示すように、下流の同時実行数が 4 の場合があり、A が生成するデータを 4 分割できます。実際に決定される下流の同時実行数は 1、2、3、4 のいずれかになり、各ノードに消費パーティション範囲が割り当てられます。たとえば、下流の同時実行数が 2 の場合はそれぞれ 2 つのデータパーティションを消費し、下流の同時実行数が 3 の場合は 1 つ消費するノードと 2 つ消費するノードが出ます。これにより、上流ノードの実行は下流の同時実行数に制約されなくなり、柔軟なデータ割り当てが可能になります。動的実行トポロジーの概念も実現できます。

自動同時実行数により 2 つの効果が得られます。まず、ユーザーが各ジョブの同時実行数を個別に設定する必要がなくなり、Flink バッチが使いやすくなります。次に、細粒度な同時実行数設定によりリソース使用率が向上し、意味のない高い同時実行数を回避できます。複数のクライアントで TPC-DS ベンチマークを使用してクラスターをテストした結果、アダプティブ同時実行数設定を有効にした後、総実行時間が 7.9% 短縮されました。

アダプティブバッチスケジューリングは後続の最適化の基盤にもなります。柔軟なデータパーティションと割り当て方法により、各データパーティションの実際のデータ量を収集できるため、データスキューによりパーティション間のデータ量に差がある場合、小規模なパーティションをマージして同じ下流ノードに処理させ、下流ノードの処理データをより均一にできます。

次に、動的実行トポロジーの導入により、実行時の情報に基づいてより適切な実行計画を動的に策定できます。たとえば、結合ノードの両端のデータサイズに応じて、どの結合モードを使用するか決定できます。

2. Specific Execution

本番環境におけるホットマシンの発生は避けられません。たとえば、混在クラスターでジョブを実行する場合やバッチジョブが集中的に実行される場合、一部のマシンが高負荷状態になり、そのノードで実行されるタスクが他のノードに比べて大幅に遅くなることがあります。これにより、ジョブ全体の実行時間が延び、ジョブの出力ベースラインが保証できなくなります。これは、一部のユーザーが Flink をバッチ処理に活用する上での障壁となっています。

そこで、Flink 1.16 で予測実行メカニズムを導入しました。

予測実行を有効にすると、Flink はバッチ処理ジョブ内で一部のタスクが他のタスクに比べて顕著に遅いことを検出した場合、そのタスクに対して新しい実行インスタンスを起動します。これらの新しい実行インスタンスは元の遅いタスクインスタンスと同じデータを消費して同じ結果を生成します。元の遅いタスクの実行インスタンスも保持されます。最初に完了したインスタンスがスケジューラによって唯一の完了インスタンスとして認識され、そのデータは下流によって検出されて消費されます。他のインスタンスはキャンセルされ、データはクリアされます。

予測実行を実装するため、Flink に以下のコンポーネントを導入しました。

Slow Task Detector は主に定期的に遅いタスクを検出して報告するために使用されます。

現在の実装では、論理ノードのある実行ノードが特に遅く、他の大半のノードの実行時間の中央値から一定のしきい値を超えた場合、遅いノードと判断されます。予測実行スケジューラは遅いノードを受け取り、遅いタスクが稼働しているマシンノードをホットマシンとして識別します。ブラックリストメカニズム (Blocklist Handler) を通じて、ホットマシンはブラックリストに登録され、後続のスケジューリングで新しいタスクがブラックリストに登録されたマシンに割り当てられなくなります。ブラックリストメカニズムは現在、Flink の最も一般的なデプロイメント方式である Yarn、K8s、スタンドアロンをサポートしています。

遅いノードの稼働中の実行インスタンス数が設定上限に達していない場合、上限まで予測実行インスタンスを起動し、ブラックリスト未登録のマシンにデプロイします。いずれかの実行インスタンスが終了すると、スケジューラは関連する他の実行インスタンスがまだ稼働中かどうかを確認し、稼働中であれば積極的にキャンセルします。

完了したインスタンスが生成したデータは下流に公開され、下流ノードのスケジューリングをトリガーします。

フレームワークレベルでは、Source ノードの予測実行をサポートしており、同じ Source の異なる同時実行インスタンスが常に同じデータを読み取れるようにします。基本的なアプローチは、キャッシュを導入して、各 Source の同時実行と各実行インスタンスが取得・処理済みのデータフラグメントを記録することです。ある実行インスタンスが現在割り当てられている Source のすべてのパーティションの処理を完了すると、新しいパーティションを要求でき、新しいパーティションもキャッシュに追加されます。

フレームワークレベルでの統合サポートにより、既存の Source のほとんどは追加修正なしで予測実行をサポートできます。新バージョンの Source を使用してカスタム SourceEvent を使用している場合にのみ、SourceEnumerator が追加のインターフェイスを実装する必要があります。実装しない場合、予測実行が有効なときに例外がスローされます。このインターフェイスは主に、ユーザー定義イベントが正しい実行インスタンスに確実に配信されるようにするためのものです。予測実行が有効な場合、複数の実行インスタンスが同時に稼働する可能性があるためです。

REST と Web UI レベルでも予測実行をサポートしています。予測実行が発生した場合、ジョブノードの詳細画面で予測実行のすべての同時実行インスタンスを確認できます。同時に、リソース概要カードではブラックリストに登録された TaskManager の数と、占有されていないがブラックリストに登録されて使用できないスロットの数も確認できます。ユーザーはこれを利用して現在のリソース使用状況を判断できます。さらに、TaskManager 画面では現在ブラックリストに登録されている TaskManager を確認できます。

現在のバージョンでは、Sink の予測実行はサポートされていません。今後は Sink ノードの予測実行を優先的にサポートする予定です。解決すべき課題は、各 Sink が確実に 1 つのデータのみをコミットし、キャンセルされた Sink が生成した他のデータをクリーンアップできるようにすることです。

また、遅いタスク検出戦略のさらなる改善も計画しています。現在、データスキューが発生すると、個々の同時実行のデータ量が他の同時実行より多くなり、実行時間も他のノードより長くなりますが、そのノードが遅いタスクとは限りません。そのため、この状況を正しく識別して処理し、無効な予測実行インスタンスの起動によるリソースの浪費を回避する必要があります。現在の初期アイデアは、各実行インスタンスが実際に処理したデータ量に基づいてタスクの実行時間を正規化することで、これは前述のアダプティブバッチスケジューラが各ノードの生産データ量を集計する機能に依存しています。

3. Hybrid Shuffle

Flink には主に 2 つのデータ交換方法があります。

ストリーミングパイプラインシャッフル:上流と下流が同時に起動し、データを空中転送するため、ディスクへの書き込みが不要で、性能的に優れています。ただし、大量のリソースを必要とし、ジョブは単一ノードの同時実行数の数倍のリソースを同時に確保する必要があり、本番バッチジョブでは対応が困難です。同時に、バッチリソースの需要により、リソースを同時に確保できないとジョブを実行できません。複数のジョブが同時にリソースを奪い合うと、リソースデッドロックが発生する可能性があります。

バッチブロッキングシャッフル:データは直接ダウンロードされ、下流は上流のデータを直接読み込みます。この交換モードではジョブがリソースにより柔軟に対応できます。理論上は上流と下流を同時に実行する必要がなく、スロットが 1 つあればジョブ全体を実行できます。ただし、性能面では劣り、下流ステージは上流ステージの完了後にしか実行できず、すべてのデータがダウンロードされるため IO オーバーヘッドが発生します。

そこで、両者の利点を組み合わせたシャッフルモードを実現したいと考えました。リソースが十分な場合はストリーミングシャッフルの性能上の優位性を発揮でき、リソースが限られている場合でもバッチシャッフルのリソース適応能力を持ち、スロットが 1 つしかなくても実行できます。同時に、リソース適応能力は自動的に切り替えられるため、ユーザーは意識する必要がなく、個別のチューニングも不要です。

そのために、ハイブリッドシャッフルを導入しました。

現在のハイブリッドシャッフルはデフォルトスケジューラに基づいて実装されているため、自動同時実行数設定や予測実行とは互換性がありません。バッチ処理をより良くサポートするには、統合が必要です。

4. Dynamic Partition Pruning

オプティマイザの重要な役割の一つは、無効で冗長な計算を回避することです。パーティションテーブルは本番環境で広く使用されています。ここでは、パーティションテーブルにおける無効なパーティションの読み取りを削減する方法を紹介します。

TPC-DS モデルの簡略化された例を使ってこの最適化を紹介します。上図に示すように、販売テーブルには SLOD_Date という名前のパーティションフィールドがあり、合計 2000 個のパーティションがあります。

上図の SQL ステートメントは、販売テーブルから SLOD_Date = 2 のデータを選択しています。パーティションプルーニングがない場合、すべてのパーティションデータを読み取ってからフィルタリングする必要があります。静的パーティションプルーニングでは、最適化フェーズでのフィルタープッシュダウンなどの最適化により、決定されたパーティションをスキャンノードに通知できます。スキャンの実行時に特定のパーティションのみを読み取ればよいため、読み取り IO が大幅に削減され、ジョブの実行が高速化されます。

上図には 2 つのテーブルがあります。ファクトテーブルの販売テーブルとディメンションテーブルの date_Dim で、これら 2 つのテーブルを結合します。フィルター条件がディメンションテーブルに適用されるため、静的パーティションプルーニングの最適化は実行できません。

ディメンションテーブルの date_Dim はすべてのデータを読み取ってフィルタリングを行います。ファクトテーブルの販売テーブルはすべてのパーティションを読み取ってから結合を行います。year = 2000 および SLOD_date = date_sk の関連データのみが出力されます。多くのパーティションデータが無効であることが推測できますが、これらのパーティションは静的最適化フェーズでは分析できず、実行フェーズでディメンションテーブルのデータに基づいて動的に分析する必要があります。これが動的パーティションプルーニングです。

動的パーティションプルーニングの処理フローは以下のとおりです。

ステップ 1:結合するディメンションテーブルのオペレータを実行します。例:Scan (date_dim) と Filter。

ステップ 2:ステップ 1 のフィルター結果をパーティションテーブルのオペレータ Scan に送信します。

ステップ 3:ステップ 2 のデータに基づいて無効なパーティションを除外し、有効なデータのみを読み取ります。

ステップ 4:ステップ 1 とステップ 3 の結果に基づいて結合を完了します。

動的パーティションプルーニングと静的パーティションプルーニングの違いは、動的パーティションプルーニングでは最適化フェーズでどのパーティションデータが有効かを特定できず、ジョブ実行後に初めて決定される点です。

Flink で動的パーティションプルーニングを実装する手順は以下のとおりです。

まず、Physical Plan に DynamicFilterDataCollector (以下 DataCollector) という特別なノードを追加して、フィルターデータを収集し、重複排除を行います。関連フィールドのみを保持してパーティションテーブル Scan に送信します。パーティションテーブル Scan はデータを取得した後、パーティションプルーニングを実行して最終的に結合を完了します。Streaming Graph では、Source オペレータ (Scan ノードに対応) には入力がありませんが、DataCollector オペレータからのデータを受信できるようにする必要があります。また、ディメンションテーブル側の date_Dim Scan と year = 2000 のフィルターには、右側の販売テーブル Scan のスケジューリングに対する依存関係がないため、右側のオペレータが先に実行され、左側のオペレータが遅れて実行される可能性があり、動的パーティションプルーニングの最適化が完了しません。

そのため、OrderEnforce ノードを導入してスケジューラにデータ依存関係を通知し、左側のオペレータが先にスケジュールされるようにして、動的パーティションプルーニングの最適化が正しく実行されるようにしました。

今後は、フレームワークレベルで上記のスケジューリング依存関係の問題を解決し、Streaming Graph をより簡潔にする計画です。

左側の date_Dim Scan と Filter が先に実行され、データが DataCollector と Join に送信されます。Source オペレータに入力がない問題を解決するために、Flink のコーディネーターメカニズムを使用します。DataCollector はデータを収集した後、コーディネーターに送信して無効なパーティションのプルーニングを完了します。パーティションテーブル Scan はコーディネーターから有効なパーティション情報を取得します。販売テーブル Scan ノードが実行された後、最終的な結合が実行されます。

DataCollector と OrderEnforce の間にはデータサイドがあり、実際のデータ転送は行われず、スケジューラに DataCollector が OrderEnforce の前に呼び出されることを通知するためにのみ使用されます。

上図は TPC-DS 10T データセットに基づく最適化前後の性能比較を示しています。青は非パーティションテーブル、赤はパーティションテーブルです。最適化後の時間節約は約 30% です。

Q&A

Q:ホットマシンで遅いタスクが発生した場合、他のマシンにインスタンスを起動して遅いタスクを再実行しますが、再度起動したインスタンスでもホットスポットが発生する可能性はありますか?

A:理論的には可能性があります。予測実行自体がリソースを時間に変換する戦略だからです。ただし、本番環境での実践により、このメカニズムは効果的であることが証明されています。追加のリソースコストとさらなるホットスポットの発生を考慮しても、トレードオフとして依然としてコスト対効果が高いと言えます。

Q:予測実行はデータ量に基づいていますか?

A:現在の戦略は実行時間に基づいています。たとえば、大半のタスクの中央実行時間が 1 分の場合、あるタスクが 1.5 分以上実行されていれば、遅いタスクと判断されます。具体的な値は設定可能です。

Q:遅いノードの検出は設定可能ですか?

A:このポリシーは現在ハードコードされており、設定による変更はまだサポートされていません。戦略が安定した後、将来的にはユーザーに開放し、二次開発やプラグインを通じて遅いタスク検出ポリシーを変更できるようにする可能性があります。

Q:予測実行メカニズムは DataStream API と Flink SQL の両方をサポートしていますか?

A:はい、両方サポートしています。

Q:ブラックリストメカニズムについて、過度なブラックリスト登録や登録の遅れによるリソースの浪費は発生しませんか?

A:現在、予測実行下でのブラックリストは保守的であり、デフォルトでは 1 分間のみブラックリストに登録されます。ただし、遅いタスクが継続して発生する場合、ブラックリスト期間は継続的に更新されるため、遅いタスクが存在するマシンノードは常にブラックリストに登録された状態になります。

Related Articles

Explore More Special Offers

  1. Short Message Service(SMS) & Mail Service

    50,000 email package starts as low as USD 1.99, 120 short messages start at only USD 1.00

phone お問い合わせ
Hi, I'm Alibaba Cloud AI Assistant!
I can help with questions and solutions.