Exemplo de Hadoop MapReduce: Primeiro Java Programa com Code

⚡ Resumo Inteligente

Os programas Hadoop MapReduce são escritos como três Java classes, um mapeador, um redutor e um driver, que são compilados, empacotados em um arquivo jar e enviados ao cluster para contabilizar as vendas por país.

  • 🔘 Conjunto de dados: O arquivo SalesJan2009.csv contém uma transação por linha, com o país na oitava coluna, separada por vírgulas.
  • ☑️ Mapeador: O SalesMapper divide cada linha e emite o país juntamente com o valor constante um.
  • Redutor: O SalesCountryReducer soma as unidades recebidas para cada país em uma única contagem de frequência.
  • 🧪 Motorista: O SalesCountryDriver nomeia o trabalho, declara os tipos de chave e valor e conecta ambas as classes.
  • 🛠️ Constituição: Compile com javac, adicione uma entrada de manifesto Main-Class e, em seguida, empacote tudo com jar cfm.
  • ⚠️ Executar: Copie o arquivo CSV para o HDFS, envie o arquivo JAR e leia a parte 00000 do diretório de saída.

Exemplo de Hadoop MapReduce criando um primeiro Java programa

Neste tutorial, você aprenderá a usar o Hadoop com exemplos de MapReduce. Os dados de entrada usados ​​são VendasJan2009.csvContém informações relacionadas às vendas, como nome do produto, preço, forma de pagamento, cidade e país do cliente. O objetivo é descobrir o número de produtos vendidos em cada país.

Primeiro programa Hadoop MapReduce

Agora neste Tutorial MapReduce, criaremos nosso primeiro Java Programa MapReduce:

A captura de tela abaixo mostra os dados brutos de VendasJan2009, onde cada linha representa uma transação e o país está na oitava coluna, separada por vírgulas.

Dados de vendas de janeiro de 2009 abertos em uma planilha mostrando as colunas de transações.

Certifique-se de ter o Hadoop instalado. Antes de iniciar o processo propriamente dito, altere o usuário para 'hduser' (o ID usado durante a configuração do Hadoop — você pode alternar para o ID de usuário usado durante sua própria configuração do Hadoop).

su - hduser_

O prompt muda para a conta hduser, conforme mostrado abaixo.

Prompt do terminal após alternar para a conta hduser

Passo 1) Crie o diretório do projeto e os arquivos de origem.

Crie um novo diretório com o nome MapReduceTutorial, conforme mostrado no exemplo de MapReduce abaixo.

sudo mkdir MapReduceTutorial

O comando mkdir cria o diretório MapReduceTutorial.

Dê permissões

sudo chmod -R 777 MapReduceTutorial

Comando chmod concedendo permissões totais no MapReduceTutorial

Crie os três Java Os arquivos de origem estão abaixo, dentro do MapReduceTutorial. Observe que todos os três usam a versão mais antiga. org.apache.hadoop.mapred API, que ainda é distribuída com o Hadoop 3.x.

SalesMapper.java

package SalesCountry;

import java.io.IOException;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.LongWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapred.*;

public class SalesMapper extends MapReduceBase implements Mapper <LongWritable, Text, Text, IntWritable> {
	private final static IntWritable one = new IntWritable(1);

	public void map(LongWritable key, Text value, OutputCollector <Text, IntWritable> output, Reporter reporter) throws IOException {

		String valueString = value.toString();
		String[] SingleCountryData = valueString.split(",");
		output.collect(new Text(SingleCountryData[7]), one);
	}
}

SalesCountryReducer.java

package SalesCountry;

import java.io.IOException;
import java.util.*;

import org.apache.hadoop.io.IntWritable;
import org.apache.hadoop.io.Text;
import org.apache.hadoop.mapred.*;

