Todos os produtos
Search
Central de documentação

Elasticsearch:Use o ES-Hadoop para gravar dados do HDFS no Elasticsearch

Última atualização: Jun 27, 2026

O ES-Hadoop é uma ferramenta desenvolvida pelo projeto open source Elasticsearch. Ela conecta o Elasticsearch ao Apache Hadoop e permite a transmissão de dados entre ambos. O ES-Hadoop combina a capacidade de busca rápida do Elasticsearch com o processamento em lote do Hadoop para viabilizar o processamento interativo de dados. Em algumas tarefas complexas de análise de dados, é necessário executar uma tarefa MapReduce para ler dados de arquivos JSON armazenados no Hadoop Distributed File System (HDFS) e gravá-los em um cluster Elasticsearch. Este tópico descreve como usar o ES-Hadoop para executar essa tarefa MapReduce.

Procedimento

  1. Preparações

    Crie um cluster Elasticsearch da Alibaba Cloud e um cluster E-MapReduce (EMR) na mesma virtual private cloud (VPC). Em seguida, ative o recurso Auto Indexing no cluster Elasticsearch e prepare os dados de teste e o ambiente Java.

  2. Etapa 1: Fazer upload do pacote JAR do ES-Hadoop para o HDFS

    Baixe o pacote do ES-Hadoop e faça o upload dele para o diretório HDFS no nó mestre do cluster EMR.

  3. Etapa 2: Configurar dependências POM

    Crie um projeto Java Maven e configure as dependências POM.

  4. Etapa 3: Compilar o código e executar uma tarefa MapReduce

    Compile o código Java usado para gravar dados no cluster Elasticsearch. Compacte o código em um pacote JAR e faça o upload do pacote para o cluster EMR. Depois, execute o código em uma tarefa MapReduce para gravar os dados.

  5. Etapa 4: Verificar os resultados

    Faça login no console Kibana do cluster Elasticsearch e consulte os dados gravados pela tarefa MapReduce.

