From Spark for batch processing to Flink for streaming and batch integration

1. ストリームバッチ統合が必要な理由

ストリームとバッチを統合するメリットは何か。
特に BI/AI/ETL の文脈において、ストリームバッチ統合にはどのような利点があるのでしょうか。
ユーザーがストリームとバッチの統合を実現できれば、以下の 4 つの明確なメリットが得られます。


コードの重複を回避し、コア処理ロジックを再利用できる
コードロジックを完全に一致させることができれば理想的ですが、実際には難しい面があります。
現在のビジネスロジックはますます長く複雑になり、要件も多岐にわたります。
異なるフレームワークやエンジンを使用すると、ユーザーは毎回ロジックを書き直す必要があり、大きな負担となり保守も困難になります。
そのため、コードの重複をできる限り回避し、ユーザーがコードロジックを再利用できるようにすることは特に重要です。

ストリームバッチ統合には 2 つの方向性があります

これら 2 つの方向性で考慮すべき点は大きく異なります。
現在、ストリーム処理では Apache Flink、バッチ処理では Apache Spark といったフレームワークが比較的成熟しており、多くのユーザーを獲得しています。
ユーザーを別の方向へ移行させる場合、たとえばストリーム処理からバッチ処理へ移行するか、バッチ処理からストリーム処理へ移行するかという 2 つのカテゴリに分けられます。
後述する 2 つの本番事例は、それぞれこの 2 つの方向性に対応しています。


メンテナンス作業量を削減する

複数のシステムを保守する必要をなくします。
システム間の差異は大きく、フレームワークやエンジンが異なるとより多くの問題を引き起こします。
社内にリアルタイム用とオフライン用の複数のパイプラインが存在すると、データの不整合が発生します。
そのため、データ整合性を可能な限り保つために、データ検証、データ精度確認、データストレージなどで多くの作業が発生します。

詳細を見る
多くのフレームワークやエンジンがあり、ビジネスロジックはリアルタイムとオフラインの両方で実行される必要があります。
そのため、ユーザーをサポートする上で習得すべき事項が多くなります。


2. 業界の現状

Apache Flink と Apache Spark はいずれもストリーム処理とバッチ処理の両方をサポートするエンジンです。
Apache Flink のストリーム処理が優れていることは認識していますが、そのバッチ処理はどの程度のレベルに達しているのでしょうか。
同時に、Apache Spark のバッチ処理は比較的優れていますが、そのストリーム処理はユーザーの既存の要件を十分に解決できるのでしょうか。


現在、さまざまなエンジンフレームワークが存在しますが、それらの上に統一されたフレームワークや、Beam API やカスタムインターフェイスのようなシンプルな物理 API を構築できないでしょうか。


Beam で考慮すべき点は、バッチ処理とストリーム処理をどの程度最適化できるかです。
Beam は現在まだ物理実装寄りの設計となっており、今後の計画を確認する必要があります。


LinkedIn を含む各社は、共通の SQL レイヤーや共通の SQL/API レイヤーを設け、その下で異なるフレームワークエンジンを実行するカスタムインターフェイスソリューションを検討しています。
ここで考慮すべき点は、Apache Spark や Apache Flink は既に比較的成熟しており大規模なユーザー基盤を持っているため、新しい API やソリューションを提案する際のユーザーの受け入れ度や、社内での保守方法を検討する必要があります。


3. 本番環境での適用事例

以下の内容では、主に Apache Flink のバッチ処理としての効果、Apache Flink と Apache Spark の簡単な比較、および LinkedIn の内部ソリューションに焦点を当てます。
2 つの本番事例を紹介します。
1 つは機械学習の特徴量エンジニアリングにおけるストリームバッチ統合、もう 1 つは複雑な ETL データフローにおけるストリームバッチ統合です。


3.1 ケース A - 機械学習の特徴量エンジニアリング

1 つ目の方向性、ストリーム処理 → バッチ処理は、ストリームバッチ統合に分類されます。


ケース A の主なロジックは、機械学習の特徴量生成時にストリーム処理からバッチ処理へどのようにストリームとバッチを統合するかです。
コアビジネスロジックは特徴量変換です。
変換プロセスとロジックが複雑であるため、標準化の好例として取り上げます。