public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {

	public void reduce(Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException {
		Text key = t_key;
		int frequencyForCountry = 0;
		while (values.hasNext()) {
			// replace type of value with the actual type of our value
			IntWritable value = (IntWritable) values.next();
			frequencyForCountry += value.get();
			
		}
		output.collect(key, new IntWritable(frequencyForCountry));
	}
}

SalesCountryDriver.java

package SalesCountry;

import org.apache.hadoop.fs.Path;
import org.apache.hadoop.io.*;
import org.apache.hadoop.mapred.*;

public class SalesCountryDriver {
    public static void main(String[] args) {
        JobClient my_client = new JobClient();
        // Create a configuration object for the job
        JobConf job_conf = new JobConf(SalesCountryDriver.class);

        // Set a name of the Job
        job_conf.setJobName("SalePerCountry");

        // Specify data type of output key and value
        job_conf.setOutputKeyClass(Text.class);
        job_conf.setOutputValueClass(IntWritable.class);

        // Specify names of Mapper and Reducer Class
        job_conf.setMapperClass(SalesCountry.SalesMapper.class);
        job_conf.setReducerClass(SalesCountry.SalesCountryReducer.class);

        // Specify formats of the data type of Input and output
        job_conf.setInputFormat(TextInputFormat.class);
        job_conf.setOutputFormat(TextOutputFormat.class);

        // Set input and output directories using command line arguments, 
        //arg[0] = name of input directory on HDFS, and arg[1] =  name of output directory to be created to store the output file.

        FileInputFormat.setInputPaths(job_conf, new Path(args[0]));
        FileOutputFormat.setOutputPath(job_conf, new Path(args[1]));

        my_client.setConf(job_conf);
        try {
            // Run the job 
            JobClient.runJob(job_conf);
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

Baixe os arquivos aqui

O arquivo se expande nos mesmos três arquivos de origem, conforme mostrado aqui.

ExtracLista de arquivos ted SalesMapper, SalesCountryReducer e SalesCountryDriver

Verifique as permissões de arquivo de todos esses arquivos

Lista extensa mostrando as permissões de arquivo nos três Java Arquivos Fonte

Se as permissões de 'leitura' estiverem faltando, conceda-as:

O comando chmod adiciona permissão de leitura ao Java Arquivos Fonte

Etapa 2) Exporte o classpath do Hadoop

Exporte o classpath conforme mostrado no exemplo do Hadoop abaixo.

export CLASSPATH="$HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-core-2.2.0.jar:$HADOOP_HOME/share/hadoop/mapreduce/hadoop-mapreduce-client-common-2.2.0.jar:$HADOOP_HOME/share/hadoop/common/hadoop-common-2.2.0.jar:~/MapReduceTutorial/SalesCountry/*:$HADOOP_HOME/lib/*"

O classpath exportado é exibido novamente no prompt.

Prompt do shell após exportar a variável CLASSPATH do Hadoop

Etapa 3) Compile o Java arquivos

Compile o Java Os arquivos (estes arquivos estão presentes no diretório Final-MapReduceHandsOn) e seus respectivos arquivos de classe serão colocados no diretório do pacote.

javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java

Saída do javac mostrando um aviso de depreciação durante a compilação

Este aviso pode ser ignorado com segurança — ele apenas informa que a API mapred está obsoleta.

Esta compilação criará um diretório no diretório atual com o nome do pacote especificado em Java arquivo de origem (ou seja, SalesCountry em nosso caso) e coloque todos os arquivos de classe compilados nele.

O diretório do pacote SalesCountry contém os três arquivos de classe compilados.

Passo 4) Crie o arquivo de manifesto

Crie um novo arquivo chamado Manifest.txt.

sudo gedit Manifest.txt

Adicione a seguinte linha:

Main-Class: SalesCountry.SalesCountryDriver

gedit Janela mostrando a entrada Main-Class no arquivo Manifest.txt

SalesCountry.SalesCountryDriver é o nome da classe principal. Observe que você precisa pressionar a tecla Enter ao final desta linha.

Passo 5) Coloque as aulas em um frasco.

Crie um arquivo Jar

jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class

O comando jar cria o arquivo ProductSalePerCountry.jar a partir do manifesto.

Verifique se o arquivo jar foi criado

Listagem no diretório confirmando a criação do arquivo ProductSalePerCountry.jar.

Etapa 6) Iniciar o Hadoop

Inicie o Hadoop

$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh

Etapa 7) Copie o arquivo de entrada para o HDFS.

Copie o arquivo SalesJan2009.csv para ~/inputMapReduce.

Agora, use o comando abaixo para copiar ~/inputMapReduce para o HDFS.

$HADOOP_HOME/bin/hdfs dfs -copyFromLocal ~/inputMapReduce /

Saída do comando copyFromLocal com aviso de biblioteca nativa

Podemos ignorar este aviso com segurança.

Verifique se um arquivo foi realmente copiado ou não.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls listando SalesJan2009.csv dentro do diretório inputMapReduce

Etapa 8) Execute o trabalho MapReduce

Execute o trabalho MapReduce

$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales

Saída do console enquanto a tarefa SalePerCountry é executada no cluster.

Isto criará um diretório de saída chamado mapreduce_output_sales em HDFS. O conteúdo deste diretório será um arquivo contendo as vendas de produtos por país.

