Todos os produtos
Search
Central de documentação

Elasticsearch:Use ES-Hadoop to ative Apache Spark to write data to and read data from Alibaba Cloud Elasticsearch

Última atualização: Jun 27, 2026

Usar o ES-Hadoop para permitir que o Apache Spark grave e leia dados no Alibaba Cloud Elasticsearch

O Elasticsearch-Hadoop (ES-Hadoop) conecta o Apache Spark ao Alibaba Cloud Elasticsearch e permite que jobs do Spark leiam e gravem dados em um cluster Elasticsearch sem conectores personalizados. Este tutorial apresenta um exemplo completo de ponta a ponta: preparação do ambiente, gravação de registros JSON em um índice do Elasticsearch, leitura desses registros e verificação dos resultados no Kibana.

Pré-requisitos

Antes de começar, verifique se você tem:

  • Um cluster do Alibaba Cloud Elasticsearch (este tutorial usa a versão V6.7.0) com o recurso Auto Indexing ativado

    Importante

    Desative o recurso Auto Indexing em ambientes de produção. Crie o índice e configure os mapeamentos antecipadamente. Este tutorial ativa o Auto Indexing apenas para fins de teste.

  • Um cluster E-MapReduce (EMR) executando na mesma virtual private cloud (VPC) do cluster Elasticsearch, com a seguinte configuração:

    • EMR Version: EMR-3.29.0

    • Required Services: Spark 2.4.5 (mantenha as configurações padrão para os demais serviços)

  • O endereço IP privado do cluster EMR adicionado à lista de permissões de endereços IP privados do cluster Elasticsearch, caso você tenha alterado a lista de permissões padrão (o padrão é 0.0.0.0/0). Para obter o endereço IP privado do cluster EMR, consulte Visualize the cluster list and cluster details. Para atualizar a lista de permissões, consulte Configure a public or private IP address whitelist for an Elasticsearch cluster.

  • JDK 1.8.0 ou superior

Para obter instruções sobre como criar um cluster Elasticsearch, consulte Crie an Alibaba Cloud Elasticsearch cluster e Access and configure an Elasticsearch cluster. Para instruções sobre como criar um cluster EMR, consulte Crie a cluster.

Como funciona

  1. Faça upload dos dados de teste JSON para o Hadoop Distributed File System (HDFS) no nó mestre do EMR.

  2. Crie um projeto Java Maven com as dependências do ES-Hadoop e do Spark e compile as classes de escrita e leitura.

  3. Empacote o código compilado em um JAR e envie-o como um job do Spark usando spark-submit.

  4. O ES-Hadoop serializa cada registro do Resilient Distributed Dataset (RDD) e o grava no índice do Elasticsearch especificado pela API REST.

  5. Verifique os dados gravados executando uma consulta no console Dev Tools do Kibana.

Preparar dados de teste

  1. Faça login no E-MapReduce console, obtenha o endereço IP do nó mestre do EMR e conecte-se via SSH à instância ECS correspondente. Para mais detalhes, consulte Log on to a cluster.

  2. Crie um arquivo chamado http_log.txt com os seguintes registros JSON:

    {"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"}
  3. Faça upload do arquivo para o HDFS:

    hadoop fs -put http_log.txt /tmp/hadoop-es

Adicionar dependências POM

Crie um projeto Java Maven e adicione as seguintes dependências ao arquivo pom.xml. Certifique-se de que cada versão corresponda à versão do serviço Alibaba Cloud associado. Por exemplo, elasticsearch-spark-20_2.11 deve corresponder à versão do seu cluster Elasticsearch, e spark-core_2.12 deve corresponder à versão do seu HDFS.

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

Gravar dados no Elasticsearch

