Architecture and practice of EMR Delta Lake in Liulishuo data access

背景

現状、オフラインコンピューティングタスクの大部分はサービス DB から取得したデータに基づいています。サービス DB へのデータアクセスの精度、安定性、適時性が、ダウンストリーム全体のオフラインコンピューティングパイプラインの精度と適時性を決定づけています。同時に、DB 内のデータと Hive 内のデータを準リアルタイムで結合クエリする業務要件もあります。

Alibaba Cloud E-MapReduce (EMR) Delta Lake の導入前は、DataX をカプセル化してビジネス DB データへのアクセスを完結させていました。Master・Slave アーキテクチャを採用し、Master は毎日実行する DataX タスクのメタデータ情報を管理します。Worker ノードは、ステータスが init の DataX タスクを継続的なプリエンプションで取得して実行し、その日の全 DataX タスクが完了するまで処理を続けます。

アーキテクチャ図はおおよそ以下の通りです:

Worker の処理プロセスは以下の通りです:

準リアルタイム要件については、スレーブデータベースを用意し、Presto コネクタを設定してスレーブデータベースに接続することで、ビジネス DB のデータと Hive のデータの準リアルタイム結合クエリを実現しています。

このアーキテクチャの利点はシンプルで実装が容易なことです。しかし、データ量の増加に伴い、以下のような課題が徐々に表面化しました:

パフォーマンスボトルネック:ビジネスの成長に伴い、SELECT を介したデータアクセスのパフォーマンスが次第に悪化します。DB のパフォーマンスボトルネックに制約され、Worker ノードの追加では解消できません。

大規模テーブルではデータベースからフルプルするしかなく、データアクセスコストがますます増大します。

サービスが準リアルタイムクエリ要件を満たせず、準リアルタイムクエリはデータベース経由のみに限定され、アクセスコストがさらに増大します。

これらの課題を解決するため、CDC リアルタイム取り込みソリューションに注目しました。

技術方式の選定

現在、業界では CDC リアルタイム取り込みのソリューションとして主に以下が存在します:CDC+Merge、CDC+Hudi、CDC+Delta Lake、CDC+Iceberg です。その中で CDC+Merge 方式はデータレイク方式が登場する以前の実践です。この方式は DB スレーブデータベースのコストを節約できますが、ビジネスの準リアルタイムクエリ要件などを満たせないため、初期段階で除外されました。Iceberg も選定当時まだ十分に成熟しておらず、業界に参考事例がなかったため、これも除外されました。最終的に CDC+Hudi と CDC+Delta Lake のいずれかを選択することになりました。

選定時、Hudi と Delta Lake は機能的に同等であったため、主に安定性、小ファイル統合、SQL サポート、クラウドベンダーのサポート、言語サポートなどの観点から検討しました。

上記の指標に加え、データプラットフォーム全体が Alibaba Cloud EMR 上に構築されているため、Delta Lake を選択すれば多くの適合作業を省略できることから、最終的に CDC+Delta Lake を選択しました。

全体アーキテクチャ

全体構成図

全体アーキテクチャは上図の通りです。取り込むデータは既存の履歴データと新規データの 2 つの部分に分けられます。既存の履歴データは DataX を使用して MySQL からエクスポートし、OSS に保存します。新規データは Binlog を使用して収集し、Delta Lake テーブルに保存します。毎朝 ETL タスクを実行する前に、まず既存データと新規データをマージし、ETL タスクはマージ後のデータを使用します。

Delta Lake データ取り込み

Binlog のリアルタイム収集には、オープンソースの Debezium を使用します。Debezium は MySQL から Binlog をリアルタイムでプルし、適切な解析を行います。各テーブルは 1 つの Topic に対応し、データベースシャーディングおよびテーブルシャーディングのデータは 1 つの Topic に統合され、Kafka に配信されてアップストリームとダウンストリームで消費されます。Binlog データが Kafka に接続された後、対応する Kafka Topic を指す Kafka Source テーブルを作成する必要があります。テーブルのフォーマットは以下の通りです:

主に使用するフィールドは value と offset です。value のフォーマットは以下の通りです:

StreamingSQL で Kafka 内のデータを処理します。主に Kafka Source テーブルから offset、value フィールド、および value フィールド内の CDC 情報(op、ts_ms の after/before フィールド、payload など)を抽出します。StreamingSQL では 5 分間のミニバッチを使用します。主な理由は、ミニバッチが小さすぎると小さなファイルが多数生成され、処理速度が徐々に低下し、読み取りパフォーマンスにも影響するためです。逆に大きすぎると準リアルタイムクエリの要件を満たせません。Delta Lake テーブルについては、after フィールドや before フィールドを解析しません。主な理由は、ビジネステーブルのスキーマが頻繁に変更されるためです。スキーマが変更されるとデータを修復する必要があり、コストがかかります。StreamingSQL 処理中、op='c' のデータは直接挿入し、json_record は after フィールドを取得します。op='u' または op='d' のデータは、Delta Lake テーブルに存在しない場合は挿入操作を実行し、存在する場合は更新操作を実行します。json_record の代入は、op='d' の場合は before フィールド、op='u' の場合は after フィールドを取得します。op='d' のデータを保持する理由は、削除されたデータが既存の履歴テーブルに存在する可能性があり、直接削除すると早朝マージ時に既存の履歴テーブルのデータが削除されないためです。

Delta Lake はタイムトラベルをサポートしていますが、CDC データを取り込む場合、データロールバック戦略を利用できません。複数バージョンのデータを保持するとストレージに影響が出るため、期限切れバージョンのデータを定期的に削除する必要があります。現在、2 時間以内のバージョンデータのみを保持しています。同時に、Delta Lake は小ファイルの自動統合機能をサポートしていないため、定期的に小ファイル統合を行う必要もあります。現在、OPTIMIZE と VACUUM を使用して 1 時間ごとに小ファイル統合と期限切れデータファイルのクリーンアップを実行しています:

現在、Hive と Presto は Spark SQL で作成された Delta Lake テーブルを直接読み取れません。しかし、モニタリングと準リアルタイムクエリで Delta Lake テーブルをクエリする必要があるため、Hive と Presto 用のクエリテーブルも作成しました。

Delta Lake データと既存データの Merge

Delta Lake のデータは新規データのみを取り込んでいるため、既存の履歴データは DataX を使用して一括でインポートします。Delta Lake テーブルは Hive で直接クエリできないため、毎朝これら 2 つのデータ部分に対してマージ操作を実行し、新しいテーブルに書き込んで Spark SQL と Hive の両方で一元的に利用できるようにしています。このモジュールのアーキテクチャはおおよそ以下の通りです:

毎日 0 時前に DeltaService API を呼び出し、Delta Lake タスクの設定に基づいて、マージタスクのタスク情報、Spark SQL スクリプト、対応する Airflow DAG ファイルを自動生成します。

マージタスクのタスク情報には主に以下の情報が含まれます:

Merge スクリプトの自動生成では、主に Delta Lake タスクの設定から MySQL テーブルのスキーマ情報を取得し、既存の Hive テーブルを削除し、スキーマ情報に基づいて Hive 外部テーブルを再作成します。次に、Delta Lake テーブルの json_record フィールドと既存データテーブルから対応するフィールド値を取得し、UNION ALL 操作を実行します。欠損値は MySQL のデフォルト値となります。UNION の後、row_key でグループ化し、ts_ms でソートして最初の 1 件を取得すると同時に operation_type='d' のデータを除外します。全体の流れは以下の通りです:

0 時以降、Airflow は Airflow DAG ファイルに基づいてマージ Spark SQL スクリプトを自動スケジュールして実行します。スクリプトが正常に実行されると、マージタスクのステータスが成功に更新されます。Airflow の ETL DAG は、マージタスクのステータスに基づいて下流の ETL タスクを自動的にスケジュールします。