Passo 9) Leia os resultados

O resultado pode ser visualizado através da interface de comando como,

$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000

Pares de país e contagem de vendas impressos pelo comando hdfs dfs -cat

Os resultados também podem ser vistos através de uma interface web como-

Abra http://localhost:50070/ em um navegador web. No Hadoop 3.x, a interface web do NameNode foi movida para a porta. 9870então use http://localhost:9870/ lá em vez disso.

Página inicial da interface web do NameNode em um navegador

Agora selecione 'Procurar no sistema de arquivos' e navegue até /mapreduce_output_sales

Navegador de arquivos HDFS exibindo o diretório mapreduce_output_sales.

Abrir parte-r-00000

O arquivo de resultado part-r-00000 foi aberto no navegador HDFS.

Explicação da classe SalesMapper

Com o processo de execução completo, as próximas três seções detalham o que cada classe realmente faz.

Nesta seção, vamos entender a implementação da classe SalesMapper.

1. Começamos especificando o nome do pacote para nossa classe. SalesCountry é o nome do nosso pacote. Observe que o resultado da compilação, SalesMapper.class, será colocado em um diretório com o nome deste pacote: SalesCountry.

Em seguida, importamos pacotes de biblioteca.

A imagem abaixo mostra uma implementação da classe SalesMapper.

Visão do editor da implementação completa da classe SalesMapper

Amostra Code Explicação:

1. Definição de classe SalesMapper-

