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

Elasticsearch:ES-Hadoop を使用した、Apache Spark から Alibaba Cloud Elasticsearch へのデータ読み書き

最終更新日:Sep 05, 2026

ES-Hadoop を使用して、Apache Spark から Alibaba Cloud Elasticsearch のデータを読み書きする方法について説明します。

Elasticsearch-Hadoop (ES-Hadoop) は Apache Spark と Alibaba Cloud Elasticsearch を橋渡しし、カスタムコネクターなしで Spark ジョブが Elasticsearch クラスターのデータを読み書きできるようにします。このチュートリアルでは、環境の準備、JSON レコードの Elasticsearch インデックスへの書き込み、データの読み取り、そして Kibana での結果の検証まで、エンドツーエンドの完全な例を解説します。

前提条件

開始する前に、以下が準備できていることを確認してください。

  • 自動インデックス作成機能が有効な Alibaba Cloud Elasticsearch クラスター (このチュートリアルでは V6.7.0 を使用)

    重要

    本番環境では、自動インデックス作成機能を無効にしてください。事前にインデックスを作成し、マッピングを設定することを推奨します。このチュートリアルでは、テスト目的でのみ自動インデックス作成を有効にしています。

  • Elasticsearch クラスターと同じ仮想プライベートクラウド (VPC) で実行され、以下の設定がされている E-MapReduce (EMR) クラスター。

    • EMR バージョン:EMR-3.29.0

    • 必須サービス:Spark 2.4.5 (他のサービスはデフォルト設定を維持)

  • デフォルトのホワイトリスト (0.0.0.0/0) を変更した場合、EMR クラスターのプライベート IP アドレスを Elasticsearch クラスターのプライベート IP アドレスホワイトリストに追加していること。EMR クラスターのプライベート IP アドレスを取得する方法については、「クラスターリストとクラスター詳細の表示」をご参照ください。ホワイトリストを更新する方法については、「Elasticsearch クラスターのパブリックまたはプライベート IP アドレスホワイトリストの設定」をご参照ください。

  • JDK 1.8.0 以降

Elasticsearch クラスターの作成手順については、「Alibaba Cloud Elasticsearch クラスターの作成」および「Elasticsearch クラスターへのアクセスと設定」をご参照ください。EMR クラスターの作成手順については、「クラスターの作成」をご参照ください。

仕組み

  1. JSON 形式のテストデータを EMR マスターノード上の Hadoop 分散ファイルシステム (HDFS) にアップロードします。

  2. ES-Hadoop と Spark の依存関係を持つ Java Maven プロジェクトをビルドし、書き込みおよび読み取りクラスをコンパイルします。

  3. コンパイルされたコードを JAR ファイルにパッケージ化し、spark-submit を使用して Spark ジョブとして送信します。

  4. ES-Hadoop は、各 Resilient Distributed Dataset (RDD) レコードをシリアル化し、REST API を使用して指定された Elasticsearch インデックスに書き込みます。

  5. Kibana の Dev Tools コンソールでクエリを実行し、書き込まれたデータを確認します。

テストデータの準備

  1. E-MapReduce コンソールにログインし、EMR マスターノードの IP アドレスを取得して、対応する ECS インスタンスに SSH で接続します。詳細については、「クラスターへのログイン」をご参照ください。

  2. http_log.txt という名前のファイルを作成し、以下の JSON レコードを書き込みます。

    {"id": 1, "name": "田中一郎", "birth": "1990-01-01", "addr": "No.969, wenyixi Rd, yuhang, hangzhou"}
    {"id": 2, "name": "鈴木二郎", "birth": "1991-01-01", "addr": "No.556, xixi Rd, xihu, hangzhou"}
    {"id": 3, "name": "佐藤三郎", "birth": "1992-01-01", "addr": "No.699 wangshang Rd, binjiang, hangzhou"}
  3. ファイルを HDFS にアップロードします。

    hadoop fs -put http_log.txt /tmp/hadoop-es

POM 依存関係の追加

Java Maven プロジェクトを作成し、以下の依存関係を pom.xml に追加します。各バージョンが対応する Alibaba Cloud サービスのバージョンと一致していることを確認してください。たとえば、elasticsearch-spark-20_2.11 は Elasticsearch クラスターのバージョンと、spark-core_2.12 は Spark および Scala のバージョンと一致させる必要があります。

