ES-Hadoop adalah alat dari Elasticsearch yang terintegrasi dengan ekosistem Hadoop dan memungkinkan transfer data antara Elasticsearch dan Hadoop. Fitur ini memungkinkan Anda menggabungkan kemampuan pencarian cepat Elasticsearch dengan kekuatan pemrosesan batch Hadoop untuk analisis data interaktif. Untuk analisis kompleks, Anda dapat menggunakan task MapReduce guna membaca file JSON dari Hadoop Distributed File System (HDFS) dan menuliskannya ke kluster Elasticsearch. Topik ini menjelaskan cara menggunakan ES-Hadoop bersama task MapReduce untuk menulis data ke kluster Elasticsearch.
Prosedur
-
Buat instance Alibaba Cloud Elasticsearch dan instance E-MapReduce (EMR) dalam virtual private cloud (VPC) yang sama. Aktifkan fitur pembuatan indeks otomatis pada instance Elasticsearch, lalu siapkan data uji dan lingkungan Java.
-
Langkah 1: Unggah paket JAR ES-Hadoop ke HDFS
Unduh paket instalasi ES-Hadoop dan unggah ke direktori HDFS pada node master kluster EMR.
-
Langkah 2: Konfigurasikan dependensi pom
Buat proyek Java Maven dan konfigurasikan dependensinya dalam file pom.xml.
-
Langkah 3: Tulis dan jalankan task MapReduce
Tulis kode Java untuk task MapReduce guna menulis data ke Elasticsearch. Kemas kode tersebut ke dalam file JAR, unggah ke kluster EMR, lalu jalankan untuk menyelesaikan task penulisan data.
-
Di Konsol Kibana instance Elasticsearch Anda, lihat data yang ditulis oleh task MapReduce.
Persiapan
-
Buat instance Alibaba Cloud Elasticsearch dan aktifkan fitur pembuatan indeks otomatis.
Untuk informasi lebih lanjut, lihat Buat instance Alibaba Cloud Elasticsearch dan Konfigurasikan parameter YML. Topik ini menggunakan contoh instance Elasticsearch 6.7.0.
PentingDalam lingkungan produksi, Anda harus menonaktifkan fitur pembuatan indeks otomatis. Buat indeks beserta pemetaannya terlebih dahulu. Karena topik ini hanya untuk tujuan pengujian, fitur pembuatan indeks otomatis diaktifkan.
-
Buat instance EMR dalam VPC yang sama dengan instance Elasticsearch.
Gunakan konfigurasi instance berikut:
-
Edition: EMR-3.29.0
-
Layanan yang Diperlukan: HDFS 2.8.5. Pertahankan pengaturan default untuk layanan lainnya.
Untuk informasi lebih lanjut, lihat Buat kluster.
PentingDaftar putih alamat IP privat untuk kluster Elasticsearch secara default diatur ke 0.0.0.0/0. Anda dapat melihat pengaturan ini di halaman konfigurasi keamanan. Jika Anda mengubah nilai default ini, Anda harus menambahkan alamat IP internal kluster EMR ke daftar putih:
-
Lihat Lihat daftar dan detail kluster untuk mendapatkan alamat IP internal kluster EMR.
-
Lihat Konfigurasikan daftar putih alamat IP publik atau privat untuk kluster Elasticsearch untuk mengatur daftar putih alamat IP privat VPC pada kluster Elasticsearch.
-
-
Siapkan data uji JSON, tulis ke file map.json, lalu unggah file tersebut ke direktori /tmp/hadoop-es di HDFS.
Topik ini menggunakan data uji berikut.
{"id": 1, "name": "zhangsan", "birth": "1990-01-01", "addr": "No.969, wenyixi Rd, yuhang, hangzhou"} {"id": 2, "name": "lisi", "birth": "1991-01-01", "addr": "No.556, xixi Rd, xihu, hangzhou"} {"id": 3, "name": "wangwu", "birth": "1992-01-01", "addr": "No.699 wangshang Rd, binjiang, hangzhou"} -
Siapkan lingkungan Java. Versi JDK harus 1.8.0 atau lebih baru.
Langkah 1: Unggah paket JAR ES-Hadoop ke HDFS
-
Unduh paket instalasi ES-Hadoop yang versinya sesuai dengan kluster Elasticsearch Anda.
Topik ini menggunakan elasticsearch-hadoop-6.7.0.zip.
-
Login ke Konsol EMR, dapatkan alamat IP node master, lalu gunakan SSH untuk login ke instance ECS yang sesuai.
Untuk informasi lebih lanjut, lihat Login ke kluster.
-
Unggah paket elasticsearch-hadoop-6.7.0.zip ke node master dan ekstrak paket tersebut untuk mendapatkan file elasticsearch-hadoop-6.7.0.jar.
-
Buat direktori HDFS dan unggah file elasticsearch-hadoop-6.7.0.jar ke direktori tersebut.
hadoop fs -mkdir /tmp/hadoop-es hadoop fs -put elasticsearch-hadoop-6.7.0/dist/elasticsearch-hadoop-6.7.0.jar /tmp/hadoop-es
Langkah 2: Konfigurasikan dependensi pom
Buat proyek Java Maven dan tambahkan dependensi berikut ke file pom.xml proyek tersebut.
<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>
Pastikan versi dalam dependensi pom sesuai dengan versi layanan yang bersangkutan. Misalnya, versi elasticsearch-hadoop-mr harus sesuai dengan versi Alibaba Cloud Elasticsearch, dan versi hadoop-hdfs harus sesuai dengan versi HDFS.
Langkah 3: Tulis dan jalankan task MapReduce
-
Tulis kode contoh.
Kode berikut membaca file JSON dari direktori /tmp/hadoop-es di HDFS. Setiap baris dari file tersebut kemudian ditulis sebagai dokumen ke Elasticsearch. Kelas EsOutputFormat menangani operasi penulisan pada fase 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); } }Tabel 1. Deskripsi parameter ES-Hadoop
Parameter
Nilai default
Deskripsi
es.nodes
localhost
Titik akhir instance Alibaba Cloud Elasticsearch. Gunakan titik akhir pribadi. Anda dapat menemukannya di halaman Informasi Dasar instance. Untuk informasi lebih lanjut, lihat Lihat informasi dasar instance.
es.port
9200
Nomor port instance Elasticsearch.
es.net.http.auth.user
elastic
Username untuk mengakses instance Elasticsearch.
CatatanJika Anda menentukan akun
elasticdalam aplikasi Anda, perubahan kata sandi akun ini di masa depan dapat menyebabkan gangguan layanan sementara akibat waktu tunda propagasi. Oleh karena itu, penggunaan akunelastictidak disarankan. Sebagai gantinya, buat pengguna khusus dengan izin yang sesuai di Konsol Kibana. Untuk informasi lebih lanjut, lihat Gunakan mekanisme RBAC Elasticsearch X-Pack untuk mengontrol akses pengguna.es.net.http.auth.pass
/
Kata sandi untuk mengakses instance Elasticsearch.
es.nodes.wan.only
false
Menentukan apakah akan melakukan node sniffing saat Anda terhubung ke kluster Elasticsearch yang menggunakan alamat IP virtual di cloud.
-
true: Diaktifkan
-
false: Parameter tidak diatur.
es.nodes.discovery
true
Menentukan apakah akan menonaktifkan penemuan node.
-
true: Dinonaktifkan
-
false: Tidak dinonaktifkan
es.input.use.sliced.partitions
true
Menentukan apakah akan menggunakan partisi slice.
-
true: Menggunakan partisi slice. Mengatur parameter ini ke
truedapat meningkatkan signifikan waktu pra-baca indeks, kadang-kadang jauh lebih lama daripada waktu kueri itu sendiri. Kami menyarankan Anda mengatur parameter ini kefalseuntuk meningkatkan kinerja kueri. -
false: Tidak menggunakan partisi slice.
es.index.auto.create
true
Menentukan apakah akan membuat indeks secara otomatis jika belum ada saat Anda menulis data dari komponen Hadoop ke kluster Elasticsearch.
-
true: Membuat indeks secara otomatis.
-
false: Pembuatan otomatis dinonaktifkan.
es.resource
/
Indeks dan tipe yang akan dibaca atau ditulis.
es.input.json
false
Menentukan apakah input berformat JSON.
-
true: Input berformat JSON.
-
false: Input tidak berformat JSON.
es.mapping.names
/
Pemetaan antara bidang tabel dan bidang indeks Elasticsearch.
es.read.metadata
false
Aktifkan properti ini jika operasi Anda melibatkan bidang internal Elasticsearch, seperti _id.
Untuk informasi lebih lanjut tentang item konfigurasi ES-Hadoop, lihat dokumentasi konfigurasi resmi.
-
-
Kemas kode ke dalam file JAR dan unggah ke mesin client EMR, seperti gateway atau node master kluster EMR.
-
Di mesin client EMR, jalankan perintah berikut untuk mengeksekusi program MapReduce.
hadoop jar es-mapreduce-1.0-SNAPSHOT.jar /tmp/hadoop-es/map.jsonCatatanGanti es-mapreduce-1.0-SNAPSHOT.jar dengan nama file JAR yang Anda unggah.
Langkah 4: Verifikasi hasil
-
Login ke Konsol Kibana instance Alibaba Cloud Elasticsearch Anda.
Untuk informasi lebih lanjut, lihat Login ke Konsol Kibana.
-
Di panel navigasi kiri, klik Dev Tools.
-
Di tab Console, jalankan perintah berikut untuk melihat data yang ditulis oleh task MapReduce.
GET maptest/_search { "query": { "match_all": {} } }Kueri yang berhasil mengembalikan hasil berikut.
{ "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" : "wangwu", "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" : "lisi", "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" : "zhangsan", "birth" : "1990-01-01", "addr" : "No.969, wenyixi Rd, yuhang, hangzhou" } } ] } }
Ringkasan
Topik ini menunjukkan cara menggunakan ES-Hadoop untuk menulis data ke Elasticsearch dengan task MapReduce, menggunakan contoh Alibaba Cloud Elasticsearch dan EMR. Anda juga dapat menggunakan task MapReduce untuk mengkueri data dari Elasticsearch. Konfigurasi untuk mengkueri mirip dengan konfigurasi untuk menulis. Untuk informasi lebih lanjut, lihat dokumentasi resmi Membaca data dari Elasticsearch.