Apache Flink is more than computing

2021 年の初め、InfoQ 編集部が企画した年間技術トレンド展望で、私たちはビッグデータ分野が「融合」(または「統合」)進化の新しい方向性への取り組みを加速させると述べました。その本質は、ビッグデータ分析の技術的な複雑さとコストを削減しながら、パフォーマンスと使いやすさに対するより高い要件を満たすことです。今日、私たちは人気のあるストリーム処理エンジン Apache Flink(以下、Flink と呼ぶ)がこのトレンドに沿って新たな一歩を踏み出しているのを目にしています。

1 月 8 日朝、Flink Forward Asia 2021 がオンラインカンファレンス形式で開幕しました。今年は Flink Forward Asia(以下、FFA と呼ぶ)が中国で開催されて 4 年目であり、Flink が Apache Software Foundation のトッププロジェクトとなって 7 年目でもあります。リアルタイム化の波の発展と深化に伴い、Flink は徐々にストリーム処理のリーダー的役割およびデファクトスタンダードへと進化してきました。その進化を振り返ると、一方では Flink はストリームコンピューティングのコア能力を継続的に最適化し、業界全体のストリームコンピューティング処理基準を絶えず向上させています。他方では、ストリーム・バッチ統合の考え方に沿って、アーキテクチャ変革とアプリケーションシナリオを徐々に推進しています。しかし、これらに加えて、Flink の長期的な発展には新たなブレークスルーが必要です。

Flink Forward Asia 2021 の基調講演で、Apache Flink 中国コミュニティの創設者であり Alibaba のオープンソースビッグデータプラットフォームの責任者である Wang Feng(Human Mowen)は、ストリーム・バッチ統合アーキテクチャの進化と実装における Flink の最新の進捗、および Flink の次の開発方向である Streaming Warehouse(以下、Streamhouse と呼ぶ)を提案しました。基調講演のタイトル「Flink Next, Beyond Stream Processing」が示すように、Flink は Stream Processing から Streaming Warehouse へと移行し、より大きなシナリオをカバーして、開発者がより多くの問題を解決できるよう支援します。ストリーミングデータウェアハウスの目標を達成するには、Flink コミュニティがストリーム・バッチ統合に適したデータストレージを拡張する必要があることを意味します。これは Flink の今年の技術革新であり、コミュニティ関連の作業は 10 月に開始されました。Flink コミュニティが来年進める重要な方向性の一つです。

では、ストリーミングデータウェアハウスをどのように理解すべきでしょうか?既存のデータアーキテクチャでどのような問題を解決しようとしているのでしょうか?なぜ Flink はこの方向を選択したのでしょうか?ストリーミングデータウェアハウスの実装パスは何でしょうか?これらの疑問を念頭に、InfoQ は Mowen 氏に独占インタビューを行い、ストリーミングデータウェアハウスの背後にある思考パスをさらに理解しました。

Flink はここ数年、ストリーム・バッチ統合を繰り返し強調してきました。つまり、同じ API セットと同じ開発パラダイムを使用してビッグデータのストリームコンピューティングとバッチコンピューティングを実現し、それによって処理プロセスと結果の一貫性を保証します。Mowen 氏によると、ストリーム・バッチ統合はむしろ技術的な概念と能力であり、それ自体はユーザーのいかなる問題も解決しません。実際のビジネスシナリオに実際に実装されて初めて、開発効率と運用効率の価値を発揮できます。ストリーミングデータウェアハウスは、ストリーム・バッチ統合の大きな方向性の下でのソリューション実装の考え方と理解できます。

ストリーム・バッチ統合の 2 つのアプリケーションシナリオ

昨年の FFA で、私たちは Tmall の独身の日における Flink のストリーム・バッチ統合の応用を見てきました。1 年が経過した今、Flink のストリーム・バッチ統合は技術アーキテクチャ進化と実装アプリケーションの両面で新たな進捗を遂げています。

技術進化のレベルでは、Flink のストリーム・バッチ統合 API とアーキテクチャ変革が完了しました。元のストリーム・バッチ統合 SQL に加えて、DataStream API と DataSet API の 2 つの API セットがさらに統合され、Java セマンティックレベルで完全なストリーム・バッチ統合 API が実現されました。1 つのコードセットでストリームストレージとバッチストレージを同時に実行できます。

今年 10 月にリリースされた Flink 1.14 バージョンでは、同じアプリケーション内有界ストリームと無制限ストリームの混在使用が既にサポートされています。Flink は現在、部分的に実行され部分的に終了しているアプリケーション(一部のオペレーターが有界入力データの処理を終了したストリーム)のチェックポイント実行をサポートしています。さらに、Flink は有界データストリームの終了処理時に最終チェックポイントをトリガーして、すべての計算結果が Sink に正常に送信されることを保証します。

バッチ実行モードでは、同じアプリケーション内での DataStream API と SQL/Table API の混在使用がサポートされるようになりました(以前は DataStream API または SQL/Table API の単独使用のみサポートされていました)。

さらに、Flink は統一された Source API と Sink API を更新し、統一 API を中心にコネクタエコシステムの統合を開始しました。新しいハイブリッドソースは複数のストレージシステム間を移行でき、Amazon S3 からの古いデータの読み取りから Apache Kafka へのシームレスな切り替えなどの操作を可能にします。

実装アプリケーションのレベルでは、2 つのより重要なアプリケーションシナリオがあります。

1 つ目は、Flink CDC に基づくフル増分統合データ統合です。

異なるデータソース間のデータ統合とデータ同期は多くのチームにとって必要ですが、従来のソリューションはしばしば複雑すぎて時間がかかります。従来のデータ統合ソリューションは通常、オフラインデータ統合とリアルタイムデータ統合に 2 つの技術スタックを使用し、多くのデータ同期ツール(Sqoop、DataX など)が関与します。これらのツールはフルまたは増分のいずれかのみを実行でき、開発者はフル増分切り替えを自分で制御する必要があるため、連携がより複雑になります。

Flink のストリーム・バッチ統合機能と Flink CDC を使用すると、単一の SQL を記述するだけで履歴データのフル同期を行い、その後自動的に増分データ転送をブレークポイントから再開して、ワンストップデータ統合を実現できます。プロセス全体でユーザーの判断や介入は必要なく、Flink がバッチ間の切り替えを自動的に完了し、データ整合性を保証します。

独立したオープンソースプロジェクトとして、Flink CDC Connectors は昨年の 7 月のオープンソース化以来、かなり高速な開発を維持しており、平均して 2 か月ごとに 1 つのバージョンがリリースされています。ますます多くの企業が Flink CDC を自社のビジネスシナリオで使用していることがわかります。少し前に InfoQ がインタビューした XTransfer もその一つです。

2 つ目のアプリケーションシナリオは、ビッグデータ分野のコアデータウェアハウスシナリオです。

ほとんどのシナリオでは Flink+Kafka を使用してリアルタイムデータストリームを処理します。つまり、リアルタイムデータウェアハウスであり、最終的な分析結果をオンラインサービスレイヤーに書き込んで表示またはさらに分析します。同時に、リアルタイムデータを補完するために、非同期のオフラインデータウェアハウス構造がバックグラウンドに存在し、毎日大規模なバッチまたはフルスケール分析を定期的に実行したり、履歴データの定期的な修正を行ったりします。

しかし、この古典的なアーキテクチャにはいくつかの明らかな問題があります。まず、リアルタイムリンクとオフラインリンクで使用される技術スタックが異なり、2 つの API セットが必要であるため、2 つの開発プロセスが必要で、開発コストが増加します。次に、リアルタイムとオフラインの技術スタックが異なるため、データ口径の一貫性を保証できません。さらに、リアルタイムリンクの中間キューデータは分析に不利です。ユーザーがリアルタイムリンクの詳細レイヤーのデータを分析したい場合、実際には非常に不便です。多くのユーザーは、この詳細レイヤーのデータを最初にエクスポートする方法を使用する可能性があります。例えば、Hive にインポートしてオフライン分析を行う方法です。しかし、リアルタイムパフォーマンスは大幅に低下します。または、クエリを高速化するためにデータを他の OLAP エンジンにインポートしますが、これによりシステムの複雑さが増し、データ整合性の保証も困難です。

Flink のストリーム・バッチ統合の概念は、前述のシナリオで完全に適用できます。Mowen 氏の視点では、Flink は業界の現在の主流データウェアハウスアーキテクチャをもう一段階進化させ、真のエンドツーエンドフルリンクのリアルタイム分析機能を実現できます。つまり、ソースでデータが変化したときにその変化をキャプチャし、レイヤーごとの分析をサポートして、すべてのデータをリアルタイムで流れさせ、流れているすべてのデータのリアルタイムクエリを可能にします。Flink の完全なストリーム・バッチ統合機能を活用することで、同じ API セットで柔軟なオフライン分析も同時にサポートできます。これにより、リアルタイム、オフライン、インタラクティブクエリ分析、ショートクエリ分析などを完全なソリューションセットに統合し、理想の「Streaming Warehouse」になることができます。

ストリーミングデータウェアハウスの理解

ストリーミングデータウェアハウスをより正確に言えば、実際には「データウェアハウスをストリーミング化する」ことです。つまり、データウェアハウス全体のデータをリアルタイムで流れるようにし、かつマイクロバッチ(ミニバッチ)ではなく純粋なストリーム方法で流れるようにします。目標は、エンドツーエンドのリアルタイム純ストリームサービス(Streaming Service)を実装し、1 セットの API を使用して流れているすべてのデータを分析することです。ソースデータが変化したとき、例えばオンラインサービスの Log やデータベースの Binlog をキャプチャすると、事前に定義されたクエリロジックまたはデータ処理ロジックに従ってデータが分析され、分析されたデータはデータウェアハウスの特定のレイヤーに落ちます。そして、最初のレイヤーから次のレイヤーへと流れ、すべてのデータウェアハウスレイヤーがすべて流れ、最終的にオンラインシステムに流れ込み、ユーザーはデータウェアハウス全体のフルリアルタイムフロー効果を見ることができます。このプロセスで、データはアクティブであり、クエリはパッシブであり、分析はデータの変化によって駆動されます。同時に、垂直方向では、各データ詳細レイヤーに対して、ユーザーはクエリを実行してアクティブにクエリでき、リアルタイムでクエリ結果を取得できます。さらに、オフライン分析シナリオとも互換性があり、API は同じで真の統合を実現します。

現在、業界にはエンドツーエンドフルストリーミングリンクの成熟したソリューションはありません。純粋なストリーミングソリューションと純粋なインタラクティブクエリソリューションは存在しますが、ユーザーは 2 つのソリューションを自分で追加する必要があり、必然的にシステムの複雑さが増します。オフラインデータウェアハウスソリューションを追加する場合は、システムの複雑さの問題はさらに大きくなります。ストリーミングデータウェアハウスがすべきことは、システム複雑さをさらに増加させることなく、高い適時性を実現し、開発者と運用担当者にとってアーキテクチャを非常にシンプルにすることです。

もちろん、ストリーミングデータウェアハウスは最終状態です。この目標を達成するには、Flink にはストリーム・バッチ統合ストレージサポートが必要です。実際、Flink 自体には分散 RocksDB がステートストアとして組み込まれていますが、このストレージはタスク内のストリームデータステートを保存する問題のみを解決できます。ストリーミングデータウェアハウスには、コンピューティングタスク間のテーブルストレージサービスが必要です。最初のタスクがそこにデータを書き込み、2 番目のタスクがそこからリアルタイムで読み取り、3 番目のタスクがユーザークエリを実行して分析できます。したがって、Flink は自身の概念に合致するストレージを拡張し、State ストレージから踏み出し、さらに外側へ進む必要があります。このために、Flink コミュニティは新しい Dynamic Table Storage を提案しました。これはストリームテーブル双対性を持つストレージソリューションです。

Flink Dynamic Table(コミュニティディスカッションについては FLIP-188 を参照)は、Flink SQL にシームレスに接続されるストリーム・バッチストレージのセットと理解できます。元々、Flink は Kafka や HBase などの外部テーブルの読み書きのみ可能でしたが、現在、同じ Flink SQL 構文セットで、元のソーステーブルやターゲットテーブルと同様に Dynamic Table を作成できます。ストリーミングデータウェアハウスのすべてのレイヤードデータを Flink Dynamic Table に配置し、Flink SQL を通じてデータウェアハウス全体のレイヤーをリアルタイムで接続し、Dynamic Table 内の異なる詳細レイヤーのデータに対してリアルタイムクエリと分析を実行できます。異なるレイヤーでバッチ ETL 処理を行うことも可能です。

データ構造の観点では、Dynamic Table 内部には 2 つのコアストレージコンポーネントがあります。File Store と Log Store です。名前が示すように、File Store は古典的な LSM アーキテクチャを使用してテーブルファイルを保存し、ストリーミング更新、削除、追加などをサポートします。同時に、オープンなカラムストレージ構造を採用し、圧縮などの最適化をサポートします。Flink SQL のバッチモードに対応し、フルバッチ読み取りをサポートします。Log Store はテーブルの操作レコードを保存し、これは不変のシーケンスです。Flink SQL のストリームモードに対応し、Flink SQL を通じて Dynamic Table の増分変化をサブスクライブしてリアルタイム分析できます。現在、プラグイン実装をサポートしています。