Preparações

  1. Crie um cluster Elasticsearch da Alibaba Cloud e ative o recurso Auto Indexing no cluster.

    Para obter mais informações, consulte

    Criar um cluster Elasticsearch da Alibaba Cloud

    e

    Acesso rápido e configuração

    . Neste tópico, utiliza-se um cluster Elasticsearch V6.7.0.

    Importante

    Em ambientes de produção, recomenda-se desativar o recurso Auto Indexing. Crie um índice e configure mapeamentos para ele antecipadamente. O cluster Elasticsearch utilizado neste tópico destina-se apenas a testes; por isso, o recurso Auto Indexing está ativado.

  2. Crie um cluster EMR na mesma VPC do cluster Elasticsearch.

    Configuração do cluster EMR:

    • Versão do EMR: Selecione EMR-3.29.0.

    • Serviços obrigatórios: O HDFS (2.8.5) é um dos serviços necessários. Mantenha as configurações padrão para os demais serviços.

    Para obter mais informações, consulte

    Criar um cluster

    .

    Importante

    A lista de permissões de endereços IP privados de um cluster Elasticsearch é definida como 0.0.0.0/0 por padrão. Visualize essa configuração na página de configuração de segurança. Caso altere esse padrão, adicione o endereço IP interno do cluster EMR à lista de permissões:

  3. Prepare dados de teste no formato JSON e grave-os no arquivo map.json. Faça o upload do arquivo para o diretório /tmp/hadoop-es do HDFS.

    Este tópico utiliza os seguintes dados de teste:

    {"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. Prepare um ambiente Java. A versão do JDK deve ser 1.8.0 ou posterior.

Etapa 1: Fazer upload do pacote JAR do ES-Hadoop para o HDFS

  1. Baixe o pacote de instalação do ES-Hadoop cuja versão corresponda à do seu cluster Elasticsearch.

    Este tópico utiliza o elasticsearch-hadoop-6.7.0.zip.

  2. Faça login no console EMR, obtenha o endereço IP do nó mestre e use SSH para fazer login na instância ECS correspondente.

    Para obter mais informações, consulte Fazer login em um cluster.

  3. Faça o upload do pacote elasticsearch-hadoop-6.7.0.zip para o nó mestre do cluster EMR. Descompacte o pacote para obter o arquivo elasticsearch-hadoop-6.7.0.jar.

  4. Crie um diretório HDFS e faça o upload do arquivo elasticsearch-hadoop-6.7.0.jar para esse diretório.

    hadoop fs -mkdir /tmp/hadoop-es
    hadoop fs -put elasticsearch-hadoop-6.7.0/dist/elasticsearch-hadoop-6.7.0.jar /tmp/hadoop-es

Etapa 2: Configurar dependências POM

Crie um projeto Java Maven e adicione as seguintes dependências POM ao arquivo pom.xml do projeto.

<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>
Importante

Certifique-se de que as versões das dependências POM sejam consistentes com as dos serviços relacionados da Alibaba Cloud. Por exemplo, a versão do elasticsearch-hadoop-mr deve corresponder à do Alibaba Cloud Elasticsearch, e a versão do hadoop-hdfs deve corresponder à do HDFS.

Etapa 3: Compilar o código e executar uma tarefa MapReduce

  1. Compile o código.

    O código abaixo lê dados dos arquivos JSON no diretório

    /tmp/hadoop-es

    do HDFS. Ele também grava cada linha de dados desses arquivos JSON como um documento no cluster Elasticsearch. O EsOutputFormat conclui a gravação de dados na 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);
        }
    }
    Tabela 1. Parâmetros do ES-Hadoop
    ParâmetroValor padrãoDescrição
    es.nodeslocalhostEndpoint usado para acessar o cluster Elasticsearch. Recomenda-se utilizar o endpoint interno, disponível na página Basic Information do cluster Elasticsearch. Para mais detalhes, consulte Informações básicas do cluster.
    es.port9200Número da porta utilizada para acessar o cluster Elasticsearch.
    es.net.http.auth.userelasticNome de usuário utilizado para acessar o cluster Elasticsearch.
    Nota

    Se você especificar a conta elastic na sua aplicação, qualquer alteração subsequente de senha dessa conta poderá causar interrupções temporárias no serviço devido ao atraso de propagação. Portanto, não é recomendado usar a conta elastic. Em vez disso, crie um usuário dedicado com as permissões adequadas no console Kibana. Para obter mais informações, consulte Usar o mecanismo RBAC do Elasticsearch X-Pack para controlar o acesso de usuários.

    es.net.http.auth.pass/Senha utilizada para acessar o cluster Elasticsearch.
    es.nodes.wan.onlyfalseDefine se o sniffing de nós deve ser ativado quando o cluster Elasticsearch utiliza um endereço IP virtual para conexões. Valores válidos:
    • true: indica que o sniffing de nós está ativado.
    • false: indica que o sniffing de nós está desativado.
    es.nodes.discoverytrueIndica se o mecanismo de descoberta de nós deve ser proibido. Valores válidos:
    • true: indica que o mecanismo de descoberta de nós está proibido.
    • false: indica que o mecanismo de descoberta de nós não está proibido.
    Nota
    es.input.use.sliced.partitionstrueEspecifica se partições devem ser utilizadas. Valores válidos:
    • true: Usa partições slice. Definir este parâmetro como true pode aumentar significativamente o tempo de pré-leitura do índice, tornando-o às vezes muito maior que o próprio tempo de consulta. Recomenda-se definir este parâmetro como false para melhorar o desempenho da consulta.

    • false: Não usa partições slice.

    es.index.auto.createtrueDetermina se o sistema cria um índice no cluster Elasticsearch ao usar o ES-Hadoop para gravar dados no cluster. Valores válidos:
    • true: indica que o sistema cria um índice no cluster Elasticsearch.
    • false: indica que o sistema não cria um índice no cluster Elasticsearch.
    es.resource/Nome e tipo do índice onde as operações de leitura ou gravação de dados ocorrem.
    es.input.jsonfalseEspecifica se os dados de entrada estão no formato JSON.
    es.mapping.names/Mapeamentos entre os nomes dos campos na tabela e aqueles no índice do cluster Elasticsearch.
    es.read.metadatafalseDefine se os metadados do documento, como _id, devem ser incluídos nos resultados. Para incluir os metadados do documento, defina o valor como true.

    Para obter mais informações sobre os itens de configuração do ES-Hadoop, consulte configuração open source do ES-Hadoop.

  2. Compacte o código em um pacote JAR e faça o upload dele para um cliente EMR, como o nó mestre do cluster EMR ou o cluster de gateway associado a este cluster EMR.

  3. No cliente EMR, execute o seguinte comando para executar a tarefa MapReduce:

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

    Substitua es-mapreduce-1.0-SNAPSHOT.jar pelo nome do arquivo JAR enviado.

Etapa 4: Verificar os resultados

  1. Faça login no console Kibana do cluster Elasticsearch.

    Para obter mais informações, consulte

    Fazer login no console Kibana

    .

  2. No painel de navegação à esquerda, clique em Dev Tools.

  3. Na aba Console da página exibida, execute o seguinte comando para consultar os dados gravados pela tarefa MapReduce:

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

    Se o comando for executado com êxito, o resultado mostrado na figura a seguir será retornado.

    Returned result

Resumo

Este tópico descreveu como usar o ES-Hadoop para gravar dados no Elasticsearch executando uma tarefa MapReduce em um cluster EMR. Também é possível executar uma tarefa MapReduce para ler dados do Elasticsearch. As configurações para operações de leitura de dados são semelhantes às de gravação. Para obter mais informações, consulte Leitura de dados do Elasticsearch na documentação open source do Elasticsearch.