Quickly build Lambda big data analysis architecture on the cloud in 5 minutes
背景
Spark China コミュニティは、Alibaba Cloud EMR 技術交流グループおよび Tablestore 技術交流グループと共同で、技術ライブ配信を開催しました。ライブ配信のテーマは「大規模構造化データのリアルタイムコンピューティングと処理」で、Tablestore ベースのデータ変更のリアルタイムキャプチャとサブスクリプション機能、およびクラウド上での Lambda アーキテクチャの軽量な実装について紹介します。ライブ配信ではデモリンクを使用しており、本記事ではそのデモリンクの簡単な操作步骤を提供します。これにより、Alibaba Cloud 上でデモシーンと同様のアーキテクチャを構築し、リアルタイムおよびオフラインのデータ処理を実現できます。
デモシーンの紹介
このデモは E コマースの注文シナリオをシミュレートし、ストリームコンピューティングにより大画面の注文シーンを実現します。大量の注文をリアルタイムに投入し、10 秒ごとに注文統計の集約と取引金額統計を実行して、リアルタイムの大画面表示を行います。全体の注文の大画面サンプルは次のとおりです。
大画面には、Alibaba Cloud の DataV を使用して Tablestore データソースに接続し実現しています。次に、注文の元データから結果の大画面データまでの生成プロセスと操作步骤をご紹介します。
全体の背景の構造は概ね次のとおりです。
ECS またはローカルで注文ジェネレーターをシミュレートし、注文データを Tablestore にリアルタイムに投入します。
Tablestore コンソールでチャネルを作成します
EMR コンソールで Spark クラスターを購入します
最新の EMR SDK をダウンロードします
以下に示すテーブル作成ステートメントと SQL コマンドを実行してリアルタイム計算を行い、結果テーブルを Tablestore に書き戻します。
DATAV を通じて結果テーブルデータのリアルタイム大画面表示を行います
ステップ 1: Alibaba Cloud 公式サイトで Tablestore コンソールにログインし、インスタンスとテーブルを作成します
インスタンスの作成後、次のようなプライマリキースキーマでテーブルを作成できます。
クライアントインジェクションプログラムを起動してランダムにデータを書き込みます。サンプルデータは次のとおりです。
Tablestore はサーバーレス形式の製品です。ユーザーは容量や仕様を購入する必要はなく、ビジネスに応じて自動的に水平スケーリングされます。
ステップ 2: Alibaba Cloud 公式サイトで EMR コンソールにログインし、Spark クラスターを購入します
Spark のクラスターサイズはビジネスニーズに応じて柔軟に選択できます。実際に 3 ノードでテストしたところ、100 万件/秒のデータをリアルタイムに集約計算できました。
ステップ 3: EMR クラスターにログインしてジョブスクリプトを実行します
EMR のマスターノードにログインし、次のコマンドを実行してストリームタスクを開始します。
1. ストリーム SQL 対話を開始します
EMR 公式サイトから最新バージョンの EMR SDK (1.8) を取得します
streaming-sql --driver-class-path emr-datasources_shaded_2.11-1.8.0.jar --jars emr-datasources_shaded_2.11-1.8.0.jar --master yarn-client --num-executors 8 --executor -memory 2g --executor-cores 2
2. ストリーミングソーステーブルを作成します
DROP TABLE IF EXISTS ots_order_test;
CREATE TABLE ots_order_test
USING tablestore
OPTIONS(
endpoint="fill in the address of Tablestore VPC",
access.key.id="",
access.key.secret="",
instance.name="",
table.name="",
tunnel.id="Find the ID of the channel you want to consume in the Tablestore console",
catalog='{"columns": {"UserId": {"col": "UserId", "type": "string"}, "OrderId": {"col": "OrderId", "type": "string "},"price": {"cols": "price", "type": "long"}, "timestamp": {"cols": "timestamp", "type": "long"}}}'
);
3. ストリーミングシンクテーブルを作成します
DROP TABLE IF EXISTS ots_order_sink_test;
CREATE TABLE ots_order_sink_test
USING tablestore
OPTIONS(
endpoint="",
access.key.id="",
access.key.secret="",
instance.name="",
table.name="",
tunnel.id="",
catalog='{"columns": {"begin": {"col": "begin", "type": "string"},"end": {"col": "end", "type": "string "}, "count": {"col": "count", "type": "long"}, "totalPrice": {"col": "totalPrice", "type": "long"}}}'
);
4. ストリーミングジョブを作成します
CREATE SCAN ots_table_stream on ots_order_test USING STREAM OPTIONS ("maxoffsetsperchannel"="10000");
CREATE STREAM job1
options(
checkpointLocation='/tmp/spark/cp/test1',
outputMode='update'
)
insert into ots_order_sink_test
SELECT CAST(window.start AS String) AS begin, CAST(window.end AS String) AS end, count(*) AS count, sum(price) AS totalPrice FROM ots_table_stream GROUP BY window(to_timestamp(timestamp / 1000000000), " 10 seconds");
Spark China コミュニティは、Alibaba Cloud EMR 技術交流グループおよび Tablestore 技術交流グループと共同で、技術ライブ配信を開催しました。ライブ配信のテーマは「大規模構造化データのリアルタイムコンピューティングと処理」で、Tablestore ベースのデータ変更のリアルタイムキャプチャとサブスクリプション機能、およびクラウド上での Lambda アーキテクチャの軽量な実装について紹介します。ライブ配信ではデモリンクを使用しており、本記事ではそのデモリンクの簡単な操作步骤を提供します。これにより、Alibaba Cloud 上でデモシーンと同様のアーキテクチャを構築し、リアルタイムおよびオフラインのデータ処理を実現できます。
デモシーンの紹介
このデモは E コマースの注文シナリオをシミュレートし、ストリームコンピューティングにより大画面の注文シーンを実現します。大量の注文をリアルタイムに投入し、10 秒ごとに注文統計の集約と取引金額統計を実行して、リアルタイムの大画面表示を行います。全体の注文の大画面サンプルは次のとおりです。
大画面には、Alibaba Cloud の DataV を使用して Tablestore データソースに接続し実現しています。次に、注文の元データから結果の大画面データまでの生成プロセスと操作步骤をご紹介します。
全体の背景の構造は概ね次のとおりです。
ECS またはローカルで注文ジェネレーターをシミュレートし、注文データを Tablestore にリアルタイムに投入します。
Tablestore コンソールでチャネルを作成します
EMR コンソールで Spark クラスターを購入します
最新の EMR SDK をダウンロードします
以下に示すテーブル作成ステートメントと SQL コマンドを実行してリアルタイム計算を行い、結果テーブルを Tablestore に書き戻します。
DATAV を通じて結果テーブルデータのリアルタイム大画面表示を行います
ステップ 1: Alibaba Cloud 公式サイトで Tablestore コンソールにログインし、インスタンスとテーブルを作成します
インスタンスの作成後、次のようなプライマリキースキーマでテーブルを作成できます。
クライアントインジェクションプログラムを起動してランダムにデータを書き込みます。サンプルデータは次のとおりです。
Tablestore はサーバーレス形式の製品です。ユーザーは容量や仕様を購入する必要はなく、ビジネスに応じて自動的に水平スケーリングされます。
ステップ 2: Alibaba Cloud 公式サイトで EMR コンソールにログインし、Spark クラスターを購入します
Spark のクラスターサイズはビジネスニーズに応じて柔軟に選択できます。実際に 3 ノードでテストしたところ、100 万件/秒のデータをリアルタイムに集約計算できました。
ステップ 3: EMR クラスターにログインしてジョブスクリプトを実行します
EMR のマスターノードにログインし、次のコマンドを実行してストリームタスクを開始します。
1. ストリーム SQL 対話を開始します
EMR 公式サイトから最新バージョンの EMR SDK (1.8) を取得します
streaming-sql --driver-class-path emr-datasources_shaded_2.11-1.8.0.jar --jars emr-datasources_shaded_2.11-1.8.0.jar --master yarn-client --num-executors 8 --executor -memory 2g --executor-cores 2
2. ストリーミングソーステーブルを作成します
DROP TABLE IF EXISTS ots_order_test;
CREATE TABLE ots_order_test
USING tablestore
OPTIONS(
endpoint="fill in the address of Tablestore VPC",
access.key.id="",
access.key.secret="",
instance.name="",
table.name="",
tunnel.id="Find the ID of the channel you want to consume in the Tablestore console",
catalog='{"columns": {"UserId": {"col": "UserId", "type": "string"}, "OrderId": {"col": "OrderId", "type": "string "},"price": {"cols": "price", "type": "long"}, "timestamp": {"cols": "timestamp", "type": "long"}}}'
);
3. ストリーミングシンクテーブルを作成します
DROP TABLE IF EXISTS ots_order_sink_test;
CREATE TABLE ots_order_sink_test
USING tablestore
OPTIONS(
endpoint="",
access.key.id="",
access.key.secret="",
instance.name="",
table.name="",
tunnel.id="",
catalog='{"columns": {"begin": {"col": "begin", "type": "string"},"end": {"col": "end", "type": "string "}, "count": {"col": "count", "type": "long"}, "totalPrice": {"col": "totalPrice", "type": "long"}}}'
);
4. ストリーミングジョブを作成します
CREATE SCAN ots_table_stream on ots_order_test USING STREAM OPTIONS ("maxoffsetsperchannel"="10000");
CREATE STREAM job1
options(
checkpointLocation='/tmp/spark/cp/test1',
outputMode='update'
)
insert into ots_order_sink_test
SELECT CAST(window.start AS String) AS begin, CAST(window.end AS String) AS end, count(*) AS count, sum(price) AS totalPrice FROM ots_table_stream GROUP BY window(to_timestamp(timestamp / 1000000000), " 10 seconds");
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