<dependencies>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-core_2.12</artifactId>
        <version>2.4.5</version>
    </dependency>
    <dependency>
        <groupId>org.apache.spark</groupId>
        <artifactId>spark-sql_2.11</artifactId>
        <version>2.4.5</version>
    </dependency>
    <dependency>
        <groupId>org.elasticsearch</groupId>
        <artifactId>elasticsearch-spark-20_2.11</artifactId>
        <version>6.7.0</version>
    </dependency>
</dependencies>

Elasticsearch へのデータ書き込み

以下の Java コードは、HDFS から http_log.txt を読み取り、各行をドキュメントとして company/_doc インデックスに書き込みます。プレースホルダーの値を実際のクラスターエンドポイントと認証情報に置き換えてください。

import java.util.Map;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.spark.SparkConf;
import org.apache.spark.SparkContext;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.api.java.function.Function;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
import org.elasticsearch.spark.rdd.api.java.JavaEsSpark;
import org.spark_project.guava.collect.ImmutableMap;

public class SparkWriteEs {
    public static void main(String[] args) {
        SparkConf conf = new SparkConf();
        conf.setAppName("Es-write");
        conf.set("es.nodes", "es-cn-n6w1o1x0w001c****.elasticsearch.aliyuncs.com"); // (1)
        conf.set("es.net.http.auth.user", "elastic");
        conf.set("es.net.http.auth.pass", "xxxxxx");                                 // (2)
        conf.set("es.nodes.wan.only", "true");                                       // (3)
        conf.set("es.nodes.discovery", "false");                                     // (4)
        conf.set("es.input.use.sliced.partitions", "false");                         // (5)

        SparkSession ss = new SparkSession(new SparkContext(conf));
        final AtomicInteger employeesNo = new AtomicInteger(0);

        JavaRDD<Map<Object, ?>> javaRDD = ss.read().text("/tmp/hadoop-es/http_log.txt") // (6)
                .javaRDD().map((Function<Row, Map<Object, ?>>) row ->
                        ImmutableMap.of("employees" + employeesNo.getAndAdd(1), row.mkString()));

        JavaEsSpark.saveToEs(javaRDD, "company/_doc");
    }
}
  1. es.nodes は、Elasticsearch クラスターの内部エンドポイントです。 クラスターの [基本情報] ページから取得します。 詳細については、「クラスターの基本情報を表示する」をご参照ください。

  2. es.net.http.auth.pass — クラスター作成時に設定したパスワードです。xxxxxx を実際のパスワードに置き換えてください。

  3. es.nodes.wan.only — Alibaba Cloud Elasticsearch は仮想 IP を使用するため、true に設定します。これにより、ノード検出が無効になり、すべてのトラフィックが es.nodes のエンドポイント経由でルーティングされます。

  4. es.nodes.discovery — Alibaba Cloud Elasticsearch では false にする必要があります。クラスターの内部ノードアドレスは外部から到達できないため、ネイティブのノード検出メカニズムは機能しません。

  5. es.input.use.sliced.partitions — false に設定して、インデックスの先行読み込みフェーズをスキップし、クエリ効率を向上させます。先行読み込みフェーズは、実際のデータクエリよりも時間がかかることがあります。

  6. /tmp/hadoop-es/http_log.txt を、実際のテストデータの HDFS パスに置き換えてください。

重要

本番環境では elastic アカウントの使用を避けてください。パスワードをリセットすると、クラスターへのアクセスが一時的にブロックされる可能性があります。代わりに、Kibana で必要なロールを持つ専用のユーザーを作成してください。詳細については、「Elasticsearch X-Pack が提供する RBAC メカニズムを使用したアクセス制御の実装」をご参照ください。

Elasticsearch からのデータ読み取り

以下のコードは、company/_doc インデックスからすべてのドキュメントを読み取り、標準出力に出力します。

import org.apache.spark.SparkConf;
import org.apache.spark.api.java.JavaPairRDD;
import org.apache.spark.api.java.JavaSparkContext;
import org.elasticsearch.spark.rdd.api.java.JavaEsSpark;
import java.util.Map;

