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 クラスターの作成手順については、「クラスターの作成」をご参照ください。
仕組み
JSON 形式のテストデータを EMR マスターノード上の Hadoop 分散ファイルシステム (HDFS) にアップロードします。
ES-Hadoop と Spark の依存関係を持つ Java Maven プロジェクトをビルドし、書き込みおよび読み取りクラスをコンパイルします。
コンパイルされたコードを JAR ファイルにパッケージ化し、
spark-submitを使用して Spark ジョブとして送信します。ES-Hadoop は、各 Resilient Distributed Dataset (RDD) レコードをシリアル化し、REST API を使用して指定された Elasticsearch インデックスに書き込みます。
Kibana の Dev Tools コンソールでクエリを実行し、書き込まれたデータを確認します。
テストデータの準備
E-MapReduce コンソールにログインし、EMR マスターノードの IP アドレスを取得して、対応する ECS インスタンスに SSH で接続します。詳細については、「クラスターへのログイン」をご参照ください。
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"}ファイルを 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");
}
}es.nodesは、Elasticsearch クラスターの内部エンドポイントです。 クラスターの [基本情報] ページから取得します。 詳細については、「クラスターの基本情報を表示する」をご参照ください。es.net.http.auth.pass— クラスター作成時に設定したパスワードです。xxxxxxを実際のパスワードに置き換えてください。es.nodes.wan.only— Alibaba Cloud Elasticsearch は仮想 IP を使用するため、trueに設定します。これにより、ノード検出が無効になり、すべてのトラフィックがes.nodesのエンドポイント経由でルーティングされます。es.nodes.discovery— Alibaba Cloud Elasticsearch ではfalseにする必要があります。クラスターの内部ノードアドレスは外部から到達できないため、ネイティブのノード検出メカニズムは機能しません。es.input.use.sliced.partitions—falseに設定して、インデックスの先行読み込みフェーズをスキップし、クエリ効率を向上させます。先行読み込みフェーズは、実際のデータクエリよりも時間がかかることがあります。/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();
}
}設定リファレンス
パラメーター | デフォルト | 説明 |
|
| Elasticsearch クラスターにアクセスするためのエンドポイント。最高のパフォーマンスを得るには、内部エンドポイントを使用してください。 |
|
| Elasticsearch クラスターにアクセスするためのポート。 |
|
| Elasticsearch 認証用のユーザー名。 |
|
| 指定されたユーザー名のパスワード。忘れた場合はリセットしてください。詳細については、「Elasticsearch クラスターのアクセスパスワードのリセット」をご参照ください。 |
|
|
|
|
|
|
|
|
|
|
|
|
|
| 読み書き対象のインデックスとタイプを |
|
| ソースデータと Elasticsearch インデックス間のフィールド名のマッピング。 |
| (なし) | 読み取り時にスクロールリクエストごとにフェッチされるドキュメントの数。 |
ES-Hadoop の設定オプションの完全なリストについては、「ES-Hadoop 設定リファレンス」をご参照ください。
Spark ジョブの送信
コンパイルされたコードを JAR ファイルにパッケージ化し、EMR マスターノードまたは関連するゲートウェイクラスターにアップロードします。
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
結果の検証
Elasticsearch クラスターの Kibana コンソールにログインします。詳細については、「Kibana コンソールへのログイン」をご参照ください。
左側のメニューで、[Dev Tools] をクリックします。
[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 のサポート」をご参照ください。