たとえば、LinkedIn のページで入力されたメンバー情報の背景などを抽出して標準化し、求人推薦などに活用します。
メンバーの ID 情報が更新されると、Kafka からの読み取りを含むフィルタリングと前処理ロジックが実行され、特徴量変換プロセス中に小規模なテーブルクエリが発生する場合があります。
このロジックは非常にシンプルで、複雑な結合操作やその他のデータ処理プロセスは伴いません。


以前のパイプラインはリアルタイムで、オフラインパイプラインから補足情報を読み取ってストリームを定期的に更新する必要がありました。
この種のバックフィルはリアルタイムクラスターに大きな負荷をかけます。
バックフィル中はバックフィルの完了を待つ必要があり、リアルタイムクラスターが停止しないようワークフローを監視する必要があります。
そのためユーザーは、オフラインでバックフィルを行うことはできないか、リアルタイムストリーム処理によるバックフィルは避けたいと要望しました。


現在、ユーザーはストリーム処理に Beam on Apache Samza を使用しています。
ユーザーは Beam API と Spark Dataset API に非常に精通しており、バックフィル以外にも Dataset API を使用して他のビジネス処理を行っています。


強調すべき点は、多くの Dataset API はオブジェクトを直接操作し、型安全性に対する要件が高いことです。
これらのユーザーに SQL や DataFrame などのワークフローへ直接移行するよう提案するのは非現実的です。
既存のビジネスロジックがオブジェクトの直接操作と変換に基づいているからです。


このケースでは、ユーザーに命令型 API の選択肢を提供できます。
業界が提供するソリューションを見てみましょう。


最初の選択肢は、統一されつつある Apache Flink DataStream API です。
プログラム評価時には Apache Flink DataSet API (非推奨) も調査しました。
DataStream API は統一でき、ストリーム処理とバッチ処理のサポートも比較的充実しています。
ただし、命令型 API であるため最適化の余地は限られており、今後も最適化が続けられる予定です。
FLIP-131: Consolidate the user-facing Dataflow SDKs/APIs (and deprecate the DataSet API) および FLIP-134: Batch execution for the DataStream API を参照してください。


2 番目の選択肢は Spark Dataset で、ユーザーにとって自然な選択肢でもあります。
Dataset API はストリーミングにも使用できますが、これは Apache Flink の Dataset や DataStream API などの物理 API とは異なります。
Spark Dataframe SQL エンジンをベースに型安全性を実現しており、最適化の面で比較的優れています。
Databricks: Introducing Apache Spark Datasets および Spark Structured Streaming Programming Guide: Unsupported-operations の記事を参照してください。


3 番目の選択肢は Beam On Spark で、現在は主に RDD ランナーを使用しています。
最適化対応ランナーのサポートはまだ困難です。
Beam の進行中の作業については、後述するケース B で詳しく説明します。
Beam Documentation - Using the Apache Spark Runner および BEAM-8470 Create a new Spark runner based on Spark Structured streaming framework を参照してください。


ユーザーのフィードバックによると、Apache Flink の DataStream (DataSet) API と Spark の Dataset API はユーザーインターフェイスの面で非常に近いとのことです。
インフラエンジニアとしてユーザーの問題解決を支援するには、API の習熟度の方が重要です。


ただし、Beam の API は Apache Flink や Apache Spark と大きく異なります。
Beam は Google のエコシステムに属しています。
以前、ユーザーの問題解決を支援した際、ユーザーのワークフローは Beam on Apache Samza 上にあり、p collections や p transformations を使用してビジネスロジックを記述していました。
出力と入力のメソッドシグネチャが大きく異なるため、既存のビジネスロジックを書き換えた Apache Flink や Apache Spark のジョブで再利用できるよう、軽量なコンバーターを開発しました。


DAG の観点では、ケース A はオブジェクトをシンプルに直接変換する非常にシンプルなビジネスプロセスです。
この場合、Apache Flink と Apache Spark のパフォーマンスは非常に近いです。


通常、Apache Flink ダッシュボード UI を使用して例外やビジネスプロセスなどを確認しますが、これは Apache Spark に対する明らかな優位点です。
Apache Spark ではドライバーログから例外を検索する必要があり、より手間がかかります。
ただし、Apache Flink にはまだ改善が必要な領域がいくつかあります。


History Server - より豊富なメトリクスなどをサポートする

