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
-
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.
-
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.
-
Etapa 2: Configurar dependências POM
Crie um projeto Java Maven e configure as dependências POM.
-
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.
-
Etapa 4: Verificar os resultados
Faça login no console Kibana do cluster Elasticsearch e consulte os dados gravados pela tarefa MapReduce.
Preparações
-
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
. Neste tópico, utiliza-se um cluster Elasticsearch V6.7.0.
ImportanteEm 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.
-
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
.
ImportanteA 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:
Consulte Visualizar lista e detalhes do cluster para obter o endereço IP interno do cluster EMR.
Consulte Configurar uma lista de permissões de endereços IP públicos ou privados para um cluster Elasticsearch para definir a lista de permissões de endereços IP privados da VPC do cluster Elasticsearch.
-
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"} 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
-
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.
-
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.
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.
-
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>
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
-
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âmetro Valor padrão Descrição es.nodes localhost Endpoint 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.port 9200 Número da porta utilizada para acessar o cluster Elasticsearch. es.net.http.auth.user elastic Nome de usuário utilizado para acessar o cluster Elasticsearch. NotaSe você especificar a conta
elasticna 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 contaelastic. 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.only false Define 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.discovery true Indica 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.
Notaes.input.use.sliced.partitions true Especifica se partições devem ser utilizadas. Valores válidos: -
true: Usa partições slice. Definir este parâmetro como
truepode 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 comofalsepara melhorar o desempenho da consulta. -
false: Não usa partições slice.
es.index.auto.create true Determina 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.json false Especifica 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.metadata false Define 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.
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.
-
No cliente EMR, execute o seguinte comando para executar a tarefa MapReduce:
hadoop jar es-mapreduce-1.0-SNAPSHOT.jar /tmp/hadoop-es/map.jsonNotaSubstitua es-mapreduce-1.0-SNAPSHOT.jar pelo nome do arquivo JAR enviado.
Etapa 4: Verificar os resultados
-
Faça login no console Kibana do cluster Elasticsearch.
No painel de navegação à esquerda, clique em Dev Tools.
-
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.

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.