File Store への書き込みは組み込み Sink にカプセル化され、書き込みの複雑さを隠蔽します。同時に、Flink のチェックポイントメカニズムと Exactly Once メカニズムはデータ整合性を保証できます。

現在、Dynamic Table の第 1 段階の実装計画が完了し、コミュニティもこの方向でさらに議論しています。コミュニティの計画によると、将来の最終状態では Dynamic Table のサービスが実現し、真に Dynamic Table Service のセットを形成し、リアルタイムでストリーム・バッチ統合ストレージを実現します。同時に、Flink コミュニティは Dynamic Table の操作とリリースを Flink の独立したサブプロジェクトとして議論しています。将来、ストリーム・バッチ統合の汎用ストレージプロジェクトとして完全に独立する可能性も排除されません。最後に、Flink CDC、Flink SQL、Flink Dynamic Table を使用することで、完全なストリーミングデータウェアハウスを構築し、リアルタイムオフライン統合体験を実現できます。全体のプロセスと効果については、以下のデモビデオを参照してください。

全体のプロセスは当初順調ですが、真にフルリアルタイムリンクを実現し十分に安定させるには、コミュニティは実装ソリューションの品質を徐々に向上させる必要があります。これには、OLAP インタラクティブシナリオでの Flink SQL の最適化、Dynamic Table ストレージパフォーマンスと整合性の最適化、Dynamic Table サービス機能の構築など多くのタスクが含まれます。ストリーミングデータウェアハウスの方向性は始まったばかりで、予備的な試みが行われています。Mowen 氏の視点では、設計に問題はありませんが、一連のエンジニアリング問題を将来解決する必要があります。これは、先進的なプロセスチップや ARM アーキテクチャを設計するようなものです。多くの人が設計できますが、歩留まりを保証する前提でチップを実際に製造することは非常に困難です。ストリーミングデータウェアハウスは、ビッグデータ分析シナリオにおける Flink の最も重要な方向となり、コミュニティもこの方向に多大な投資を行います。

Flink はコンピューティングを超える

ビッグデータのリアルタイム変革の大きなトレンドの下で、Flink は 1 つのことだけでなく、より多くのことができます。

業界の元々の Flink の位置づけは、ストリームプロセッサまたはストリームコンピューティングエンジンとしてのものが多かったですが、そうではありません。Mowen 氏によると、Flink はネイティブコンピューティングだけではありません。狭義のコンピューティングと考えるかもしれませんが、広義では、Flink には既にストレージがあります。「Flink がストリームコンピューティングで包囲網を突破できるのは、ステートフルストレージに依存しており、これは Storm よりも大きな優位性です。」

現在、Flink はさらに進んで、より広範なリアルタイム問題をカバーするソリューションを実現することを望んでいますが、元のストレージでは不十分です。しかし、外部ストレージシステムや他のエンジンシステムは Flink の目標と特徴と完全に一致せず、Flink とうまく統合できません。例えば、Flink は Hudi や Iceberg を含むデータレイクと既に統合されており、リアルタイムレイク取り込みとレイク取り込みのリアルタイム増分分析をサポートしています。しかし、これらのシナリオでも Flink のフルリアルタイム優位性を十分に活用することはできません。データレイクストレージフォーマットの本質は依然としてミニバッチであり、Flink もミニバッチモードに低下するからです。これは Flink が最も見たいアーキテクチャでも、Flink に最も適したアーキテクチャでもありません。したがって、当然、Flink のストリーム・バッチ統合概念に合致するストレージシステムを開発する必要があります。

Mowen 氏の視点では、ビッグデータコンピューティング分析エンジンセットには、自身の概念をサポートするストレージ技術システムのサポートがなければ、究極の体験を持つデータ分析ソリューションセットを提供することは不可能です。これは、優れたアルゴリズムもそれに対応するデータ構造が伴って初めて、最適な効率で問題を解決できることと似ています。

