Tablestore+Delta Lake
背景紹介
近年、HTAP (ハイブリッドトランザクション分析処理) の普及が進んでいます。ストレージとコンピューティングを組み合わせることで、従来の大規模な構造化データ分析だけでなく、高速なトランザクション更新の書き込みにも対応できます。これはデータ集約型システム向けの実証済みアーキテクチャです。
Tablestore は、Alibaba Cloud が独自に開発した NoSQL マルチモデルデータベースで、PB レベルのストレージ、数千万 TPS、ミリ秒レベルのレイテンシを実現する大規模構造化データのストレージと高速クエリ・分析サービスを提供します。Tablestore の基盤エンジンの力を借りることで、OLTP シナリオの要件を十分に満たすことができます。Delta Lake は、デルタ対応のデータレイクに似ており、ベースデータを列指向ストレージで保存し、新規デルタデータを行形式で保存することで、データ操作の ACID と CRUD をサポートし、Spark のビッグデータエコシステムと完全に互換性があります。Delta Lake と Spark エコシステムを組み合わせることで、OLAP シナリオの要件を十分に満たすことができます。下図に Tablestore と Delta Lake を組み合わせた HTAP シナリオの概要論理構造図を示します。構造化ビッグデータ分析プラットフォームの設計に関する詳細は、「Structrured Big Data Analysis Platform Design」の記事を参照してください。
事前準備
1:Alibaba Cloud E-MapReduce コンソールにログインする
2:Hadoop クラスターを作成する (既に作成済みの場合は省略してください)
3:E-MapReduce クラスターに Tablestore インスタンスをデプロイする
4:同じ VPC 環境に配置する
ステップ 1 Tablestore ソーステーブルを作成する
有効化の手順については、公式ドキュメントを参照してください。このデモで作成するテーブルの名前は Source です。テーブルのスキーマは下図に示す通りで、PKString と PkInt の 2 つの主キーがあり、データ型はそれぞれ String 型と Integer 型です。
テーブル Source に増分チャネルを作成します。下図に示すように、チャネルリストにチャネル名、ID、タイプが表示されます。
技術的な補足:
トンネルサービス (Tunnel Service) は、Tablestore のデータインターフェイスに基づく完全増分統合サービスで、以下の 3 種類のトンネルを提供します。
全量:データテーブル内の履歴ストックデータの消費と処理
増分:データテーブル内の新規追加データの消費と処理
全量 + 増分:まずデータテーブルの履歴ストックデータをすべて消費し、その後に新規追加データを消費する
トンネルサービスの詳細については、Tablestore 公式サイトのドキュメントを参照してください。
ステップ 2 関連する JAR パッケージを取得して Hadoop クラスターにアップロードする
環境に必要な JAR パッケージを取得します。
クラスター管理ページで、作成した Hadoop クラスターのクラスター ID をクリックし、クラスターとサービスの管理ページに入ります。
左側のナビゲーションツリーでホストリストを選択し、右側で Hadoop クラスターの emr-header-1 ホストの IP 情報を確認します。
SSH クライアントでコマンドウィンドウを開き、Hadoop クラスターの emr-header-1 ホストにログインします。
すべての JAR パッケージを emr-header-1 ノードの任意のディレクトリにアップロードします。
ステップ 3 Spark Streaming ジョブを実行する
emr デモを基に修正したコードを例に、コンパイルして JAR パッケージを生成し、Hadoop クラスターの emr-header-1 ホストにアップロードします (ステップ 2 を参照)。コードの変更が大きいため、この記事には完全なコードは含まれていません。今後、emr デモの公式プロジェクトに統合される予定です。
この例では Tablestore テーブルをデータソースとして使用し、Tablestore CDC テクノロジー、Tablestore Streaming Source、Delta Sink を組み合わせることで、Tablestore から Delta Lake への完全なデータ連携を実演します。
次のコマンドで Spark Streaming ジョブを開始し、Tablestore Source テーブルのデータを Delta Lake Table にリアルタイムで同期するリスナープログラムを起動します。
各パラメーターの説明は次のとおりです。
ステップ 4 データ CRUD の例
まず、Tablestore に 2 行のデータを挿入します。この例では、2 つの主キー (PkString、PkInt) と 6 つの属性列 (col1、col2、col3、timestamp、col5、col6) を含む 8 列の同期列を作成します。Tablestore はスキーマフリー構造のため、属性列を自由に追加でき、Tablestore の Spark Source が自動的に属性列をフィルター処理します。次の 2 つの図に示すように、2 行のデータを挿入すると、Delta Table でも即座に 2 行を読み取ることができ、データの一貫性が保たれていることを確認できます。
次に、Tablestore で行の更新と挿入の操作を実行します。次の 2 つの図に示すように、マイクロバッチによるデータ同期の短い待機後、Tablestore のデータ変更がリアルタイムで Delta Table に反映されます。
Tablestore のデータをすべて削除します。次の 2 つの図に示すように、Delta Table も同期して空になります。
クラスター上では、Delta Table はデフォルトで HDFS に保存されます。下図に示すように、_delta_log ディレクトリに保存される JSON ファイルがトランザクションログであり、Parquet 形式のファイルが基盤データファイルです。
近年、HTAP (ハイブリッドトランザクション分析処理) の普及が進んでいます。ストレージとコンピューティングを組み合わせることで、従来の大規模な構造化データ分析だけでなく、高速なトランザクション更新の書き込みにも対応できます。これはデータ集約型システム向けの実証済みアーキテクチャです。
Tablestore は、Alibaba Cloud が独自に開発した NoSQL マルチモデルデータベースで、PB レベルのストレージ、数千万 TPS、ミリ秒レベルのレイテンシを実現する大規模構造化データのストレージと高速クエリ・分析サービスを提供します。Tablestore の基盤エンジンの力を借りることで、OLTP シナリオの要件を十分に満たすことができます。Delta Lake は、デルタ対応のデータレイクに似ており、ベースデータを列指向ストレージで保存し、新規デルタデータを行形式で保存することで、データ操作の ACID と CRUD をサポートし、Spark のビッグデータエコシステムと完全に互換性があります。Delta Lake と Spark エコシステムを組み合わせることで、OLAP シナリオの要件を十分に満たすことができます。下図に Tablestore と Delta Lake を組み合わせた HTAP シナリオの概要論理構造図を示します。構造化ビッグデータ分析プラットフォームの設計に関する詳細は、「Structrured Big Data Analysis Platform Design」の記事を参照してください。
事前準備
1:Alibaba Cloud E-MapReduce コンソールにログインする
2:Hadoop クラスターを作成する (既に作成済みの場合は省略してください)
3:E-MapReduce クラスターに Tablestore インスタンスをデプロイする
4:同じ VPC 環境に配置する
ステップ 1 Tablestore ソーステーブルを作成する
有効化の手順については、公式ドキュメントを参照してください。このデモで作成するテーブルの名前は Source です。テーブルのスキーマは下図に示す通りで、PKString と PkInt の 2 つの主キーがあり、データ型はそれぞれ String 型と Integer 型です。
テーブル Source に増分チャネルを作成します。下図に示すように、チャネルリストにチャネル名、ID、タイプが表示されます。
技術的な補足:
トンネルサービス (Tunnel Service) は、Tablestore のデータインターフェイスに基づく完全増分統合サービスで、以下の 3 種類のトンネルを提供します。
全量:データテーブル内の履歴ストックデータの消費と処理
増分:データテーブル内の新規追加データの消費と処理
全量 + 増分:まずデータテーブルの履歴ストックデータをすべて消費し、その後に新規追加データを消費する
トンネルサービスの詳細については、Tablestore 公式サイトのドキュメントを参照してください。
ステップ 2 関連する JAR パッケージを取得して Hadoop クラスターにアップロードする
環境に必要な JAR パッケージを取得します。
クラスター管理ページで、作成した Hadoop クラスターのクラスター ID をクリックし、クラスターとサービスの管理ページに入ります。
左側のナビゲーションツリーでホストリストを選択し、右側で Hadoop クラスターの emr-header-1 ホストの IP 情報を確認します。
SSH クライアントでコマンドウィンドウを開き、Hadoop クラスターの emr-header-1 ホストにログインします。
すべての JAR パッケージを emr-header-1 ノードの任意のディレクトリにアップロードします。
ステップ 3 Spark Streaming ジョブを実行する
emr デモを基に修正したコードを例に、コンパイルして JAR パッケージを生成し、Hadoop クラスターの emr-header-1 ホストにアップロードします (ステップ 2 を参照)。コードの変更が大きいため、この記事には完全なコードは含まれていません。今後、emr デモの公式プロジェクトに統合される予定です。
この例では Tablestore テーブルをデータソースとして使用し、Tablestore CDC テクノロジー、Tablestore Streaming Source、Delta Sink を組み合わせることで、Tablestore から Delta Lake への完全なデータ連携を実演します。
次のコマンドで Spark Streaming ジョブを開始し、Tablestore Source テーブルのデータを Delta Lake Table にリアルタイムで同期するリスナープログラムを起動します。
各パラメーターの説明は次のとおりです。
ステップ 4 データ CRUD の例
まず、Tablestore に 2 行のデータを挿入します。この例では、2 つの主キー (PkString、PkInt) と 6 つの属性列 (col1、col2、col3、timestamp、col5、col6) を含む 8 列の同期列を作成します。Tablestore はスキーマフリー構造のため、属性列を自由に追加でき、Tablestore の Spark Source が自動的に属性列をフィルター処理します。次の 2 つの図に示すように、2 行のデータを挿入すると、Delta Table でも即座に 2 行を読み取ることができ、データの一貫性が保たれていることを確認できます。
次に、Tablestore で行の更新と挿入の操作を実行します。次の 2 つの図に示すように、マイクロバッチによるデータ同期の短い待機後、Tablestore のデータ変更がリアルタイムで Delta Table に反映されます。
Tablestore のデータをすべて削除します。次の 2 つの図に示すように、Delta Table も同期して空になります。
クラスター上では、Delta Table はデフォルトで HDFS に保存されます。下図に示すように、_delta_log ディレクトリに保存される JSON ファイルがトランザクションログであり、Parquet 形式のファイルが基盤データファイルです。
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