public class SalesMapper extends MapReduceBase implements Mapper<LongWritable, Text, Text, IntWritable> {

Toda classe de mapeamento deve estender a classe MapReduceBase e implementar a interface Mapper.

2. Definindo a função 'mapa'-

public void map(LongWritable key,
         Text value,
OutputCollector<Text, IntWritable> output,
Reporter reporter) throws IOException

A parte principal da classe Mapper é um método 'map()' que aceita quatro argumentos.

A cada chamada ao método 'map()', um par chave-valor ('chave' e 'valor' neste código) é passado.

O método 'map()' começa dividindo o texto de entrada recebido como argumento. Ele divide cada linha em campos.

String valueString = value.toString();
String[] SingleCountryData = valueString.split(",");

Aqui, ',' é usado como delimitador.

Em seguida, um par é formado usando um registro no 7º índice da matriz 'SingleCountryData' e o valor '1'.

output.collect(new Text(SingleCountryData[7]), one);

Estamos selecionando o registro no 7º índice porque precisamos dos dados do país e eles estão localizados no 7º índice da matriz 'SingleCountryData'.

Observe que nossos dados de entrada estão no formato abaixo (onde o país está no 7º índice, com 0 como índice inicial):

Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude

A saída do mapeador é novamente um par chave-valor emitido pelo método 'collect()' do 'OutputCollector'.

Explicação da classe SalesCountryReducer

Nesta seção, vamos entender a implementação da classe SalesCountryReducer.

1. Começamos especificando o nome do pacote para nossa classe. SalesCountry é o nome do nosso pacote. Observe que o resultado da compilação, SalesCountryReducer.class, será salvo em um diretório com o nome deste pacote: SalesCountry.

Em seguida, importamos pacotes de biblioteca.

A imagem abaixo mostra uma implementação da classe SalesCountryReducer.

Visão do editor da implementação completa da classe SalesCountryReducer

Code Explicação:

1. Definição da classe SalesCountryReducer-

public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {

Aqui, os dois primeiros tipos de dados, 'Text' e 'IntWritable', são os tipos de dados do par chave-valor de entrada para o redutor.

A saída do mapeador está na forma de , A saída do mapeador torna-se a entrada do redutor. Portanto, para estar em conformidade com o tipo de dados, Text e IntWritable são usados ​​aqui.

Os dois últimos tipos de dados, 'Text' e 'IntWritable', são os tipos de dados de saída gerados pelo redutor na forma de um par chave-valor.

Toda classe redutora deve estender a classe MapReduceBase e implementar a interface Reducer.

2. Definindo a função 'reduzir'-

public void reduce( Text t_key,
             Iterator<IntWritable> values,                           
             OutputCollector<Text,IntWritable> output,
             Reporter reporter) throws IOException {

A entrada para o método reduce() é uma chave com uma lista de múltiplos valores.

Por exemplo, no nosso caso, será-

, , , , , .

Isso é fornecido ao redutor como

Assim, para aceitar argumentos desse formato, os dois primeiros tipos de dados são utilizados, a saber, Texto e Iterador. O texto é um tipo de dados chave-iterador. É um tipo de dado para uma lista de valores para essa chave.

O próximo argumento é do tipo OutputCollector. que coleta a saída da fase redutora.

O método reduce() começa copiando o valor da chave e inicializando a contagem de frequência com 0.

Text key = t_key;
int frequencyForCountry = 0;

Em seguida, usando um laço 'while', iteramos pela lista de valores associados à chave e calculamos a frequência final somando todos os valores.

 while (values.hasNext()) {
            // replace type of value with the actual type of our value
            IntWritable value = (IntWritable) values.next();
            frequencyForCountry += value.get();
            
        }

Agora, enviamos o resultado para o coletor de saída na forma de chave e contagem de frequência obtida.

O código abaixo faz isso-

output.collect(key, new IntWritable(frequencyForCountry));

Explicação da classe SalesCountryDriver

Nesta seção, vamos entender a implementação da classe SalesCountryDriver.

1. Começamos especificando o nome do pacote para nossa classe. SalesCountry é o nome do nosso pacote. Observe que o resultado da compilação, SalesCountryDriver.class, será colocado em um diretório com o nome deste pacote: SalesCountry.

Aqui está uma linha especificando o nome do pacote seguido do código para importar pacotes de biblioteca.

Declaração de embalagem e declarações de importação da SalesCountryDriver

2. Defina uma classe de driver que criará um novo trabalho de cliente, objeto de configuração e anunciará classes Mapeadora e Redutora.

A classe driver é responsável por configurar nosso trabalho MapReduce para ser executado em HadoopNesta classe, especificamos o nome da tarefa, o tipo de dados de entrada/saída e os nomes das classes de mapeamento e redução.

O código do driver define o nome do trabalho, as classes de chave e valor e os formatos.

3. No trecho de código abaixo, definimos diretórios de entrada e saída que são usados ​​para consumir o conjunto de dados de entrada e produzir saída, respectivamente.

arg[0] e arg[1] são os argumentos da linha de comando passados ​​com um comando fornecido no MapReduce hands-on, ou seja,

$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales

Código do driver que passa os caminhos de entrada e saída da linha de comando.

4. Acione nosso trabalho

O código abaixo inicia a execução da tarefa MapReduce.

try {
    // Run the job 
    JobClient.runJob(job_conf);
} catch (Exception e) {
    e.printStackTrace();
}

Perguntas Frequentes

Todas as três classes importam org.apache.hadoop.mapred, a API MapReduce original. O Hadoop 3.x ainda a inclui e a executa, mas novos projetos normalmente são escritos usando org.apache.hadoop.mapreduce, que substitui JobConf e JobClient por Configuration e Job.

Modelos treinados com base no histórico de tarefas anteriores preveem o tempo de execução, sugerem tamanhos de divisão e número de redutores e identificam distorções antes da conclusão da execução. Eles também sinalizam tarefas cujos contadores de registros perdidos ou tarefas com falha se desviam do intervalo usual.

O Copilot executa bem as partes repetitivas: assinaturas de classe, genéricos, importações e chamadas de configuração do driver. A lógica de negócios, como qual coluna contém o país, ainda precisa ser verificada em relação ao esquema real por um desenvolvedor.

O Hadoop se recusa a gravar em um diretório de saída existente, portanto, os resultados finais nunca são sobrescritos. Exclua o diretório `mapreduce_output_sales` recursivamente com `hdfs dfs -rm -r`, ou passe um caminho de saída diferente na próxima execução.

A ferramenta jar ignora uma linha final do manifesto que não possui um terminador de linha, então a classe Main-Class é descartada silenciosamente. O arquivo jar é então compilado sem erros, mas falha em tempo de execução porque nenhuma classe principal é registrada.

Sim. O exemplo fixa a versão 2.2.0 em todos os nomes de arquivos JAR. Em outra versão, substitua por essa versão ou simplesmente use a saída do comando `hadoop classpath`, que lista todos os arquivos JAR necessários para a distribuição instalada.

Não. Uma instalação pseudo-distribuída de nó único executa o arquivo sem alterações, pois os comandos pressupõem apenas que o HDFS e o YARN estejam em execução. O mesmo arquivo JAR é submetido a um cluster real com vários nós sem qualquer alteração de código.

Altere o índice da matriz no mapeador de 7 para 5, pois "City" é a sexta coluna do arquivo SalesJan2009.csv. Recompile, reconstrua o arquivo JAR e execute o trabalho em um novo diretório de saída.

Resuma esta postagem com: