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
ImportanteDesative 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
Faça upload dos dados de teste JSON para o Hadoop Distributed File System (HDFS) no nó mestre do EMR.
Crie um projeto Java Maven com as dependências do ES-Hadoop e do Spark e compile as classes de escrita e leitura.
Empacote o código compilado em um JAR e envie-o como um job do Spark usando
spark-submit.O ES-Hadoop serializa cada registro do Resilient Distributed Dataset (RDD) e o grava no índice do Elasticsearch especificado pela API REST.
Verifique os dados gravados executando uma consulta no console Dev Tools do Kibana.
Preparar dados de teste
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.
-
Crie um arquivo chamado
http_log.txtcom 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"} -
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");
}
}
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.es.net.http.auth.pass— a senha definida durante a criação do cluster. Substituaxxxxxxpela sua senha real.es.nodes.wan.only— defina comotrueporque o Alibaba Cloud Elasticsearch usa um IP virtual. Isso desativa a descoberta de nós e direciona todo o tráfego pelo endpoint especificado emes.nodes.es.nodes.discovery— deve serfalsepara 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.es.input.use.sliced.partitions— defina comofalsepara 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.Substitua
/tmp/hadoop-es/http_log.txtpelo caminho real dos seus dados de teste no HDFS.
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 |
|
|
|
Endpoint para acessar o cluster Elasticsearch. Use o endpoint interno para obter o melhor desempenho. |
|
|
|
Porta para acessar o cluster Elasticsearch. |
|
|
|
Nome de usuário para autenticação no Elasticsearch. |
|
|
|
Senha do nome de usuário especificado. Caso tenha esquecido, redefina-a. Consulte Reset the access password for an Elasticsearch cluster. |
|
|
|
Quando definido como |
|
|
|
Quando definido como |
|
|
|
Quando definido como |
|
|
|
Quando definido como |
|
|
|
O índice e o tipo de onde ler ou onde gravar, no formato |
|
|
|
Mapeamentos de nomes de campos entre os dados de source e o índice do Elasticsearch. |
|
|
(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
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.
-
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.jarLer dados:
cd /usr/lib/spark-current ./bin/spark-submit --master yarn --executor-cores 1 --class "ReadES" /usr/local/spark_es.jarImportanteSubstitua
/usr/local/spark_es.jarpelo 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:

Verificar resultados
Faça login no console Kibana do seu cluster Elasticsearch. Consulte Log on to the Kibana console.
No painel de navegação à esquerda, clique em Dev Tools.
-
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:

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.