public class ReadES {
    public static void main(String[] args) {
        SparkConf conf = new SparkConf()
                .setAppName("readEs")
                .setMaster("local[*]")
                .set("es.nodes", "es-cn-n6w1o1x0w001c****.elasticsearch.aliyuncs.com")
                .set("es.port", "9200")
                .set("es.net.http.auth.user", "elastic")
                .set("es.net.http.auth.pass", "xxxxxx")
                .set("es.nodes.wan.only", "true")
                .set("es.nodes.discovery", "false")
                .set("es.input.use.sliced.partitions", "false")
                .set("es.resource", "company/_doc")
                .set("es.scroll.size", "500");

        JavaSparkContext sc = new JavaSparkContext(conf);
        JavaPairRDD<String, Map<String, Object>> rdd = JavaEsSpark.esRDD(sc);

        for (Map<String, Object> item : rdd.values().collect()) {
            System.out.println(item);
        }

        sc.stop();
    }
}

設定リファレンス

パラメーター

デフォルト

説明

es.nodes

localhost

Elasticsearch クラスターにアクセスするためのエンドポイント。最高のパフォーマンスを得るには、内部エンドポイントを使用してください。

es.port

9200

Elasticsearch クラスターにアクセスするためのポート。

es.net.http.auth.user

elastic

Elasticsearch 認証用のユーザー名。

es.net.http.auth.pass

/

指定されたユーザー名のパスワード。忘れた場合はリセットしてください。詳細については、「Elasticsearch クラスターのアクセスパスワードのリセット」をご参照ください。

es.nodes.wan.only

false

true の場合、ノード検出を無効にし、すべてのトラフィックを es.nodes 経由でルーティングします。Elasticsearch が仮想 IP を使用する場合に必要です。

es.nodes.discovery

true

false の場合、ES-Hadoop が追加のクラスターノードを検出するのを防ぎます。クラスターの内部ノードアドレスは外部から到達できないため、Alibaba Cloud Elasticsearch では false にする必要があります。

es.input.use.sliced.partitions

true

false の場合、インデックスの先行読み込みフェーズをスキップします。先行読み込みフェーズは実際のデータクエリよりも時間がかかることがあるため、クエリ効率を向上させるには false に設定してください。

es.index.auto.create

true

true の場合、インデックスが存在しない場合に ES-Hadoop が自動的にインデックスを作成します。本番環境ではこれを無効にし、代わりに明示的なマッピングでインデックスを作成してください。

es.resource

/

読み書き対象のインデックスとタイプを index/type 形式で指定します。

es.mapping.names

/

ソースデータと Elasticsearch インデックス間のフィールド名のマッピング。

es.scroll.size

(なし)

読み取り時にスクロールリクエストごとにフェッチされるドキュメントの数。

ES-Hadoop の設定オプションの完全なリストについては、「ES-Hadoop 設定リファレンス」をご参照ください。

