All Products
Search
Document Center

Elasticsearch:Gunakan ES-Hadoop untuk menulis data dari HDFS ke Elasticsearch

Last Updated:Sep 02, 2026

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

  1. Persiapan

    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.

  2. Langkah 1: Unggah paket JAR ES-Hadoop ke HDFS

    Unduh paket instalasi ES-Hadoop dan unggah ke direktori HDFS pada node master kluster EMR.

  3. Langkah 2: Konfigurasikan dependensi pom

    Buat proyek Java Maven dan konfigurasikan dependensinya dalam file pom.xml.

  4. 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.

  5. Langkah 4: Verifikasi hasil

    Di Konsol Kibana instance Elasticsearch Anda, lihat data yang ditulis oleh task MapReduce.

Persiapan

  1. 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.

    Penting

    Dalam 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.

  2. 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.

    Penting

    Daftar 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:

  3. 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"}
  4. Siapkan lingkungan Java. Versi JDK harus 1.8.0 atau lebih baru.

Langkah 1: Unggah paket JAR ES-Hadoop ke HDFS

  1. Unduh paket instalasi ES-Hadoop yang versinya sesuai dengan kluster Elasticsearch Anda.

    Topik ini menggunakan elasticsearch-hadoop-6.7.0.zip.

  2. 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.

  3. Unggah paket elasticsearch-hadoop-6.7.0.zip ke node master dan ekstrak paket tersebut untuk mendapatkan file elasticsearch-hadoop-6.7.0.jar.

  4. 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>
Penting

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

  1. 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.

    Catatan

    Jika Anda menentukan akun elastic dalam aplikasi Anda, perubahan kata sandi akun ini di masa depan dapat menyebabkan gangguan layanan sementara akibat waktu tunda propagasi. Oleh karena itu, penggunaan akun elastic tidak 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 true dapat meningkatkan signifikan waktu pra-baca indeks, kadang-kadang jauh lebih lama daripada waktu kueri itu sendiri. Kami menyarankan Anda mengatur parameter ini ke false untuk 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.

  2. Kemas kode ke dalam file JAR dan unggah ke mesin client EMR, seperti gateway atau node master kluster EMR.

  3. Di mesin client EMR, jalankan perintah berikut untuk mengeksekusi program MapReduce.

    hadoop jar es-mapreduce-1.0-SNAPSHOT.jar /tmp/hadoop-es/map.json
    Catatan

    Ganti es-mapreduce-1.0-SNAPSHOT.jar dengan nama file JAR yang Anda unggah.

Langkah 4: Verifikasi hasil

  1. Login ke Konsol Kibana instance Alibaba Cloud Elasticsearch Anda.

    Untuk informasi lebih lanjut, lihat Login ke Konsol Kibana.

  2. Di panel navigasi kiri, klik Dev Tools.

  3. 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.