How Delta Lake helps cloud users solve the problem of real-time data storage
1. CDC の概要
CDC は Change Data Capture の略称であり、変更データキャプチャを意味します。たとえば、初期段階ではツールを使用してビジネスデータをデータウェアハウスやデータレイクにインポートしていました。その後、データをインポートする際に、データの動的な変化を反映し、増分インポートを実行し、これらの変更データをできるだけ早くキャプチャして、より迅速に後続の分析をタイムリーに実行できるようにすることが求められるようになりました。CDC 技術は、このような変化するデータをキャプチャするのに役立ちます。
ビッグデータシナリオでは、一般的に使用されるツールは Sqoop です。これはバッチモードのツールであり、ビジネスライブラリからデータウェアハウスにデータをインポートするために使用できます。インポート前に、ビジネスライブラリのデータ内で時間的な変化を反映できるフィールドを選択し、タイムスタンプに基づいて変更されたデータをデータウェアハウスにインポートする必要がある点に注意が必要です。これが Sqoop の制限事項です。さらに、このツールには以下の欠点があります。
ソースライブラリに負荷をかける
遅延が大きい(呼び出し頻度に依存)
削除イベントを処理できず、ソースデータベースで削除されたデータをデータウェアハウスで同期的に削除できない
スキーマの変更に追随できない。ソースライブラリのスキーマが変更されると、データウェアハウスのテーブルモデルの再構築とインポートが必要になる
Sqoop を使用する以外に、binlog を使用してデータ同期を行う方法もあります。ソースライブラリは、挿入、更新、削除などの操作を実行する際に binlog を生成します。binlog を Kafka に入力し、Kafka から binlog を読み取って、1 つずつ解析した後に対応する操作を実行するだけです。ただし、この方法では、上流の頻繁な更新/削除操作に対応するために、ダウンストリームが比較的頻繁な更新/削除操作をサポートする必要があります。ここでは、KUDU または HBASE をターゲットストレージとして選択できます。ただし、KUDU と HBASE はデータウェアハウスではないため、フルデータの保存はできません。そのため、それらのデータは定期的に Hive にインポートする必要があります。下図に示す通りです。この方法には、複数コンポーネントの運用・メンテナンスの負荷が高い、マージロジックが複雑などの欠点がある点に注意が必要です。
2. Spark Streaming SQL & Delta に基づく CDC ソリューション
(1) Spark Streaming SQL
Spark Streaming SQL は、Alibaba Computing Platform Division の EMR チームが Spark Streaming を基に開発した SQL サポートであり、コミュニティ版には含まれていません。Spark Streaming SQL はこの CDC ソリューションに必須ではありませんが、特に SQL の使用に慣れているユーザーにとってより使いやすいものとなっています。そのため、EMR チームが Spark Streaming SQL のサポートを開発しました。下図に示すように、EMR の Spark Streaming SQL は、DDL、DML、SELECT など、多くの側面で SQL 構文のサポートを実装しています。以下にいくつかを紹介します。
(1) CREATE SCAN & CREATE STREAM
以下に示す例では、Kafka のテーブルから一部のデータを選択することが目標です。設計目標は、バッチとストリームの両方に対応できる限りサポートすることです。通常の SQL では、select は実際には読み取り操作を生成しますが、ここではバッチとストリーミングを区別するために、明示的な create scan が必要です。データソースからは、バッチ読み取りかストリーミング読み取りかを区別できないためです。バッチの場合は USING batch を使用し、ストリーミングの場合は USING stream を使用します。
バッチスキャンの場合、スキャンを作成した後、scan から直接 select でき、scan をテーブルとして扱うことができます。ただし、ストリーミングの場合、このスキャンを読み取るには多くのパラメータを設計する必要があります。ジョブを開始する必要があるためです。そのため、下図に示す create stream 構文があります。これは本質的に select 構文のカプセル化です。
(2) MERGE INTO
もう 1 つのコア構文は MERGE INTO であり、Delta Lake の CDC ソリューションで非常に重要な役割を果たします。MERGE INTO の構文はより複雑で、下図に示す通りです。MERGE INTO の mergeCondition は、ソーステーブルとターゲットテーブルの 1 対 1 対応である必要がある点に注意が必要です。そうでない場合、1 つのソースレコードが複数のターゲットレコードに対応すると、システムはどれを操作すればよいか判断できません。そのため、ここでは mergeCondition がプライマリキー結合、またはプライマリキー結合と同等の効果であることが実際には要求されます。
上で紹介した構文以外にも、DELAY、TUMBLING などの他の UDF も実装しており、Spark Streaming SQL をより簡単に使用できるようにしています。
(2) Delta Lake
データレイクは近年の注目技術です。初期段階では、比較的成熟したデータウェアハウスシステムが使用され、データは ETL を通じてデータウェアハウスにインポートされていました。データウェアハウスの典型的な用途は BI レポートなどの分析シナリオであり、シナリオは比較的限られていました。モバイルインターネット時代では、データソースがより豊富で多様化し、データ構造は構造化データに限定されず、データの用途も分析に限定されなくなったため、データレイクが登場しました。データを事前に処理しない、または簡単な処理のみ行ってデータレイクにインポートし、その後、スクリーニング、フィルタリング、変換などの transformation 操作を実行するため、データウェアハウス時代の ETL がデータレイク時代の ELT へと変化しました。
データレイクの典型的なアーキテクチャは、上位層に 1 つまたは複数の分析エンジンまたは他のコンピューティングフレームワーク、下位層に分散ストレージシステムを持つ構成です。下図の左側に示す通りです。ただし、この元のデータレイクの使用方法には、トランザクションサポートの欠如、データ品質検証の欠如など、管理が不足しており、すべてのデータ管理は完全に手動で保証されていました。
Delta Lake は、統一ストレージレイヤーの上に管理レイヤーを配置し、データレイクのデータを手動で管理する際の課題を解決します。管理レイヤーの追加により、まずメタデータ管理を導入できます。メタデータ管理があれば、データにスキーマがある場合、スキーマを管理し、データ保存プロセス中にデータ品質を検証し、一致しないデータを排除できます。さらに、メタデータを管理した後、ACID トランザクションも実装できます。これはトランザクションの特性です。管理レイヤーがない状態で並行操作が実行されると、複数の操作が相互に影響を与える可能性があります。たとえば、あるユーザーが問い合わせを行っている間に、別のユーザーが削除操作を実行する場合などです。トランザクションのサポートがあれば、このような状況を回避できます。トランザクションのサポートにより、各操作はスナップショットを生成し、すべての操作はスナップショットのシーケンスを生成します。これにより、時間の遡及、つまりタイムトラベルが可能になります。
データウェアハウス、データレイク、Delta Lake の主な機能の比較を下図に示します。Delta Lake はデータウェアハウスとデータレイクの利点を組み合わせ、管理レイヤーを導入して両者の欠点のほとんどを解決していることがわかります。
(3) Spark Streaming SQL & Delta に基づく CDC ソリューション
それでは、Spark Streaming SQL & Delta に基づく CDC ソリューションをどのように実装するかというトピックに戻りましょう。下図に示すように、まず binlog から Kafka へ送る点は以前の方法と同じです。以前の方法との違いは、Kafka 内の binlog を HBASE や KUDU に再生するのではなく、直接 Delta Lake に送る点です。このソリューションは使いやすく、追加の運用・メンテナンスが不要で、マージロジックも簡単に実装でき、ほぼリアルタイムのデータストリームとなります。
上記のソリューションの具体的な操作手順を下図に示します。本質的には、各ミニバッチを継続的に MERGE INTO でターゲットテーブルに適用していくことです。Spark Streaming のミニバッチスケジューリングは秒レベルで設定できるため、このソリューションでほぼリアルタイムのデータ同期を実現します。
プログラムの実際の実装過程で、いくつかの問題にも直面しました。最も重要なのは小ファイルの問題です。たとえば、5 秒ごとにバッチを実行すると、1 日で非常に多くのバッチが発生し、大量の小ファイルが生成される可能性があります。これはテーブルのクエリパフォーマンスに深刻な影響を与えます。小ファイルの問題に対するソリューションは以下の通りです。
スケジューリングバッチ間隔を延長する:リアルタイム要件がそれほど高くない場合、スケジューリングバッチ間隔を延長して小ファイルの生成頻度を下げることができます
小ファイルをマージする:小ファイルをマージして小ファイルの数を減らします。構文は以下の通りです:
OPTIMIZE WHERE where_clause]
アダプティブ実行:アダプティブ実行により、いくつかの小さな reduce タスクを組み合わせることができ、それによって小ファイルの数を減らすことができます
小ファイルマージの optimize トリガーについては、2 つの方法を実装しました。1 つ目は自動化された optimize で、各ミニバッチの実行後にマージが必要かどうかをチェックする方法です。必要なければ次のミニバッチに進みます。判断ルールは多くあり、小ファイルが一定数に達した、または合計ファイルサイズが一定サイズに達した場合などにマージを実行します。もちろん、マージ時にはすでに比較的大きなファイルを除外するなどの最適化も行っています。自動化された optimize 方式は、一定数のバッチが経過するたびにマージ操作が必要であり、データ取り込みに一定の影響を与える可能性があります。そのため、2 つ目の方法として、定期的に optimize を実行する方法があります。この方法はリアルタイムのデータ取り込みに影響を与えません。ただし、定期的に optimize を実行する方法には、トランザクション競合の問題、つまり optimize とストリームの間の競合があります。この場合、Delta 内部のトランザクションコミットメカニズムを最適化して、insert フローが失敗しないようにしました。optimize の前に update/delete が実行され、optimize が成功した場合、成功後にリトライプロセスを追加して、ストリームが途切れないようにする必要があります。
OPTIMIZE の実装も比較的複雑で、ビンパッキングメカニズムとアダプティブメカニズムを開発しました。達成される効果は、OPTIMIZE 後にすべてのファイル(最後を除く)がターゲットサイズ(128M など)に達することで、再パーティションの有無に関わらず同様です。
3. 今後の取り組み
今後、以下の側面が作業目標となります。
(1) 自動スキーマ検出
Delta Lake を使用するユーザーは、ビジネスデータにのみ接するのではなく、機械データにも接する可能性があります。多くのシナリオで、機械データのフィールドが変化する可能性があります。このシナリオのユーザーには、自動スキーマ検出メカニズムが急務です。次の段階での目標は、binlog 解析中に新しいフィールド、変更されたフィールドなどを自動検出し、Delta テーブルに反映させることです。
(2) ストリーミングマージパフォーマンス(Merge on Read)
前述の通り、Spark Streaming SQL & Delta の CDC ソリューションは本質的にストリーム処理を開始し、ミニバッチに従ってデータをターゲットテーブルにマージします。マージの実現は実際には結合です。テーブルが徐々に大きくなると、マージのパフォーマンスはますます悪化し、パフォーマンスに深刻な影響を与えます。この問題を解決する方法は、Merge on Read の方法を採用することで、HIVE の方法に似ており、次の目標となります。
(3) より使いやすいエクスペリエンス
前述の CDC ソリューションでも、ユーザーはある程度の専門知識といくつかの手動作業を行う必要があることがわかります。次のステップでは、より使いやすいエクスペリエンスを提供し、ユーザーの負担をさらに軽減することを目指します。
CDC は Change Data Capture の略称であり、変更データキャプチャを意味します。たとえば、初期段階ではツールを使用してビジネスデータをデータウェアハウスやデータレイクにインポートしていました。その後、データをインポートする際に、データの動的な変化を反映し、増分インポートを実行し、これらの変更データをできるだけ早くキャプチャして、より迅速に後続の分析をタイムリーに実行できるようにすることが求められるようになりました。CDC 技術は、このような変化するデータをキャプチャするのに役立ちます。
ビッグデータシナリオでは、一般的に使用されるツールは Sqoop です。これはバッチモードのツールであり、ビジネスライブラリからデータウェアハウスにデータをインポートするために使用できます。インポート前に、ビジネスライブラリのデータ内で時間的な変化を反映できるフィールドを選択し、タイムスタンプに基づいて変更されたデータをデータウェアハウスにインポートする必要がある点に注意が必要です。これが Sqoop の制限事項です。さらに、このツールには以下の欠点があります。
ソースライブラリに負荷をかける
遅延が大きい(呼び出し頻度に依存)
削除イベントを処理できず、ソースデータベースで削除されたデータをデータウェアハウスで同期的に削除できない
スキーマの変更に追随できない。ソースライブラリのスキーマが変更されると、データウェアハウスのテーブルモデルの再構築とインポートが必要になる
Sqoop を使用する以外に、binlog を使用してデータ同期を行う方法もあります。ソースライブラリは、挿入、更新、削除などの操作を実行する際に binlog を生成します。binlog を Kafka に入力し、Kafka から binlog を読み取って、1 つずつ解析した後に対応する操作を実行するだけです。ただし、この方法では、上流の頻繁な更新/削除操作に対応するために、ダウンストリームが比較的頻繁な更新/削除操作をサポートする必要があります。ここでは、KUDU または HBASE をターゲットストレージとして選択できます。ただし、KUDU と HBASE はデータウェアハウスではないため、フルデータの保存はできません。そのため、それらのデータは定期的に Hive にインポートする必要があります。下図に示す通りです。この方法には、複数コンポーネントの運用・メンテナンスの負荷が高い、マージロジックが複雑などの欠点がある点に注意が必要です。
2. Spark Streaming SQL & Delta に基づく CDC ソリューション
(1) Spark Streaming SQL
Spark Streaming SQL は、Alibaba Computing Platform Division の EMR チームが Spark Streaming を基に開発した SQL サポートであり、コミュニティ版には含まれていません。Spark Streaming SQL はこの CDC ソリューションに必須ではありませんが、特に SQL の使用に慣れているユーザーにとってより使いやすいものとなっています。そのため、EMR チームが Spark Streaming SQL のサポートを開発しました。下図に示すように、EMR の Spark Streaming SQL は、DDL、DML、SELECT など、多くの側面で SQL 構文のサポートを実装しています。以下にいくつかを紹介します。
(1) CREATE SCAN & CREATE STREAM
以下に示す例では、Kafka のテーブルから一部のデータを選択することが目標です。設計目標は、バッチとストリームの両方に対応できる限りサポートすることです。通常の SQL では、select は実際には読み取り操作を生成しますが、ここではバッチとストリーミングを区別するために、明示的な create scan が必要です。データソースからは、バッチ読み取りかストリーミング読み取りかを区別できないためです。バッチの場合は USING batch を使用し、ストリーミングの場合は USING stream を使用します。
バッチスキャンの場合、スキャンを作成した後、scan から直接 select でき、scan をテーブルとして扱うことができます。ただし、ストリーミングの場合、このスキャンを読み取るには多くのパラメータを設計する必要があります。ジョブを開始する必要があるためです。そのため、下図に示す create stream 構文があります。これは本質的に select 構文のカプセル化です。
(2) MERGE INTO
もう 1 つのコア構文は MERGE INTO であり、Delta Lake の CDC ソリューションで非常に重要な役割を果たします。MERGE INTO の構文はより複雑で、下図に示す通りです。MERGE INTO の mergeCondition は、ソーステーブルとターゲットテーブルの 1 対 1 対応である必要がある点に注意が必要です。そうでない場合、1 つのソースレコードが複数のターゲットレコードに対応すると、システムはどれを操作すればよいか判断できません。そのため、ここでは mergeCondition がプライマリキー結合、またはプライマリキー結合と同等の効果であることが実際には要求されます。
上で紹介した構文以外にも、DELAY、TUMBLING などの他の UDF も実装しており、Spark Streaming SQL をより簡単に使用できるようにしています。
(2) Delta Lake
データレイクは近年の注目技術です。初期段階では、比較的成熟したデータウェアハウスシステムが使用され、データは ETL を通じてデータウェアハウスにインポートされていました。データウェアハウスの典型的な用途は BI レポートなどの分析シナリオであり、シナリオは比較的限られていました。モバイルインターネット時代では、データソースがより豊富で多様化し、データ構造は構造化データに限定されず、データの用途も分析に限定されなくなったため、データレイクが登場しました。データを事前に処理しない、または簡単な処理のみ行ってデータレイクにインポートし、その後、スクリーニング、フィルタリング、変換などの transformation 操作を実行するため、データウェアハウス時代の ETL がデータレイク時代の ELT へと変化しました。
データレイクの典型的なアーキテクチャは、上位層に 1 つまたは複数の分析エンジンまたは他のコンピューティングフレームワーク、下位層に分散ストレージシステムを持つ構成です。下図の左側に示す通りです。ただし、この元のデータレイクの使用方法には、トランザクションサポートの欠如、データ品質検証の欠如など、管理が不足しており、すべてのデータ管理は完全に手動で保証されていました。
Delta Lake は、統一ストレージレイヤーの上に管理レイヤーを配置し、データレイクのデータを手動で管理する際の課題を解決します。管理レイヤーの追加により、まずメタデータ管理を導入できます。メタデータ管理があれば、データにスキーマがある場合、スキーマを管理し、データ保存プロセス中にデータ品質を検証し、一致しないデータを排除できます。さらに、メタデータを管理した後、ACID トランザクションも実装できます。これはトランザクションの特性です。管理レイヤーがない状態で並行操作が実行されると、複数の操作が相互に影響を与える可能性があります。たとえば、あるユーザーが問い合わせを行っている間に、別のユーザーが削除操作を実行する場合などです。トランザクションのサポートがあれば、このような状況を回避できます。トランザクションのサポートにより、各操作はスナップショットを生成し、すべての操作はスナップショットのシーケンスを生成します。これにより、時間の遡及、つまりタイムトラベルが可能になります。
データウェアハウス、データレイク、Delta Lake の主な機能の比較を下図に示します。Delta Lake はデータウェアハウスとデータレイクの利点を組み合わせ、管理レイヤーを導入して両者の欠点のほとんどを解決していることがわかります。
(3) Spark Streaming SQL & Delta に基づく CDC ソリューション
それでは、Spark Streaming SQL & Delta に基づく CDC ソリューションをどのように実装するかというトピックに戻りましょう。下図に示すように、まず binlog から Kafka へ送る点は以前の方法と同じです。以前の方法との違いは、Kafka 内の binlog を HBASE や KUDU に再生するのではなく、直接 Delta Lake に送る点です。このソリューションは使いやすく、追加の運用・メンテナンスが不要で、マージロジックも簡単に実装でき、ほぼリアルタイムのデータストリームとなります。
上記のソリューションの具体的な操作手順を下図に示します。本質的には、各ミニバッチを継続的に MERGE INTO でターゲットテーブルに適用していくことです。Spark Streaming のミニバッチスケジューリングは秒レベルで設定できるため、このソリューションでほぼリアルタイムのデータ同期を実現します。
プログラムの実際の実装過程で、いくつかの問題にも直面しました。最も重要なのは小ファイルの問題です。たとえば、5 秒ごとにバッチを実行すると、1 日で非常に多くのバッチが発生し、大量の小ファイルが生成される可能性があります。これはテーブルのクエリパフォーマンスに深刻な影響を与えます。小ファイルの問題に対するソリューションは以下の通りです。
スケジューリングバッチ間隔を延長する:リアルタイム要件がそれほど高くない場合、スケジューリングバッチ間隔を延長して小ファイルの生成頻度を下げることができます
小ファイルをマージする:小ファイルをマージして小ファイルの数を減らします。構文は以下の通りです:
OPTIMIZE WHERE where_clause]
アダプティブ実行:アダプティブ実行により、いくつかの小さな reduce タスクを組み合わせることができ、それによって小ファイルの数を減らすことができます
小ファイルマージの optimize トリガーについては、2 つの方法を実装しました。1 つ目は自動化された optimize で、各ミニバッチの実行後にマージが必要かどうかをチェックする方法です。必要なければ次のミニバッチに進みます。判断ルールは多くあり、小ファイルが一定数に達した、または合計ファイルサイズが一定サイズに達した場合などにマージを実行します。もちろん、マージ時にはすでに比較的大きなファイルを除外するなどの最適化も行っています。自動化された optimize 方式は、一定数のバッチが経過するたびにマージ操作が必要であり、データ取り込みに一定の影響を与える可能性があります。そのため、2 つ目の方法として、定期的に optimize を実行する方法があります。この方法はリアルタイムのデータ取り込みに影響を与えません。ただし、定期的に optimize を実行する方法には、トランザクション競合の問題、つまり optimize とストリームの間の競合があります。この場合、Delta 内部のトランザクションコミットメカニズムを最適化して、insert フローが失敗しないようにしました。optimize の前に update/delete が実行され、optimize が成功した場合、成功後にリトライプロセスを追加して、ストリームが途切れないようにする必要があります。
OPTIMIZE の実装も比較的複雑で、ビンパッキングメカニズムとアダプティブメカニズムを開発しました。達成される効果は、OPTIMIZE 後にすべてのファイル(最後を除く)がターゲットサイズ(128M など)に達することで、再パーティションの有無に関わらず同様です。
3. 今後の取り組み
今後、以下の側面が作業目標となります。
(1) 自動スキーマ検出
Delta Lake を使用するユーザーは、ビジネスデータにのみ接するのではなく、機械データにも接する可能性があります。多くのシナリオで、機械データのフィールドが変化する可能性があります。このシナリオのユーザーには、自動スキーマ検出メカニズムが急務です。次の段階での目標は、binlog 解析中に新しいフィールド、変更されたフィールドなどを自動検出し、Delta テーブルに反映させることです。
(2) ストリーミングマージパフォーマンス(Merge on Read)
前述の通り、Spark Streaming SQL & Delta の CDC ソリューションは本質的にストリーム処理を開始し、ミニバッチに従ってデータをターゲットテーブルにマージします。マージの実現は実際には結合です。テーブルが徐々に大きくなると、マージのパフォーマンスはますます悪化し、パフォーマンスに深刻な影響を与えます。この問題を解決する方法は、Merge on Read の方法を採用することで、HIVE の方法に似ており、次の目標となります。
(3) より使いやすいエクスペリエンス
前述の CDC ソリューションでも、ユーザーはある程度の専門知識といくつかの手動作業を行う必要があることがわかります。次のステップでは、より使いやすいエクスペリエンスを提供し、ユーザーの負担をさらに軽減することを目指します。
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
