Exemplu Hadoop MapReduce: Primul Java Program cu Code
โก Rezumat inteligent
Programele Hadoop MapReduce sunt scrise รฎn trei Java clase, un mapper, un reducer ศi un driver, care sunt compilate, รฎmpachetate รฎntr-un fiศier jar ศi trimise clusterului pentru a numฤra vรขnzฤrile pe ศarฤ.
รn acest tutorial, veศi รฎnvฤศa sฤ utilizaศi Hadoop cu exemple MapReduce. Datele de intrare utilizate sunt VรขnzฤriJan2009.csvConศine informaศii legate de vรขnzฤri, cum ar fi numele produsului, preศul, metoda de platฤ, oraศul ศi ศara clientului. Scopul este de a afla numฤrul de produse vรขndute รฎn fiecare ศarฤ.
Primul program Hadoop MapReduce
Acum รฎn asta Tutorial MapReduce, vom crea primul nostru Java Programul MapReduce:
Captura de ecran de mai jos prezintฤ datele brute SalesJan2009, unde fiecare linie reprezintฤ o tranzacศie, iar ศara se aflฤ รฎn a opta coloanฤ separatฤ prin virgulฤ.
Asiguraศi-vฤ cฤ aveศi Hadoop instalat. รnainte de a รฎncepe procesul propriu-zis, schimbaศi numele utilizatorului รฎn โhduserโ (ID-ul folosit la configurarea Hadoop - puteศi comuta la ID-ul de utilizator folosit รฎn timpul propriei configurฤri Hadoop).
su - hduser_
Solicitarea se schimbฤ รฎn contul hduser, aศa cum se aratฤ mai jos.
Pasul 1) Creaศi directorul proiectului ศi fiศierele sursฤ
Creaศi un director nou cu numele MapReduceTutorial, aศa cum se aratฤ รฎn exemplul MapReduce de mai jos.
sudo mkdir MapReduceTutorial
Acordaศi permisiuni
sudo chmod -R 777 MapReduceTutorial
Creaศi cele trei Java fiศierele sursฤ de mai jos, รฎn MapReduceTutorial. Reศineศi cฤ toate trei folosesc versiunea mai veche org.apache.hadoop.mapred API, care este รฎncฤ inclus รฎn 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(); } } }
Arhiva se extinde รฎn aceleaศi trei fiศiere sursฤ, aศa cum se aratฤ aici.
Verificaศi permisiunile pentru toate aceste fiศiere
Dacฤ lipsesc permisiunile de โcitireโ, atunci acordaศi-le:
Pasul 2) Exportaศi calea de clasฤ Hadoop
Exportaศi calea de clasฤ aศa cum se aratฤ รฎn exemplul Hadoop de mai jos.
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/*"
Calea de clasฤ exportatฤ este readusฤ la prompt.
Pasul 3) Compilaศi Java fiศiere
Compilaศi Java fiศiere (aceste fiศiere sunt prezente รฎn directorul Final-MapReduceHandsOn). Fiศierele lor de clasฤ vor fi plasate รฎn directorul pachetului.
javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java
Acest avertisment poate fi ignorat รฎn siguranศฤ โ raporteazฤ doar cฤ API-ul mapred este depreciat.
Aceastฤ compilare va crea un director รฎn directorul curent numit cu numele pachetului specificat รฎn Java fiศierul sursฤ (de exemplu, SalesCountry รฎn cazul nostru) ศi puneศi toate fiศierele de clasฤ compilate รฎn el.
Pasul 4) Creaศi fiศierul manifest
Creaศi un fiศier nou Manifest.txt
sudo gedit Manifest.txt
Adฤugaศi urmฤtoarea linie la aceasta:
Main-Class: SalesCountry.SalesCountryDriver
SalesCountry.SalesCountryDriver este numele clasei principale. Reศineศi cฤ trebuie sฤ apฤsaศi tasta Enter la sfรขrศitul acestei linii.
Pasul 5) รmpachetaศi clasele รฎntr-un borcan
Creaศi un fiศier Jar
jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class
Verificaศi dacฤ fiศierul jar este creat
Pasul 6) Porniศi Hadoop
Porniศi Hadoop
$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh
Pasul 7) Copiaศi fiศierul de intrare รฎn HDFS
Copiaศi fiศierul SalesJan2009.csv รฎn ~/inputMapReduce
Acum folosiศi comanda de mai jos pentru a copia ~/inputMapReduce รฎn HDFS.
$HADOOP_HOME/bin/hdfs dfs -copyFromLocal ~/inputMapReduce /
Putem ignora รฎn siguranศฤ acest avertisment.
Verificaศi dacฤ un fiศier este copiat sau nu.
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
Pasul 8) Executaศi jobul MapReduce
Rulaศi jobul MapReduce
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
Aceasta va crea un director de ieศire numit mapreduce_output_sales on HDFS. Conศinutul acestui director va fi un fiศier care conศine vรขnzฤrile de produse pe ศarฤ.
Pasul 9) Citiศi rezultatele
Rezultatul poate fi vฤzut prin interfaศa de comandฤ ca fiind,
$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000
Rezultatele pot fi vฤzute ศi printr-o interfaศฤ web ca-
Operatii Deschise http://localhost:50070/ รฎntr-un browser web. Pe Hadoop 3.x, interfaศa web NameNode a fost mutatฤ รฎn port 9870, deci foloseศte http://localhost:9870/ acolo รฎn schimb.
Acum selectaศi โRฤsfoiศi sistemul de fiศiereโ ศi navigaศi la /mapreduce_output_sales
Deschideศi piesa-r-00000
Explicaศia clasei SalesMapper
Avรขnd รฎn vedere cฤ jobul se desfฤศoarฤ de la un capฤt la altul, urmฤtoarele trei secศiuni prezintฤ ce face de fapt fiecare clasฤ.
รn aceastฤ secศiune, vom รฎnศelege implementarea clasei SalesMapper.
1. รncepem prin a specifica un nume de pachet pentru clasa noastrฤ. SalesCountry este numele pachetului nostru. Reศineศi cฤ rezultatul compilฤrii, SalesMapper.class, va intra รฎntr-un director numit cu acest nume de pachet: SalesCountry.
Dupฤ aceasta, importฤm pachete de bibliotecฤ.
Mai jos este prezentatฤ o implementare a clasei SalesMapper -
Eลantion Code Explicaลฃie:
1. Definiศia clasei SalesMapper-
public class SalesMapper extends MapReduceBase implements Mapper<LongWritable, Text, Text, IntWritable> {
Fiecare clasฤ mapper trebuie sฤ fie extinsฤ din clasa MapReduceBase ศi trebuie sฤ implementeze interfaศa Mapper.
2. Definirea funcศiei โhartฤโ-
public void map(LongWritable key, Text value, OutputCollector<Text, IntWritable> output, Reporter reporter) throws IOException
Partea principalฤ a clasei Mapper este metoda โmap()โ care acceptฤ patru argumente.
La fiecare apel al metodei โmap()โ, se transmite o pereche cheie-valoare (โkeyโ ศi โvalueโ รฎn acest cod).
Metoda 'map()' รฎncepe prin divizarea textului de intrare primit ca argument. Aceasta รฎmparte fiecare linie รฎn cรขmpuri.
String valueString = value.toString(); String[] SingleCountryData = valueString.split(",");
Aici, se foloseศte โ,โ ca delimitator.
Dupฤ aceasta, se formeazฤ o pereche folosind o รฎnregistrare la al 7-lea index al matricei โSingleCountryDataโ ศi o valoare โ1โ.
output.collect(new Text(SingleCountryData[7]), one);
Alegem รฎnregistrarea la al 7-lea index deoarece avem nevoie de date despre ศarฤ, iar acestea se aflฤ la al 7-lea index din matricea โSingleCountryDataโ.
Vฤ rugฤm sฤ reศineศi cฤ datele noastre de intrare sunt รฎn formatul urmฤtor (unde ศara se aflฤ la al 7-lea indice, cu 0 ca indice de pornire) -
Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude
O ieศire a funcศiei mapper este din nou o pereche cheie-valoare emisฤ folosind metoda โcollect()โ a funcศiei โOutputCollectorโ.
Explicaศia clasei SalesCountryReducer
รn aceastฤ secศiune, vom รฎnศelege implementarea clasei SalesCountryReducer.
1. รncepem prin a specifica numele pachetului pentru clasa noastrฤ. SalesCountry este numele pachetului nostru. Reศineศi cฤ rezultatul compilฤrii, SalesCountryReducer.class, va intra รฎntr-un director numit cu acest nume de pachet: SalesCountry.
Dupฤ aceasta, importฤm pachete de bibliotecฤ.
Instantaneul de mai jos prezintฤ o implementare a clasei SalesCountryReducer-
Code Explicaลฃie:
1. Definiศia clasei SalesCountryReducer-
public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {
Aici, primele douฤ tipuri de date, โTextโ ศi โIntWritableโ, reprezintฤ tipul de date al cheii-valoare de intrare a reductorului.
Rezultatul mapperului este sub forma , Aceastฤ ieศire a mapper-ului devine intrare pentru reducer. Aศadar, pentru a se alinia cu tipul sฤu de date, se folosesc aici tipuri de date Text ศi IntWritable.
Ultimele douฤ tipuri de date, โTextโ ศi โIntWritableโ, sunt tipurile de date de ieศire generate de reductor sub forma unei perechi cheie-valoare.
Fiecare clasฤ Reducer trebuie sฤ fie extinsฤ din clasa MapReduceBase ศi trebuie sฤ implementeze interfaศa Reducer.
2. Definirea funcศiei โreduceโ-
public void reduce( Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException {
O intrare pentru metoda reduce() este o cheie cu o listฤ de valori multiple.
De exemplu, รฎn cazul nostru, va fi...
, , , , , .
Acest lucru este dat reductorului ca
Deci, pentru a accepta argumente de aceastฤ formฤ, se utilizeazฤ primele douฤ tipuri de date, ศi anume Text ศi Iterator. Textul este un tip de date cheie ศi iterator. este un tip de date pentru o listฤ de valori pentru cheia respectivฤ.
Urmฤtorul argument este de tip OutputCollector care colecteazฤ ieศirea fazei reducฤtoare.
Metoda reduce() รฎncepe prin copierea valorii cheii ศi iniศializarea numฤrului de frecvenศฤ la 0.
Text key = t_key;int frequencyForCountry = 0;
Apoi, folosind bucla โwhileโ, iterฤm prin lista de valori asociate cheii ศi calculฤm frecvenศa finalฤ prin รฎnsumarea tuturor valorilor.
while (values.hasNext()) { // replace type of value with the actual type of our value IntWritable value = (IntWritable) values.next(); frequencyForCountry += value.get(); }
Acum, trimitem rezultatul cฤtre colectorul de ieศire sub forma unei chei ศi a numฤrului de frecvenศฤ obศinut.
Codul de mai jos face asta -
output.collect(key, new IntWritable(frequencyForCountry));
Explicaศia clasei SalesCountryDriver
รn aceastฤ secศiune, vom รฎnศelege implementarea clasei SalesCountryDriver.
1. รncepem prin a specifica un nume de pachet pentru clasa noastrฤ. SalesCountry este numele pachetului nostru. Reศineศi cฤ rezultatul compilฤrii, SalesCountryDriver.class, va intra รฎntr-un director numit cu acest nume de pachet: SalesCountry.
Aici este o linie care specificฤ numele pachetului urmat de cod pentru a importa pachete de bibliotecฤ.
2. Definiศi o clasฤ de driver care va crea un nou loc de muncฤ client, un obiect de configurare ศi va face publicitate claselor Mapper ศi Reducer.
Clasa de ศoferi este responsabilฤ pentru setarea jobului nostru MapReduce pentru a rula Hadoopรn aceastฤ clasฤ, specificฤm numele jobului, tipul de date de intrare/ieศire ศi numele claselor mapper ศi reducer.
3. รn fragmentul de cod de mai jos, setฤm directoarele de intrare ศi de ieศire care sunt folosite pentru a consuma setul de date de intrare ศi, respectiv, a produce ieศiri.
arg[0] ศi arg[1] sunt argumentele din linia de comandฤ transmise odatฤ cu o comandฤ datฤ รฎn MapReduce hands-on, adicฤ,
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
4. Declanศeazฤ-ne jobul
Codul de mai jos porneศte execuศia jobului MapReduce-
try { // Run the job JobClient.runJob(job_conf); } catch (Exception e) { e.printStackTrace(); }





















