すべてのプロダクト
Search
ドキュメントセンター

Elasticsearch:DataWorks を使用した MaxCompute データの Alibaba Cloud Elasticsearch への同期

最終更新日:Aug 21, 2026

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 クラスターにのみ同期できます。セルフマネージドの Elasticsearch クラスターはサポートされていません。

  • MaxCompute プロジェクト、Alibaba Cloud Elasticsearch クラスター、および DataWorks ワークスペースは、同じリージョンにある必要があります。

  • Alibaba Cloud Elasticsearch クラスター、MaxCompute プロジェクト、および DataWorks ワークスペースは、同じタイムゾーンである必要があります。そうでない場合、時間関連のデータを同期するときにタイムゾーンの不一致が発生する可能性があります。

課金

手順

手順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/1

    • category:製品カテゴリ、例:アウターウェア、生鮮食品、電化製品、バスルーム

    • brand:ブランド名、例:ブランドAからブランドG

    • buyer_id:購入者ID、例:user1からuser13

    • trans_num:トランザクション数

    • trans_amount:トランザクション額

    • click_cnt:クリック数

    • pt:パーティションフィールド、値は 1

手順2:専用リソースグループの購入と設定

Data Integration 用の専用リソースグループを購入し、VPC およびワークスペースに関連付けます。専用リソースグループは、高速で安定したデータ転送を保証します。

  1. DataWorks コンソールにログインします。

  2. 上部のメニューでリージョンを選択します。左側のメニューで、リソースグループ をクリックします。

  3. 専用リソースグループ タブで、レガシーリソースグループの作成 > Data Integration リソースグループ を選択します。

  4. DataWorks Exclusive Resources (Subscription) 購入ページで、[専用リソースタイプ] を 専用 Data Integration リソース に設定し、リソースグループの名前を入力して、今すぐ購入 をクリックします。

    設定の詳細については、「手順1:リソースグループの購入」をご参照ください。

  5. 作成した専用リソースグループを見つけ、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アドレス許可リストの設定」をご参照ください。

  6. ページの左上隅にある戻るアイコンをクリックして、Resource List ページに戻ります。

  7. 作成した専用リソースグループを見つけ、Actions 列の 所属ワークスペースの変更 をクリックして、ターゲットワークスペースに関連付けます。

    詳細については、「手順2:ワークスペースの関連付け」をご参照ください。

手順3:データソースの追加

Data Integration で、MaxCompute と Alibaba Cloud Elasticsearch をデータソースとして追加します。

  1. Data Integration ページに移動します。

    1. DataWorks コンソールにログインします。

    2. 左側のメニューで、ワークスペース をクリックします。

    3. ターゲットワークスペースを見つけ、Actions 列で ショートカット > Data Integration を選択します。

  2. 左側のメニューで、データソース をクリックします。

  3. MaxCompute データソースを追加します。

    1. ソースインスタンス ページで、ソースインスタンスの追加 をクリックします。

    2. ソースインスタンスの追加 ページで、MaxCompute データソースタイプを見つけて選択します。

    3. [MaxCompute データソースを追加] ダイアログボックスで、基本情報 セクションのパラメーターを設定します。

      詳細については、「MaxComputeデータソースの設定」をご参照ください。

    4. 接続設定 セクションで、接続テスト をクリックします。「[到達可能]」というステータスは、接続が成功したことを示します。

    5. 完了 をクリックします。

  4. 同じ手順に従って、Elasticsearch データソースを追加します。詳細については、「Elasticsearchデータソースの設定」をご参照ください。

手順4:バッチ同期タスクの設定と実行

バッチ同期タスクは、専用リソースグループを使用して実行されます。リソースグループは、Data Integration のデータソースからデータを取得し、Alibaba Cloud Elasticsearch クラスターにデータを書き込みます。

説明
  1. DataWorks の データ開発 ページに移動します。

    1. DataWorks コンソールにログインします。

    2. 左側のメニューで、ワークスペース をクリックします。

    3. ターゲットワークスペースを見つけ、Actions 列で ショートカット > データ開発 を選択します。

  2. オフライン同期ノードを作成します。

    1. Data Development (image) タブで、新規作成 > 業務フローの作成 を選択します。

    2. 作成したワークフローを右クリックし、ノードの作成 > Data Integration > オフライン同期 を選択します。

    3. ノードの作成 ダイアログボックスで、ノード名を入力し、OK をクリックします。

  3. ネットワークとリソースグループを設定します。

    1. データソース セクションで、データソース を MaxCompute(ODPS) に設定し、ソースインスタンス名 でソースデータソースを選択します。

    2. リソースグループ セクションで、専用リソースグループを選択します。

    3. データ宛先 セクションで、データ宛先 を Elasticsearch に設定し、ソースインスタンス名 で宛先データソースを選択します。

  4. Next step をクリックします。

  5. タスクを設定します。

    1. データソース セクションで、ソーステーブルを選択します。

    2. データ宛先 セクションで、宛先のパラメーターを設定します。

    3. フィールドマッピング セクションで、ソースフィールド を ターゲットフィールド にマッピングします。

    4. チャンネル制御 セクションで、チャネルパラメーターを設定します。

    設定の詳細については、「ウィザードモードでの同期タスクの設定」をご参照ください。

  6. タスクを実行します。

    1. (任意) タスクのスケジューリングプロパティを設定します。右側のペインで スケジューリング設定 をクリックし、必要に応じてスケジューリングパラメーターを設定します。パラメーターの詳細については、「スケジューリング設定」をご参照ください。

    2. ノード領域の右上隅にある保存アイコンをクリックして、タスクを保存します。

    3. ノード領域の右上隅にあるコミットアイコンをクリックします。

      タスクのスケジューリングプロパティを設定した場合、タスクはスケジュールされた間隔で自動的に実行されます。ノード領域の右上隅にある実行アイコンをクリックして、タスクをすぐに実行することもできます。

      実行ログのメッセージ 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 コンソールで、同期されたデータを表示し、クエリを実行します。

  1. 対象の Alibaba Cloud Elasticsearch クラスターの Kibana コンソールにログインします。

    詳細については、「Kibanaコンソールへのログイン」をご参照ください。

  2. Kibana ページの左上隅にある menu.png アイコンをクリックし、[開発ツール] を選択します。

  3. [コンソール] で、次のコマンドを実行して同期されたデータを表示します。

    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
            ]
          }
        ]
      }
    }
  4. 次のコマンドを実行して、ドキュメント内の category フィールドと brand フィールドを検索します。

    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {} },
    "_source": ["category", "brand"]
    }
  5. 次のコマンドを実行して、category が 「生鮮」 であるドキュメントを検索します。

    POST /odps_index/_search?pretty
    {
    "query": { "match": {"category":"生鮮"} }
    }
  6. 次のコマンドを実行して、trans_num フィールドでドキュメントをソートします。

    POST /odps_index/_search?pretty
    {
    "query": { "match_all": {} },
    "sort": { "trans_num": { "order": "desc" } }
    }

    その他のコマンドとアクセス方法については、「Elastic.coヘルプセンター」をご参照ください。