MaxCompute (ODPS) の大規模なデータセットで、情報の取得、多次元クエリの実行、または統計分析を行うには、Alibaba Cloud Elasticsearch を使用します。このトピックでは、DataWorks の Data Integration を使用して、MaxCompute から Alibaba Cloud Elasticsearch クラスターに大量のデータを数分で同期する方法について説明します。
背景
DataWorks は、ビッグデータエンジンを基盤とするエンドツーエンドのビッグデータ開発およびガバナンスプラットフォームです。 データ開発、タスクスケジューリング、データ管理などの機能が統合されています。 DataWorks の同期タスクを使用すると、さまざまなデータソースから Alibaba Cloud Elasticsearch にデータを迅速に同期できます。
-
サポートされているデータソース:
-
Alibaba Cloud データベース (MySQL、PostgreSQL、SQL Server、MongoDB、HBase)
-
Alibaba Cloud PolarDB-X (DRDS からアップグレード)
-
Alibaba Cloud MaxCompute
-
Alibaba Cloud OSS
-
Alibaba Cloud Tablestore
-
HDFS、Oracle、FTP、DB2、およびその他のサポートされているデータベースタイプのセルフマネージドバージョン
-
-
シナリオ:
-
ビッグデータの Alibaba Cloud Elasticsearch へのオフライン同期。 データベース全体または特定のテーブルのすべてのデータを同期できます。 詳細については、「MySQL データベース全体の Elasticsearch へのオフライン同期」をご参照ください。
-
ビッグデータの Alibaba Cloud Elasticsearch へのリアルタイム同期。 この方法は、完全同期と増分同期の両方をサポートしています。 詳細については、「MySQL データベース全体の Elasticsearch へのリアルタイム同期」をご参照ください。
-
前提条件
-
MaxCompute プロジェクトを作成済みであること。詳細については、「MaxComputeプロジェクトの作成」をご参照ください。
-
Alibaba Cloud Elasticsearch クラスターを作成し、その自動インデックス作成機能を有効にしてあること。詳細については、「Alibaba Cloud Elasticsearchクラスターの作成」および「YMLファイルの設定」をご参照ください。
-
DataWorks ワークスペースを作成済みであること。詳細については、「ワークスペースの作成」をご参照ください。
-
データは Alibaba Cloud Elasticsearch クラスターにのみ同期できます。セルフマネージドの Elasticsearch クラスターはサポートされていません。
-
MaxCompute プロジェクト、Alibaba Cloud Elasticsearch クラスター、および DataWorks ワークスペースは、同じリージョンにある必要があります。
-
Alibaba Cloud Elasticsearch クラスター、MaxCompute プロジェクト、および DataWorks ワークスペースは、同じタイムゾーンである必要があります。そうでない場合、時間関連のデータを同期するときにタイムゾーンの不一致が発生する可能性があります。
課金
-
Alibaba Cloud Elasticsearch インスタンスの料金については、「Elasticsearch の課金項目」をご参照ください。
-
Data Integration リソースグループの料金については、「リソースグループの料金」をご参照ください。
手順
手順1:ソースデータの準備
MaxCompute テーブルを作成し、テストデータをテーブルにインポートします。詳細については、「テーブルの作成」および「データのインポート」をご参照ください。
このトピックでは、次のテーブルスキーマとサンプルデータを使用します。
-
テーブルスキーマ
テーブルには 7 つのフィールドと 1 つのパーティションフィールドが含まれています。フィールドは次のように定義されます。
create_time(string、プライマリキー)category(string)brand(string)buyer_id(string)trans_num(bigint)trans_amount(double)click_cnt(bigint)
パーティションフィールドは
pt(bigint) です。 -
サンプルデータ
ソーステーブルには次のフィールドが含まれています。
create_time:トランザクション日、例:2020/6/1category:製品カテゴリ、例:アウターウェア、生鮮食品、電化製品、バスルームbrand:ブランド名、例:ブランドAからブランドGbuyer_id:購入者ID、例:user1からuser13trans_num:トランザクション数trans_amount:トランザクション額click_cnt:クリック数pt:パーティションフィールド、値は 1
手順2:専用リソースグループの購入と設定
Data Integration 用の専用リソースグループを購入し、VPC およびワークスペースに関連付けます。専用リソースグループは、高速で安定したデータ転送を保証します。
-
DataWorks コンソールにログインします。
-
上部のメニューでリージョンを選択します。左側のメニューで、リソースグループ をクリックします。
-
専用リソースグループ タブで、 を選択します。
-
DataWorks Exclusive Resources (Subscription) 購入ページで、[専用リソースタイプ] を 専用 Data Integration リソース に設定し、リソースグループの名前を入力して、今すぐ購入 をクリックします。
設定の詳細については、「手順1:リソースグループの購入」をご参照ください。
-
作成した専用リソースグループを見つけ、Actions 列の ネットワーク設定 をクリックして VPC に関連付けます。詳細については、「VPCのバインド」をご参照ください。
説明このトピックでは、Data Integration の専用リソースグループを使用して VPC 経由でデータを同期します。インターネット経由でデータを同期する方法については、「許可リストの設定」をご参照ください。
データ同期を有効にするには、専用リソースグループが Alibaba Cloud Elasticsearch クラスターが存在する VPC に接続されている必要があります。したがって、専用リソースグループを Alibaba Cloud Elasticsearch クラスターの Virtual Private Cloud (VPC)、Zone、および vSwitch に関連付ける必要があります。この情報を表示するには、「Elasticsearchクラスターの基本情報の表示」をご参照ください。
重要専用リソースグループを VPC に関連付けた後、Alibaba Cloud Elasticsearch クラスターの VPC プライベート IP 許可リストに vSwitch CIDR ブロック を追加する必要があります。詳細については、「ElasticsearchクラスターのパブリックまたはプライベートIPアドレス許可リストの設定」をご参照ください。
-
ページの左上隅にある戻るアイコンをクリックして、Resource List ページに戻ります。
-
作成した専用リソースグループを見つけ、Actions 列の 所属ワークスペースの変更 をクリックして、ターゲットワークスペースに関連付けます。
詳細については、「手順2:ワークスペースの関連付け」をご参照ください。
手順3:データソースの追加
Data Integration で、MaxCompute と Alibaba Cloud Elasticsearch をデータソースとして追加します。
-
Data Integration ページに移動します。
-
DataWorks コンソールにログインします。
-
左側のメニューで、ワークスペース をクリックします。
-
ターゲットワークスペースを見つけ、Actions 列で を選択します。
-
-
左側のメニューで、データソース をクリックします。
-
MaxCompute データソースを追加します。
-
ソースインスタンス ページで、ソースインスタンスの追加 をクリックします。
-
ソースインスタンスの追加 ページで、MaxCompute データソースタイプを見つけて選択します。
-
[MaxCompute データソースを追加] ダイアログボックスで、基本情報 セクションのパラメーターを設定します。
詳細については、「MaxComputeデータソースの設定」をご参照ください。
-
接続設定 セクションで、接続テスト をクリックします。「[到達可能]」というステータスは、接続が成功したことを示します。
-
完了 をクリックします。
-
-
同じ手順に従って、Elasticsearch データソースを追加します。詳細については、「Elasticsearchデータソースの設定」をご参照ください。
手順4:バッチ同期タスクの設定と実行
バッチ同期タスクは、専用リソースグループを使用して実行されます。リソースグループは、Data Integration のデータソースからデータを取得し、Alibaba Cloud Elasticsearch クラスターにデータを書き込みます。
-
バッチ同期タスクは、ウィザードモードまたはスクリプトモードで設定できます。このトピックでは、ウィザードモードについて説明します。詳細については、「スクリプトモードでの同期タスクの設定」および「Elasticsearch Writer」をご参照ください。
-
以下の手順は レガシ Data Development (DataStudio) ページで実行されます。
-
DataWorks の データ開発 ページに移動します。
-
DataWorks コンソールにログインします。
-
左側のメニューで、ワークスペース をクリックします。
-
ターゲットワークスペースを見つけ、Actions 列で を選択します。
-
-
オフライン同期ノードを作成します。
-
Data Development (
) タブで、 を選択します。 -
作成したワークフローを右クリックし、 を選択します。
-
ノードの作成 ダイアログボックスで、ノード名を入力し、OK をクリックします。
-
-
ネットワークとリソースグループを設定します。
-
データソース セクションで、データソース を MaxCompute(ODPS) に設定し、ソースインスタンス名 でソースデータソースを選択します。
-
リソースグループ セクションで、専用リソースグループを選択します。
-
データ宛先 セクションで、データ宛先 を Elasticsearch に設定し、ソースインスタンス名 で宛先データソースを選択します。
-
-
Next step をクリックします。
-
タスクを設定します。
-
データソース セクションで、ソーステーブルを選択します。
-
データ宛先 セクションで、宛先のパラメーターを設定します。
-
フィールドマッピング セクションで、ソースフィールド を ターゲットフィールド にマッピングします。
-
チャンネル制御 セクションで、チャネルパラメーターを設定します。
設定の詳細については、「ウィザードモードでの同期タスクの設定」をご参照ください。
-
-
タスクを実行します。
-
(任意) タスクのスケジューリングプロパティを設定します。右側のペインで スケジューリング設定 をクリックし、必要に応じてスケジューリングパラメーターを設定します。パラメーターの詳細については、「スケジューリング設定」をご参照ください。
-
ノード領域の右上隅にある保存アイコンをクリックして、タスクを保存します。
-
ノード領域の右上隅にあるコミットアイコンをクリックします。
タスクのスケジューリングプロパティを設定した場合、タスクはスケジュールされた間隔で自動的に実行されます。ノード領域の右上隅にある実行アイコンをクリックして、タスクをすぐに実行することもできます。
実行ログのメッセージ
Shell run successfully!は、タスクが正常に実行されたことを示します。タスク実行ログのサンプルは次のとおりです。2023-10-31 16:52:35 INFO Exit code of the Shell command 0 2023-10-31 16:52:35 INFO --- Invocation of Shell command completed --- 2023-10-31 16:52:35 INFO Shell run successfully! 2023-10-31 16:52:35 INFO Current task status: FINISH 2023-10-31 16:52:35 INFO Cost time is: 33.106s
-
手順5:データ同期結果の確認
Kibana コンソールで、同期されたデータを表示し、クエリを実行します。
-
対象の Alibaba Cloud Elasticsearch クラスターの Kibana コンソールにログインします。
詳細については、「Kibanaコンソールへのログイン」をご参照ください。
-
Kibana ページの左上隅にある
アイコンをクリックし、[開発ツール] を選択します。 -
[コンソール] で、次のコマンドを実行して同期されたデータを表示します。
POST /odps_index/_search?pretty { "query": { "match_all": {}} }説明odps_indexは、データ同期スクリプトで設定したindexフィールドの値です。同期が成功すると、次の結果が返されます。
--- Response --- { "took" : 2, "timed_out" : false, "_shards" : { "total" : 1, "successful" : 1, "skipped" : 0, "failed" : 0 }, "hits" : { "total" : 13, "max_score" : null, "hits" : [ { "_index" : "odps_index", "_type" : "_doc", "_id" : "2020/6/7 8:00", "_score" : null, "_source" : { "trans_num" : 88, "click_cnt" : 80, "category" : "アウターウェア", "buyer_id" : "user7", "trans_amount" : 150.0, "brand" : "ブランドE" }, "sort" : [ 88 ] }, { "_index" : "odps_index", "_type" : "_doc", "_id" : "2020/6/11 8:00", "_score" : null, "_source" : { "trans_num" : 22, "click_cnt" : 70, "category" : "バスルーム", "buyer_id" : "user11", "trans_amount" : 4500.0, "brand" : "ブランドG" }, "sort" : [ 22 ] } ] } } -
次のコマンドを実行して、ドキュメント内の
categoryフィールドとbrandフィールドを検索します。POST /odps_index/_search?pretty { "query": { "match_all": {} }, "_source": ["category", "brand"] } -
次のコマンドを実行して、
categoryが「生鮮」であるドキュメントを検索します。POST /odps_index/_search?pretty { "query": { "match": {"category":"生鮮"} } } -
次のコマンドを実行して、
trans_numフィールドでドキュメントをソートします。POST /odps_index/_search?pretty { "query": { "match_all": {} }, "sort": { "trans_num": { "order": "desc" } } }その他のコマンドとアクセス方法については、「Elastic.coヘルプセンター」をご参照ください。