Apache Spark History Server UI が表示するメトリクスは比較的豊富で、ユーザーのパフォーマンス分析に大きく役立ちます。
Apache Flink がバッチ処理を行う場合も Apache Spark のユーザーと同程度のメトリクス情報を確認できるようにすれば、ユーザーの開発難易度を下げ、開発効率を向上できます。


より良いバッチ運用保守ツール

LinkedIn が 2、3 年前から取り組んでいることを共有します。
LinkedIn では毎日 20 万件のジョブがクラスター上で実行されており、バッチユーザーが自身のジョブを運用保守できるよう、より優れたツールが必要です。
Dr. Elephant と GridBench を提供し、ユーザーが自身のジョブをデバッグして運用できるように支援しています。


Dr. Elephant はオープンソースで、ユーザーがジョブをより良くデバッグし、問題を特定して改善提案を行うことができます。
また、テストクラスターから本番クラスターへ移行する前に、Dr. Elephant が生成したレポートの評価結果のスコアに基づいて本番環境への移行可否を判断します。


GridBench は主に CPU メソッドのホットスポット分析など、データの統計分析を行い、ユーザーのジョブの最適化と改善を支援します。
GridBench も将来的にオープンソース化する予定で、Apache Flink を含むさまざまなエンジンフレームワークをサポートし、Apache Flink のジョブを GridBench でより良く評価できるようになります。
GridBench Talk: Project Optimum: Spark Performance at LinkedIn Scale を参照してください。


ユーザーは GridBench や Dr. Elephant が生成するレポートだけでなく、コマンドラインからもアプリケーションの CPU 時間やリソース消費など、ジョブの最も基本的な情報を確認でき、異なる Apache Spark ジョブと Apache Flink ジョブの差異を比較分析することも可能です。


以上が Apache Flink のバッチ処理で改善が必要な 2 つの領域です。


3.2 ケース B - 複雑な ETL データフロー

2 つ目の方向性、バッチ処理 → ストリーム処理は、ストリームバッチ統合に分類されます。


ETL データフローのコアロジックは比較的複雑で、セッションウィンドウ集計ウィンドウを含み、1 時間ごとのユーザーページビューを計算し、異なるジョブに分割します。
中間のメタデータテーブルでページキーを共有し、最初のジョブは 00 分時点を処理し、2 番目のジョブは 01 分時点を処理してセッション化操作を行い、最終的に結果をオープンセッションとクローズセッションに分けて出力することで、各時間のデータを増分処理します。


このワークフローは元々 Apache Spark SQL によるオフライン増分処理で、純粋なオフライン増分処理でした。
ユーザーがジョブをオンラインに移行してリアルタイム処理を行いたい場合、Beam On Apache Samza などのリアルタイムワークフローを再構築する必要があります。
構築プロセス中、ユーザーと非常に密接に連絡を取り合い、ユーザーは多くの問題に直面しました。
開発ロジック全体の再利用、2 つのビジネスロジックが同じ結果を生成することの保証、データの最終保存先などです。
移行には長い時間がかかり、最終的な結果もあまり良いものではありませんでした。


さらに、ユーザーのジョブロジックは Apache Hive と Apache Spark の両方を使用して記述された多数の大規模で複雑な UDF を使用しており、この移行も非常に重い作業です。
ユーザーは Apache Spark SQL と Apache Spark DataFrame API に精通しています。


上図の黒い実線はリアルタイム処理プロセスを示し、灰色の矢印は主にバッチ処理プロセスを示しており、Lambda アーキテクチャに相当します。


ケース B のジョブには多数の結合とセッションウィンドウが含まれており、以前は Apache Spark SQL を使用してジョブを開発していました。
明らかに宣言型 API からアプローチする必要があります。
現在、3 つのソリューションを提供しています。


最初の選択肢は Apache Flink Table API/SQL で、ストリーム処理とバッチ処理の両方を実行できます。
同じ SQL で包括的な機能サポートを備え、ストリーム処理とバッチ処理の両方に最適化も施されています。
Alibaba Cloud Blog: What's All Involved with Blink Merging with Apache Flink?
および FLINK-11439 INSERT INTO flink_sql SELECT * FROM blink_sql の記事を参照してください。


2 番目の選択肢は Apache Spark DataFrame API/SQL で、同じインターフェイスでバッチ処理とストリーム処理を使用できますが、Apache Spark のストリーム処理サポートはまだ十分ではありません。