O código Java a seguir lê o arquivo http_log.txt do HDFS e grava cada linha como um documento no índice company/_doc. Substitua os valores de espaço reservado pelo endpoint e pelas credenciais reais do seu cluster.

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");
    }
}
  1. es.nodes — o endpoint interno do seu cluster Elasticsearch. Obtenha-o na página Basic Information do cluster. Consulte Visualize the basic information of a cluster.

  2. es.net.http.auth.pass — a senha definida durante a criação do cluster. Substitua xxxxxx pela sua senha real.

  3. es.nodes.wan.only — defina como true porque o Alibaba Cloud Elasticsearch usa um IP virtual. Isso desativa a descoberta de nós e direciona todo o tráfego pelo endpoint especificado em es.nodes.

  4. es.nodes.discovery — deve ser false para o Alibaba Cloud Elasticsearch. Os endereços internos dos nós do cluster não são acessíveis externamente, portanto o mecanismo nativo de descoberta de nós não funciona.

  5. es.input.use.sliced.partitions — defina como false para ignorar a fase de leitura antecipada do índice e melhorar a eficiência da consulta. A fase de leitura antecipada pode levar mais tempo do que a própria consulta de dados.

  6. Substitua /tmp/hadoop-es/http_log.txt pelo caminho real dos seus dados de teste no HDFS.

Importante

Evite usar a conta elastic em produção. Redefinir a senha dessa conta pode bloquear temporariamente o acesso ao cluster. Em vez disso, crie um usuário dedicado com a função necessária no Kibana. Consulte Use the RBAC mechanism provided by Elasticsearch X-Pack to implement access control.

Ler dados do Elasticsearch

O código a seguir lê todos os documentos do índice company/_doc e os imprime na saída padrão (stdout):

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();
    }
}

Referência de configuração

Parâmetro

Padrão

Descrição

es.nodes

localhost

Endpoint para acessar o cluster Elasticsearch. Use o endpoint interno para obter o melhor desempenho.

es.port

9200

Porta para acessar o cluster Elasticsearch.

es.net.http.auth.user

elastic

Nome de usuário para autenticação no Elasticsearch.

es.net.http.auth.pass

/

Senha do nome de usuário especificado. Caso tenha esquecido, redefina-a. Consulte Reset the access password for an Elasticsearch cluster.

es.nodes.wan.only

false

Quando definido como true, desativa a descoberta de nós e roteia todo o tráfego por es.nodes. Obrigatório quando o Elasticsearch usa um IP virtual.

es.nodes.discovery

true

Quando definido como false, impede que o ES-Hadoop descubra nós adicionais do cluster. Deve ser false no Alibaba Cloud Elasticsearch, pois os endereços internos dos nós não são acessíveis externamente.

es.input.use.sliced.partitions

true

Quando definido como false, ignora a fase de leitura antecipada do índice. Defina como false para melhorar a eficiência da consulta, pois a leitura antecipada pode demorar mais que a consulta real dos dados.

es.index.auto.create

true

Quando definido como true, o ES-Hadoop cria automaticamente o índice caso ele não exista. Desative esta opção em produção e crie índices com mapeamentos explícitos.

es.resource

/

O índice e o tipo de onde ler ou onde gravar, no formato index/type.

es.mapping.names

/

Mapeamentos de nomes de campos entre os dados de source e o índice do Elasticsearch.

es.scroll.size

(nenhum)

Quantidade de documentos buscados por solicitação scroll durante a leitura.

Para a lista completa de opções de configuração do ES-Hadoop, consulte ES-Hadoop configuration reference.

Enviar os jobs do Spark

  1. Empacote o código compilado em um JAR e faça upload para o nó mestre do EMR ou para um cluster de gateway associado.

  2. No cliente EMR, execute os comandos abaixo para enviar os jobs.

    Gravar dados:

    cd /usr/lib/spark-current
    ./bin/spark-submit --master yarn --executor-cores 1 --class "SparkWriteEs" /usr/local/spark_es.jar

    Ler dados:

    cd /usr/lib/spark-current
    ./bin/spark-submit --master yarn --executor-cores 1 --class "ReadES" /usr/local/spark_es.jar
    Importante

    Substitua /usr/local/spark_es.jar pelo caminho real onde você fez upload do JAR.

    Após a conclusão do job de leitura, a saída exibe cada documento recuperado do índice do Elasticsearch:

    Returned result

Verificar resultados

  1. Faça login no console Kibana do seu cluster Elasticsearch. Consulte Log on to the Kibana console.

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

  3. Na aba Console, execute a consulta a seguir para confirme que os dados foram gravados com sucesso:

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

    Uma resposta bem-sucedida retorna os três documentos gravados pelo job do Spark:

    Query result

Próximos passos

O ES-Hadoop oferece suporte a muito mais do que apenas gravações via Java RDD. Após a integração com o Spark, também é possível utilizar Spark Datasets, Spark Streaming, Scala e Spark SQL. Para mais detalhes, consulte Apache Spark support.