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.
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.
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.
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
Dê permissões
sudo chmod -R 777 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(); } } }
O arquivo se expande nos mesmos três arquivos de origem, conforme mostrado aqui.
Verifique as permissões de arquivo de todos esses arquivos
Se as permissões de 'leitura' estiverem faltando, conceda-as:
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.
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
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.
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
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
Verifique se o arquivo jar foi criado
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 /
Podemos ignorar este aviso com segurança.
Verifique se um arquivo foi realmente copiado ou não.
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
Etapa 8) Execute o trabalho MapReduce
Execute o trabalho MapReduce
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
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
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.
Agora selecione 'Procurar no sistema de arquivos' e navegue até /mapreduce_output_sales
Abrir parte-r-00000
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.
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.
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.
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.
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
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(); }





















