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 が直接クエリできない問題の解決推進。
現状、オフラインコンピューティングタスクの大部分はサービス 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
-
A detailed explanation of Hadoop core architecture HDFS
Knowledge Base Team
-
What Does IOT Mean
Knowledge Base Team
-
6 Optional Technologies for Data Storage
Knowledge Base Team
-
What Is Blockchain Technology
Knowledge Base Team
Explore More Special Offers
-
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