3 番目の選択肢は Beam Schema Aware API/SQL です。
Beam は物理 API の側面が強く、Schema Aware API/SQL は現在初期段階の作業が行われているため、現時点では検討対象外とします。
したがって、以降の分析結果と経験は主に Apache Flink Table API/SQL と Apache Spark DataFrame API/SQL の比較から得られています。
Beam Design Document - Schema-Aware PCollections および Beam User Guide - Beam SQL overview の記事を参照してください。


ユーザーの視点から見ると、Apache Flink Table API/SQL と Apache Spark DataFrame API/SQL は非常に近いです。
キーワード、ルール、結合の書き方など、比較的小さな差異がありますが、ユーザーに一定の混乱をもたらす可能性があります。
使い方を間違っているのではないかと気になるでしょう。


Apache Flink と Apache Spark の両方とも Apache Hive と良く統合されており、Apache Hive UDF の再利用などが可能です。
これにより、ケース B の UDF 移行における負荷が半減しました。


パイプラインモードでの Apache Flink のパフォーマンスは Apache Spark よりも明らかに優れています。
ディスクへの書き込みの有無がパフォーマンスに大きな影響を与えることは想像に難くありません。
大量のディスク書き込みが必要な場合、各ステージでデータをディスクに落としてから再度読み取る必要があるため、処理パフォーマンスはディスクに書き込まないパイプラインモードよりも劣ります。
パイプラインは短期間の処理に適しており、20 分から 40 分の処理では依然として比較的大きな優位性があります。
パイプラインがより長くなると、フォールトトレランスはバッチモードに及びません。
Apache Spark のバッチパフォーマンスは依然として Apache Flink より優れています。
この領域は社内の事例に基づいて評価する必要があります。


Apache Flink のウィンドウサポートは他のエンジンよりも明らかに豊富で、セッションウィンドウなどはユーザーにとって非常に便利です。
以前、ユーザーはセッションウィンドウを実現するために多くの UDF を記述していました。
増分処理や全セッションの構築、レコードの抽出処理などです。
現在、セッションウィンドウオペレーターを直接使用することで、多くの開発コストを削減できました。
同時に、グループ集計などのウィンドウ操作もストリームバッチでサポートされています。


UDF はエンジンフレームワーク間で移行する際の最大の障壁です。
UDF が Apache Hive で記述されていれば移行は容易です。
Apache Flink と Apache Spark の両方が Apache Hive UDF を非常によくサポートしているからです。
しかし、UDF が Apache Flink や Apache Spark で記述されている場合、どのエンジンフレームワークに移行しても非常に大きな問題に直面します。
たとえば、OLAP 準リアルタイムクエリのために Presto に移行する場合などです。


UDF の再利用を実現するため、LinkedIn では transport プロジェクトを内部で開発し、オープンソースとして GitHub に公開しています。
LinkedIn が公開したブログ Transport: Towards Logical Independence Using Translatable Portable UDFs を参照してください。


transport はすべてのエンジンフレームワーク向けの User API を提供し、共通の関数開発インターフェイスを備え、Presto、Apache Hive、Apache Spark、Apache Flink など異なるエンジンフレームワークに応じた UDF を自動生成します。


共通の UDF API を使用してすべてのエンジンフレームワークを接続することで、ユーザーは独自のビジネスロジックを再利用できます。
ユーザーは簡単に利用できます。


ユーザーの SQL 移行問題に対処するため、ユーザーは以前 Apache Spark SQL でジョブを開発していましたが、ストリームバッチ統合を使用して Apache Flink SQL に変更したいと考えていました。
現在も多くのエンジンフレームワークが存在します。
LinkedIn はオープンソースとして GitHub に公開されている coral ソリューションを開発し、Facebook でのトークも行い、transport UDF とともに分離レイヤーを提供することで、ユーザーがクロスエンジンの移行と独自のビジネスロジックの再利用をより良く実現できるようにしています。


Coral の実行プロセスを見てみましょう。
まず、馴染みのある ASCII SQL とテーブル属性がジョブスクリプトで定義され、Coral IR ツリー構造が生成され、最終的に各エンジンの物理プランに変換されます。



ケース B の分析では、ストリームとバッチが統合されています。
クラスターの業務量が特に多い場合、ユーザーはバッチ処理のパフォーマンス、安定性、成功率を非常に重視します。
その中で、Shuffle Service はバッチ処理のパフォーマンスに大きな影響を与えます。


4. Apache Spark と Apache Flink における Shuffle Service の比較

インメモリシャッフルは Apache Spark と Apache Flink の両方でサポートされており、高速ですが、スケーラビリティはサポートされていません。


ハッシュベースシャッフルは Apache Spark と Apache Flink の両方でサポートされています。
インメモリシャッフルと比較してフォールトトレランスのサポートは優れていますが、スケーラビリティはサポートされていません。


ソートベースシャッフルは大規模なシャッフルに対してスケーラビリティをサポートします。
ディスクから少しずつ読み込んでソートマッチを行い、その後読み戻します。
FLIP-148: Introduce Sort-Based Blocking Shuffle to Apache Flink でもサポートされています。


External Shuffle Service は、クラスターが非常に混雑している場合、たとえば動的リソーススケジューリング時に重要です。
シャッフルのパフォーマンスとリソース依存性の分離がより良く、分離後のリソーススケジューリングがより効率的になります。
FLINK-11805 A Common External Shuffle Service Framework は現在再オープン中です。


Disaggregate Shuffle とビッグデータ分野では Cloud Native が提唱されており、Shuffle Service の設計でもコンピューティングとストレージの分離を考慮する必要があります。
FLIP-10653 Introduce Pluggable Shuffle Service Architecture では、プラグイン可能な Shuffle Service アーキテクチャが導入されています。


Apache Spark は Shuffle Service に比較的大きな改善を加えました。この作業は LinkedIn が主導する magnet プロジェクトでもあり、Apache Spark 向けのスケーラブルで高性能なシャッフルアーキテクチャに関する論文 introducing-magnet (Magnet: A scalable and performant shuffle architecture for Apache Spark) が LinkedIn ブログ 2020 に収録されました。Magnet はディスク読み書きの効率を明らかに向上させています。比較的小さなランダム読み取りから比較的大きなシーケンシャル読み取りに変更し、ランダムなシャッフルデータの読み取りの代わりにマージ処理を行うことで、ランダム IO の問題を回避しています。

シャッフルの安定性とスケーラビリティの問題は Magent Shuffle Service により軽減されています。それ以前は、ジョブの失敗率が高いなど多くのシャッフルの問題が発見されていました。Apache Flink をバッチ処理に使用し、以前 Apache Spark でバッチ処理を行っていたユーザーを支援するには、シャッフルの部分により多くの労力を割く必要があります。

シャッフルの可用性に関しては、best-effort 方式でシャッフルブロックをプッシュし、大きなブロックをスキップして結果整合性と精度を保証します。
シャッフル一時データのコピーを作成し、精度を確保します。
プッシュプロセスが特に遅い場合、早期終了技術が適用されます。

Vanilla Shuffle と比較して、Magent シャッフルはシャッフルデータの読み取り待ち時間をほぼ 100%、タスク実行時間をほぼ 50%、エンドツーエンドのタスク持続時間をほぼ 30% 削減しています。

5. まとめ

LinkedIn は、Apache Flink がストリーム処理とバッチ処理の両方で明らかな優位性を持ち、より統一され継続的に最適化されていることを高く評価し、喜んでいます。

Apache Flink のバッチ処理機能は改善が必要です。History Server、メトリクス、デバッグなどです。ユーザーは開発時にユーザーコミュニティのソリューションを参照する必要があります。ユーザーが便利に使えるよう、エコシステム全体を構築する必要があります。

Apache Flink は Shuffle Service と大規模クラスターのオフラインワークフローにより多くのエネルギーを投入し、ワークフローの成功率を確保する必要があります。また、スケールが増大した場合のユーザーサポートの改善とクラスターの健全性監視も必要です。

フレームワークエンジンが増える中、ユーザーにより統一されたインターフェイスを提供することが理想的です。この領域の課題は比較的大きく、開発や運用保守も含まれます。LinkedIn の経験によると、依然として多くの問題があり、単一のソリューションですべてのユーザーのユースケースをカバーすることは不可能で、Coral や transport UDF のようなアプローチでも一部の機能や表現を完全にカバーすることは困難です。

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.