なぜ Flink をストリーミングデータウェアハウスと言う方が適切なのでしょうか?これは Flink の概念によって決定されます。Flink のコア概念は、ストリーミングを優先してデータ処理問題を解決することです。データウェアハウス全体のデータをリアルタイムで流れるようにするには、ストリーミングが不可欠です。データが流れた後、集約データのフローテーブル双対性と Flink のストリーム・バッチ統合分析機能により、フロー内の任意のリンクのデータを分析できます。ショートクエリの秒レベル分析でも、オフライン ETL 分析でも、Flink は対応する機能を持っています。Mowen 氏によると、Flink のストリーム・バッチ統合の最大の制限は、中間にサポートするストレージデータ構造がないことで、これによりシナリオの実装が困難になります。ストレージとデータ構造を補完する限り、ストリーム・バッチ統合の多くの化学反応が自然に現れます。

Flink の独自構築データストレージシステムは、ビッグデータエコシステムの既存データストレージプロジェクトに一定の影響を与えるのでしょうか?これに関して、Mowen 氏は、Flink コミュニティが新しいストリーム・バッチ統合ストレージ技術をリリースしたのは、自身のストリーム・バッチ統合コンピューティングニーズをよりよく満たすためだと説明しました。ストレージとデータに対してオープンプロトコル、オープン API と SDK を維持し、将来このプロジェクトを独立して開発する計画もあります。さらに、Flink は業界の主流ストレージプロジェクトとの接続を積極的に維持し、外部エコシステムの互換性と開放性を維持します。

ビッグデータエコシステムの異なるコンポーネント間の境界がますます曖昧になっています。Mowen 氏は、現在のトレンドは単一コンポーネント能力から統合ソリューションへ移行していると考えています。「誰もが実際にはこのトレンドに従っています。例えば、多くのデータベースプロジェクトを見ることができます。それらは元々 OLTP でしたが、後に OLAP を追加し、最後に HTAP と呼びました。実際、これらはロウストレージとカラムストレージの組み合わせで、分析をサポートしてユーザーに完全なデータ分析体験を提供します。」Mowen 氏はさらに付け加えました。「現在、多くのシステムが境界を拡張し始めており、リアルタイムからオフラインへ、またはオフラインからリアルタイムへ、相互に浸透しています。そうでなければ、ユーザーはさまざまな技術コンポーネントを自分で組み合わせ、さまざまな複雑さに対処する必要があり、閾値がますます高くなっています。したがって、統合トレンドは非常に明白です。誰が誰を統合するかに正誤はなく、重要なことは良い統合方法でユーザーに最高の体験を提供できるかどうかです。それを実行した者が最終的なユーザーを獲得します。コミュニティには活力と継続的な発展が必要であり、自身が最も得意とする分野を最大化するだけでは不十分で、ユーザーニーズとシナリオに基づいて境界を継続的に革新し突破する必要があります。ほとんどのユーザーのニーズは、単一能力の 95 点から 100 点までのギャップにあるわけではありません。」

Mowen 氏の推定によると、比較的成熟したストリーミングデータウェアハウスソリューションを形成するには約 1 年かかります。Flink をリアルタイムコンピューティングエンジンとして採用しているユーザーにとっては、新しいストリーミングデータウェアハウスソリューションを試すことが自然に適しており、ユーザーインターフェースは Flink SQL と完全に互換性があります。報道によると、最新の Flink 1.15 バージョンでは最初のプレビューバージョンがリリースされる予定で、Flink を使用しているユーザーは最初に試すことができます。Mowen 氏によると、Flink ベースのストリーミングデータウェアハウスはまだ始まったばかりで、技術ソリューションにはさらなるイテレーションが必要であり、成熟するまでに磨き上げる時間が必要です。より多くの企業や開発者が自身のニーズで構築に参加することを望んでいます。それがオープンソースコミュニティの価値です。

エピローグ

数多くのオープンソースエコロジカルコンポーネントの問題とビッグデータの高度なアーキテクチャ複雑さは長年批判されてきました。現在、業界はある程度コンセンサスに達しているようで、統合と融合を通じてデータアーキテクチャの進化を簡素化の方向に推進することです。企業によって異なる言い方と実装パスがありますが。

Mowen 氏の視点では、オープンソースエコロジーが繁栄することは正常です。各技術コミュニティには独自の専門分野がありますが、ビジネスシナリオの問題を本当に解決するには、ワンストップソリューションが必要で、ユーザーに使いやすい体験を提供します。したがって、彼は全体的なトレンドが統合と融合の方向に進むことに同意していますが、可能性は一意ではありません。将来、すべてのコンポーネントを統合する責任を負う専用システムがある可能性も、各システムが徐々に統合へと進化していく可能性もあります。どちらの可能性が最終的な結果になるのか、おそらく時間の答えを待つしかないでしょう。

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.