Build a streaming data lake with Flink Hudi

1. 背景

ニアリアルタイム

2016 年以降、Apache Hudi コミュニティは Hudi の UPSERT 機能を活用して、ニアリアルタイムシナリオのユースケースを模索してきました [1]。MR / Spark のバッチ処理モデルを通じて、ユーザーは時間レベルのデータを HDFS / OSS に投入できます。純粋なリアルタイムシナリオでは、ストリームコンピューティングエンジン Flink と KV / OLAP ストレージのアーキテクチャにより、エンドツーエンドの秒レベル(5 分間)リアルタイム分析を実現できます。しかし、秒レベル(5 分レベル)から時間レベルの間にはまだ大量のユースケースが存在し、これをニアリアルタイムと呼びます。

実践では、ニアリアルタイムに分類されるケースが多数あります:

分レベルのダッシュボード表示;
各種 BI 分析(OLAP);
機械学習向けの分レベルの特徴抽出。
増分コンピューティング
ニアリアルタイムへの対応は、現在比較的オープンです。

ストリーム処理はレイテンシが低いですが、SQL のパターンが比較的固定されており、クエリ側(インデックス、アドホック)の機能が不足しています;
バッチ処理のデータウェアハウスは機能が豊富ですが、データ遅延が大きいという課題があります。
そこで Hudi コミュニティは、ミニバッチベースの増分コンピューティングモデルを提案しました:

増分データセット => 増分計算結果を保存結果にマージ => 外部ストレージ

このモデルでは、レイクに保存されたスナップショットを通じて増分データセット(2 つのコミット間のデータセット)をプルし、Spark / Hive などのバッチ処理フレームワークで増分結果(シンプルな count など)を計算してから、保存結果にマージします。

主要な問題
増分モデルが解決すべき核心的な問題:

UPSERT 機能:Kudu や Hive ACID と同様に、Hudi も分レベルの更新機能を提供します;
増分消費:Hudi はレイクに保存された複数のスナップショットを通じて増分プルを提供します。
ミニバッチベースの増分コンピューティングモデルは、一部シナリオのレイテンシを改善し、コンピューティングコストを削減できますが、大きな制約があります。SQL パターンに要件があることです。バッチで計算が実行されるため、バッチ計算自体は状態を保持せず、計算される指標は容易にマージ可能である必要があります。シンプルな count や sum は実現できますが、avg や count distinct は全量データをプルして再計算する必要があります。

ストリームコンピューティングとリアルタイムデータウェアハウスの普及に伴い、Hudi コミュニティも積極的に変化を受け入れ、ストリームコンピューティングを通じて元のミニバッチベースの増分コンピューティングモデルを継続的に最適化および進化させています。バージョン 0.7 ではストリーミングデータレイク取り込みが導入され、バージョン 0.9 ではネイティブ CDC 形式がサポートされました。

2. 増分 ETL

DB データのレイクへの取り込み

CDC テクノロジーの成熟に伴い、Debezium などの CDC ツールがますます普及しており、Hudi コミュニティもストリーム書き込みとストリーム読み取りの機能を統合しました。ユーザーは Flink SQL を通じて CDC データを Hudi ストレージにリアルタイムで書き込めます:

Flink CDC コネクタ経由で DB データを直接 Hudi にインポートできます;
CDC データを先に Kafka にインポートしてから、Kafka コネクタ経由で Hudi にインポートすることもできます。
2 番目の方式の方がフォールトトレランスとスケーラビリティに優れています。

データレイク CDC
リリース予定のバージョン 0.9 では、Hudi が CDC 形式をネイティブでサポートし、1 つのレコードのすべての変更記録を保存できます。これにより、Hudi とストリームコンピューティングシステムの連携がより緊密になり、CDC データをストリームで読み取れます [2]:

ソース CDC ストリームのすべてのメッセージ変更は、レイクに取り込まれた後保存され、ストリーミング消費に使用されます。Flink のステートフル計算が計算結果(state)をリアルタイムで蓄積し、計算変更を Hudi へのストリーム書き込みを通じて Hudi レイクストレージに同期します。その後、Hudi に保存されたチェンジログを Flink がストリーム消費することで、次のレベルのステートフル計算を実現します。ニアリアルタイムのエンドツーエンド ETL パイプライン:


このアーキテクチャにより、エンドツーエンドの ETL レイテンシが分レベルに短縮され、各レイヤーのストレージ形式はコンパクションを通じてカラムナーストレージ(Parquet、ORC)に圧縮し、OLAP 分析機能を提供できます。データレイクのオープン性により、圧縮後の形式をさまざまなクエリエンジン(Flink、Spark、Presto、Hive など)に接続できます。

Hudi データレイクテーブルは 2 つの形式を持ちます:

テーブル形式:効率的なカラムナーストレージ形式を提供しながら、最新のスナップショット結果をクエリします
ストリーミング形式:変更をストリーム消費し、任意の時点からチェンジログを指定してストリーム読み取りできます

3. デモ

デモを通じて Hudi テーブルの 2 つの形式を実演します。

環境構築
Flink SQL Client
Hudi Master の hudi-flink-bundle jar をパッケージ化
Flink 1.13.1

事前に debezium-json 形式の CDC データを用意します

Flink SQL Client でテーブルを作成し、CDC データファイルを読み取ります

SELECT を実行して結果を観察すると、計 20 レコードあり、途中に UPDATE が含まれ、最後のメッセージは DELETE であることが確認できます

Hudi テーブルを作成します。ここではテーブル形式を MERGE_ON_READ に設定し、changelog.enabled 属性を有効にします

クエリ
INSERT 文でデータを Hudi にインポートし、ストリーム読み取りモードを有効にして、クエリを実行して結果を観察します

Hudi が各行の変更記録を保持していることが確認できます。変更ログの操作タイプも含まれます。ここでは TABLE HINTS 機能を有効にして、テーブルパラメータの動的設定を容易にします。

次にバッチ読み取りモードを使用し、クエリを実行して出力結果を観察すると、中間の変更がマージされていることが確認できます。

集計

Bounded Source 読み取りモードで count(*) を計算します
ストリーム読み取りモードで count(*) を計算します


バッチモードとストリームモードの計算結果が一致していることが確認できます。

現在のデータレイク CDC 形式はまだ急速な開発期間にあり、コミュニティも本番環境でのユースケースを積極的に推進しています。Hudi のシナリオやケースに興味のある方は、コードをスキャンしてグループに参加してください。

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.