ES-Hadoop は、Elasticsearch が提供する Hadoop エコシステムと連携するためのツールで、Elasticsearch と Hadoop の間でデータを移動します。これにより、Elasticsearch の高速な検索機能と Hadoop のバッチ処理能力を組み合わせて、インタラクティブなデータ処理を実現できます。複雑な分析では、MapReduce タスクを使用して Hadoop 分散ファイルシステム (HDFS) から JSON ファイルを読み取り、Elasticsearch クラスターに書き込むことができます。本トピックでは、ES-Hadoop と MapReduce タスクを使用して Elasticsearch クラスターにデータを書き込む方法について説明します。
操作手順
-
同じ仮想プライベートクラウド (VPC) に Alibaba Cloud Elasticsearch インスタンスと E-MapReduce (EMR) インスタンスを作成します。次に、Elasticsearch インスタンスのインデックス自動作成機能を有効にし、テストデータと Java 環境を準備します。
-
ステップ 1: ES-Hadoop JAR パッケージを HDFS にアップロードする
ES-Hadoop インストールパッケージをダウンロードし、EMR クラスターのマスターノード上の HDFS ディレクトリにアップロードします。
-
Java Maven プロジェクトを作成し、pom 依存関係を設定します。
-
ステップ 3: MapReduce タスクを記述して実行する
Elasticsearch にデータを書き込むための MapReduce タスクの Java コードを記述します。コードを JAR ファイルにパッケージ化し、EMR クラスターにアップロードして実行し、データ書き込みタスクを完了します。
-
Elasticsearch インスタンスの Kibana コンソールで、MapReduce タスクによって書き込まれたデータを確認します。
事前準備
-
Alibaba Cloud Elasticsearch インスタンスを作成し、インデックスの自動作成機能を有効にします。
詳細については、「Alibaba Cloud Elasticsearch インスタンスの作成」および「YML パラメーターの設定」をご参照ください。本トピックでは、Elasticsearch 6.7.0 インスタンスを例とします。
重要本番環境では、インデックスの自動作成機能を無効にする必要があります。事前にインデックスとそのマッピングを作成してください。本トピックはテスト目的のみであるため、インデックスの自動作成機能を有効にしています。
-
Elasticsearch インスタンスと同じ VPC に EMR インスタンスを作成します。
次のインスタンス設定を使用します:
-
エディション:EMR-3.29.0
-
必須サービス:HDFS 2.8.5。その他のサービスはデフォルト設定のままにします。
詳細については、「クラスターの作成」をご参照ください。
重要ES クラスターのプライベート IP アドレスホワイトリストは、デフォルトで 0.0.0.0/0 に設定されています。この設定はセキュリティ設定ページで確認できます。このデフォルト設定を変更する場合は、EMR クラスターの内部 IP アドレスをホワイトリストに追加する必要があります。
-
EMR クラスターの内部 IP アドレスを取得するには、「クラスターリストと詳細の表示」をご参照ください。
-
ES クラスターの VPC プライベート IP アドレスホワイトリストを設定するには、「Elasticsearch クラスターのパブリックまたはプライベート IP アドレスホワイトリストの設定」をご参照ください。
-
-
JSON テストデータを準備し、map.json ファイルに書き込み、HDFS の /tmp/hadoop-es ディレクトリにファイルをアップロードします。
本トピックでは、次のテストデータを使用します:
{"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"} -
Java 環境を準備します。JDK のバージョンは 1.8.0 以降が必要です。
ステップ 1: ES-Hadoop JAR パッケージを HDFS にアップロードする
-
ES クラスターのバージョンに対応するES-Hadoop インストールパッケージをダウンロードします。
本トピックでは、elasticsearch-hadoop-6.7.0.zip を使用します。
-
EMR コンソールにログインし、マスターノードの IP アドレスを取得してから、SSH を使用して対応する ECS インスタンスにログインします。
詳細については、「クラスターへのログイン」をご参照ください。
-
elasticsearch-hadoop-6.7.0.zip パッケージをマスターノードにアップロードし、パッケージを解凍して elasticsearch-hadoop-6.7.0.jar ファイルを取得します。
-
HDFS ディレクトリを作成し、elasticsearch-hadoop-6.7.0.jar ファイルをそのディレクトリにアップロードします。
hadoop fs -mkdir /tmp/hadoop-es hadoop fs -put elasticsearch-hadoop-6.7.0/dist/elasticsearch-hadoop-6.7.0.jar /tmp/hadoop-es
ステップ 2: pom 依存関係を設定する
Java Maven プロジェクトを作成し、プロジェクトの pom.xml ファイルに次の pom 依存関係を追加します。
<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-shade-plugin</artifactId>
<version>2.4.1</version>
<executions>
<execution>
<phase>package</phase>
<goals>
<goal>shade</goal>
</goals>
<configuration>
<transformers>
<transformer
implementation="org.apache.maven.plugins.shade.resource.ManifestResourceTransformer">
<mainClass>WriteToEsWithMR</mainClass>
</transformer>
</transformers>
</configuration>
</execution>
</executions>
</plugin>
</plugins>
</build>
<dependencies>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-hdfs</artifactId>
<version>2.8.5</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-mapreduce-client-jobclient</artifactId>
<version>2.8.5</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-common</artifactId>
<version>2.8.5</version>
</dependency>
<dependency>
<groupId>org.apache.hadoop</groupId>
<artifactId>hadoop-auth</artifactId>
<version>2.8.5</version>
</dependency>
<dependency>
<groupId>org.elasticsearch</groupId>
<artifactId>elasticsearch-hadoop-mr</artifactId>
<version>6.7.0</version>
</dependency>
<dependency>
<groupId>commons-httpclient</groupId>
<artifactId>commons-httpclient</artifactId>
<version>3.1</version>
</dependency>
</dependencies>
pom 依存関係のバージョンが、対応するサービスのバージョンと一致していることを確認してください。たとえば、elasticsearch-hadoop-mr のバージョンは Alibaba Cloud Elasticsearch のバージョンと、hadoop-hdfs のバージョンは HDFS のバージョンと一致している必要があります。
ステップ 3: MapReduce タスクを記述して実行する
-
サンプルコードを記述します。
次のコードは、HDFS の /tmp/hadoop-es ディレクトリから JSON ファイルを読み取ります。次に、これらのファイルの各行をドキュメントとして Elasticsearch に書き込みます。
EsOutputFormatクラスは、Map フェーズで書き込み操作を完了します。import java.io.IOException; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.conf.Configured; import org.apache.hadoop.fs.Path; import org.apache.hadoop.io.NullWritable; import org.apache.hadoop.io.Text; import org.apache.hadoop.mapreduce.Job; import org.apache.hadoop.mapreduce.Mapper; import org.apache.hadoop.mapreduce.lib.input.FileInputFormat; import org.apache.hadoop.mapreduce.lib.input.TextInputFormat; import org.apache.hadoop.util.GenericOptionsParser; import org.elasticsearch.hadoop.mr.EsOutputFormat; import org.apache.hadoop.util.Tool; import org.apache.hadoop.util.ToolRunner; public class WriteToEsWithMR extends Configured implements Tool { public static class EsMapper extends Mapper<Object, Text, NullWritable, Text> { private Text doc = new Text(); @Override protected void map(Object key, Text value, Context context) throws IOException, InterruptedException { if (value.getLength() > 0) { doc.set(value); System.out.println(value); context.write(NullWritable.get(), doc); } } } public int run(String[] args) throws Exception { Configuration conf = new Configuration(); String[] otherArgs = new GenericOptionsParser(conf, args).getRemainingArgs(); conf.setBoolean("mapreduce.map.speculative", false); conf.setBoolean("mapreduce.reduce.speculative", false); conf.set("es.nodes", "es-cn-4591jumei000u****.elasticsearch.aliyuncs.com"); conf.set("es.port","9200"); conf.set("es.net.http.auth.user", "elastic"); conf.set("es.net.http.auth.pass", "xxxxxx"); conf.set("es.nodes.wan.only", "true"); conf.set("es.nodes.discovery","false"); conf.set("es.input.use.sliced.partitions","false"); conf.set("es.resource", "maptest/_doc"); conf.set("es.input.json", "true"); Job job = Job.getInstance(conf); job.setInputFormatClass(TextInputFormat.class); job.setOutputFormatClass(EsOutputFormat.class); job.setMapOutputKeyClass(NullWritable.class); job.setMapOutputValueClass(Text.class); job.setJarByClass(WriteToEsWithMR.class); job.setMapperClass(EsMapper.class); FileInputFormat.setInputPaths(job, new Path(otherArgs[0])); return job.waitForCompletion(true) ? 0 : 1; } public static void main(String[] args) throws Exception { int ret = ToolRunner.run(new WriteToEsWithMR(), args); System.exit(ret); } }表 1. ES-Hadoop パラメーターの説明
パラメーター
デフォルト値
説明
es.nodes
localhost
Alibaba Cloud Elasticsearch インスタンスのエンドポイント。プライベートエンドポイントを使用します。インスタンスの基本情報ページで確認できます。詳細については、「インスタンスの基本情報の表示」をご参照ください。
es.port
9200
Elasticsearch インスタンスのポート番号。
es.net.http.auth.user
elastic
Elasticsearch インスタンスにアクセスするためのユーザー名。
説明アプリケーションで
elasticアカウントを指定すると、このアカウントのパスワード変更後、伝播遅延により一時的なサービス中断が発生する可能性があります。そのため、elasticアカウントの使用は推奨されません。代わりに、Kibana コンソールで適切な権限を持つ専用ユーザーを作成してください。詳細については、「Elasticsearch X-Pack の RBAC メカニズムを使用したユーザーアクセスの制御」をご参照ください。es.net.http.auth.pass
/
Elasticsearch インスタンスにアクセスするためのパスワード。
es.nodes.wan.only
false
ノードスニッフィングを無効にするかどうかを指定します。
-
true: ノードスニッフィングを無効にします。
-
false: ノードスニッフィングを有効にします (デフォルト)。
es.nodes.discovery
true
ノード検出を有効にするかどうかを指定します。
-
true: 有効
-
false: 無効
es.input.use.sliced.partitions
true
スライスパーティションを使用するかどうかを指定します。
-
true: スライスパーティションを使用します。このパラメーターを
trueに設定すると、インデックスの事前読み取り時間が大幅に増加し、クエリ時間自体よりもはるかに長くなることがあります。クエリパフォーマンスを向上させるため、このパラメーターをfalseに設定することを推奨します。 -
false: スライスパーティションを使用しません。
es.index.auto.create
true
Hadoop コンポーネントから Elasticsearch クラスターにデータを書き込む際に、インデックスが存在しない場合に自動的に作成するかどうかを指定します。
-
true: インデックスを自動的に作成します。
-
false: インデックスを自動的に作成しません。
es.resource
/
読み取りまたは書き込みを行うインデックスとタイプ。
es.input.json
false
入力が JSON 形式かどうかを指定します。
-
true: 入力は JSON 形式です。
-
false: 入力は JSON 形式ではありません。
es.mapping.names
/
テーブルフィールドと Elasticsearch インデックスフィールド間のマッピング。
es.read.metadata
false
操作に [_id] などの Elasticsearch 内部フィールドが含まれる場合は、このプロパティを有効にします。
ES-Hadoop の設定項目の詳細については、公式設定ドキュメントをご参照ください。
-
-
コードを JAR ファイルにパッケージ化し、ゲートウェイや EMR クラスターのマスターノードなどの EMR クライアントマシンにアップロードします。
-
EMR クライアントマシンで、次のコマンドを実行して MapReduce プログラムを実行します。
hadoop jar es-mapreduce-1.0-SNAPSHOT.jar /tmp/hadoop-es/map.json説明es-mapreduce-1.0-SNAPSHOT.jarは、アップロードした JAR ファイルの名前に置き換えてください。
ステップ 4: 結果を検証する
-
Alibaba Cloud Elasticsearch インスタンスの Kibana コンソールにログインします。
詳細については、「Kibana コンソールへのログイン」をご参照ください。
-
左側のナビゲーションペインで、[開発者ツール] をクリックします。
-
[コンソール] タブで、次のコマンドを実行して、MapReduce タスクによって書き込まれたデータを表示します。
GET maptest/_search { "query": { "match_all": {} } }クエリが成功すると、次の結果が返されます。
{ "took" : 8, "timed_out" : false, "_shards" : { "total" : 5, "successful" : 5, "skipped" : 0, "failed" : 0 }, "hits" : { "total" : 3, "max_score" : 1.0, "hits" : [ { "_index" : "maptest", "_type" : "_doc", "_id" : "V8D0KnUB0sZ7Mms_YRGu", "_score" : 1.0, "_source" : { "id" : 3, "name" : "佐藤三郎", "birth" : "1992-01-01", "addr" : "No.699 wangshang Rd, binjiang, hangzhou" } }, { "_index" : "maptest", "_type" : "_doc", "_id" : "WMD0KnUB0sZ7Mms_YRGu", "_score" : 1.0, "_source" : { "id" : 2, "name" : "鈴木二郎", "birth" : "1991-01-01", "addr" : "No.556, xixi Rd, xihu, hangzhou" } }, { "_index" : "maptest", "_type" : "_doc", "_id" : "WcD0KnUB0sZ7Mms_YRGu", "_score" : 1.0, "_source" : { "id" : 1, "name" : "田中一郎", "birth" : "1990-01-01", "addr" : "No.969, wenyixi Rd, yuhang, hangzhou" } } ] } }
まとめ
本トピックでは、Alibaba Cloud Elasticsearch と EMR を例に、ES-Hadoop を使用して MapReduce タスクで Elasticsearch にデータを書き込む方法を説明しました。MapReduce タスクを使用して Elasticsearch からデータをクエリすることもできます。クエリの設定は書き込みの設定と同様です。詳細については、Elasticsearch からのデータの読み取りを扱った公式ドキュメントをご参照ください。