Delta Lake データモニタリング

Delta Lake データのモニタリングは主に 2 つの目的があります:データの遅延監視とデータの損失監視です。主に MySQL と Delta Lake テーブル間、および Kafka Topic と CDC でアクセスされる Delta Lake テーブル間を対象とします。

Kafka Topic と CDC でアクセスされる Delta Lake テーブル間の遅延監視:15 分ごとに Kafka Topic から各パーティションの最大オフセットに対応する MySQL の row_key フィールドの内容を取得し、監視対象の MySQL テーブル delta_kafka_monitor_info に格納します。次に、delta_kafka_monitor_info から前サイクルの row_key フィールドの内容を取得し、Delta Lake テーブルで照会します。見つからない場合はデータ遅延またはデータ損失が発生しているため、アラームを発します。

MySQL と Delta Lake 間の監視:2 つの方式があります。1 つはプローブスキームで、15 分ごとに MySQL から最大 ID を取得します。データベースとテーブルについては 1 つのテーブルのみを監視し、delta_mysql_monitor_info に格納します。次に、delta_mysql_monitor_info から前サイクルの最大 ID を取得し、Delta Lake テーブルで照会します。照会できない場合はデータ遅延またはデータ損失が発生していることを示し、アラームを発します。もう 1 つは直接 count(id) を使用する方法で、単一データベース・単一テーブルとデータベースシャーディング・テーブルシャーディングの場合に分かれます。メタデータは MySQL テーブル id_based_mysql_delta_monitor_info に格納され、主に min_id、max_id、mysql_count の 3 つのフィールドを含みます。単一データベース・単一テーブルの場合、5 分ごとに Delta Lake テーブルから min_id と max_id の間のカウント値を取得し、mysql_count と比較します。mysql_count より少ない場合はデータ損失または遅延があることを示し、アラームを発します。次に MySQL から max(id) と max(id) の間のカウント値を取得し、id_based_mysql_delta_monitor_info テーブルを更新します。データベースシャーディング・テーブルシャーディングの場合、シャーディングルールに基づいて各テーブルに対応する id_based_mysql_delta_monitor_info 情報を生成し、30 分ごとに監視します。ルールは単一データベース・単一テーブルの場合と同じです。

課題

ビジネステーブルのスキーマは頻繁に変更されます。Delta Lake テーブルが CDC のフィールド情報を直接解析する場合、データが見つからない場合にタイムリーに修復できないと、後期のデータ修復コストが大きくなります。現在はフィールドを解析せず、早朝マージ時まで解析を延期しています。

データ量の増加に伴い、StreamingSQL タスクのパフォーマンスが次第に悪化しています。現在は StreamingSQL 処理遅延への対応として、大量の遅延アラームが発生した後、Delta Lake の既存データを昨日のマージ後のデータで置き換え、Delta Lake テーブルを削除し、チェックポイントデータを削除し、KafkaSource テーブルデータを最初から再消費します。Delta Lake テーブルのデータ量を削減することで、StreamingSQL の負荷を軽減します。

Hive と Presto は Spark SQL で作成された Delta Lake テーブルを直接クエリできません。現在は Hive と Presto 向けに、Hive と Presto でクエリ可能な外部テーブルを作成していますが、これらのテーブルは Spark SQL ではクエリできません。そのため、上位の ETL アプリケーションはコードを変更せずに Hive、Spark SQL、Presto エンジン間を自由に切り替えることができません。

メリット

DB スレーブデータベースのコストを節約しました。CDC+Delta Lake の採用により、コストを約 80% 削減しました。

早朝の DB データアクセスの時間コストが大幅に削減され、特別な要件がないすべての DB データアクセスを 1 時間以内に完了できるようになりました。

今後の計画

Delta Lake テーブルのデータ量増加に伴い、StreamingSQL タスクのパフォーマンスが次第に悪化する課題への対応。

Spark SQL で作成された Delta Lake テーブルを Hive と Presto が直接クエリできない問題の解決推進。

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.