Spark ジョブの送信

  1. コンパイルされたコードを JAR ファイルにパッケージ化し、EMR マスターノードまたは関連するゲートウェイクラスターにアップロードします。

  2. EMR クライアントで、以下のコマンドを実行してジョブを送信します。

    データの書き込み:

    cd /usr/lib/spark-current
    ./bin/spark-submit --master yarn --executor-cores 1 --class "SparkWriteEs" /usr/local/spark_es.jar

    データの読み取り:

    cd /usr/lib/spark-current
    ./bin/spark-submit --master yarn --executor-cores 1 --class "ReadES" /usr/local/spark_es.jar
    重要

    /usr/local/spark_es.jar を、JAR をアップロードした実際のパスに置き換えてください。

    読み取りジョブが完了すると、出力には Elasticsearch インデックスから取得された各ドキュメントが表示されます。

    20/11/09 14:52:19 INFO [Executor task launch worker for task 0] Executor: Adding file:/tmp/spark-383278cd-c582-4c1c-xxx                    /userFiles-8d64xxx
    20/11/09 14:52:19 INFO [Executor task launch worker for task 0] Executor: Finished task 0.0 in stage 0.0 (TID 0). 1274 bytes result sent to driver
    20/11/09 14:52:19 INFO [task-result-getter-0] TaskSetManager: Finished task 0.0 in stage 0.0 (TID 0) in 544 ms on localhost (executor driver) (1/1)
    20/11/09 14:52:19 INFO [task-result-getter-0] TaskSchedulerImpl: Removed TaskSet 0.0, whose tasks have all completed, from pool
    20/11/09 14:52:19 INFO [dag-scheduler-event-loop] DAGScheduler: ResultStage 0 (collect at ReadES.java:26) finished in 0.691 s
    20/11/09 14:52:19 INFO [main] DAGScheduler: Job 0 finished: collect at ReadES.java:26, took 0.761197 s
    {employees0={"id": 1, "name": "田中一郎", "birth": "1990-01-01", "addr": "No.969, wenyixi Rd, yuhang, hangzhou"}}
    {employees1={"id": 2, "name": "鈴木二郎", "birth": "1991-01-01", "addr": "No.556, xixi Rd, xihu, hangzhou"}}
    {employees2={"id": 3, "name": "佐藤三郎", "birth": "1992-01-01", "addr": "No.699 wangshang Rd, binjiang, hangzhou"}}
    20/11/09 14:52:19 INFO [main] SparkUI: Stopped Spark Web UI at http://xxx:4041
    20/11/09 14:52:19 INFO [dispatcher-event-loop-1] MapOutputTrackerMasterEndpoint: MapOutputTrackerMasterEndpoint stopped!
    20/11/09 14:52:19 INFO [main] MemoryStore: MemoryStore cleared
    20/11/09 14:52:19 INFO [main] BlockManager: BlockManager stopped
    20/11/09 14:52:19 INFO [main] BlockManagerMaster: BlockManagerMaster stopped
    20/11/09 14:52:19 INFO [dispatcher-event-loop-1] OutputCommitCoordinator$OutputCommitCoordinatorEndpoint: OutputCommitCoordinator stopped!
    20/11/09 14:52:19 INFO [main] SparkContext: Successfully stopped SparkContext
    20/11/09 14:52:19 INFO [pool-1-thread-1] ShutdownHookManager: Shutdown hook called
    20/11/09 14:52:19 INFO [pool-1-thread-1] ShutdownHookManager: Deleting directory /tmp/spark-383278cd-c582-4c1c-b4b6-xxx
    20/11/09 14:52:19 INFO [pool-1-thread-1] ShutdownHookManager: Deleting directory /tmp/spark-72515c56-2ba4-4b36-942c-xxx

結果の検証

  1. Elasticsearch クラスターの Kibana コンソールにログインします。詳細については、「Kibana コンソールへのログイン」をご参照ください。

  2. 左側のメニューで、[Dev Tools] をクリックします。

  3. [Console] タブで、以下のクエリを実行してデータが正常に書き込まれたことを確認します。

    GET company/_search
    {
      "query": {
        "match_all": {}
      }
    }

    レスポンスが成功すると、Spark ジョブによって書き込まれた 3 つすべてのドキュメントが返されます。

    GET company/_search
    {
      "query": {
        "match_all": {}
      }
    }
    
    --- Response ---
    
    {
      "took" : 0,
      "timed_out" : false,
      "_shards" : {
        "total" : 1,
        "successful" : 1,
        "skipped" : 0,
        "failed" : 0
      },
      "hits" : {
        "total" : 3,
        "max_score" : 1.0,
        "hits" : [
          {
            "_index" : "company",
            "_type" : "_doc",
            "_id" : "1MDGq3UBNM177Uedd0uo",
            "_score" : 1.0,
            "_source" : {
              "employees0" : "{\"id\": 1, \"name\": \"田中一郎\", \"birth\": \"1990-01-01\", \"addr\": \"No.969, wenyixi Rd, yuhang, hangzhou\"}"
            }
          },
          {
            "_index" : "company",
            "_type" : "_doc",
            "_id" : "1cDGq3UBNM177Uedd0uo",
            "_score" : 1.0,
            "_source" : {
              "employees1" : "{\"id\": 2, \"name\": \"鈴木二郎\", \"birth\": \"1991-01-01\", \"addr\": \"No.556, xixi Rd, xihu, hangzhou\"}"
            }
          },
          {
            "_index" : "company",
            "_type" : "_doc",
            "_id" : "1sDGq3UBNM177Uedd0uo",
            "_score" : 1.0,
            "_source" : {
              "employees2" : "{\"id\": 3, \"name\": \"佐藤三郎\", \"birth\": \"1992-01-01\", \"addr\": \"No.699 wangshang Rd, binjiang, hangzhou\"}"
            }
          }
        ]
      }
    }

次のステップ

ES-Hadoop は、Java RDD の書き込みだけでなく、さらに多くの機能をサポートしています。Spark と統合すると、Spark データセット、Spark Streaming、Scala、および Spark SQL も使用できます。詳細については、「Apache Spark のサポート」をご参照ください。