Hadoop MapReduce-eksempel: Først Java Program med Code
⚡ Smart oppsummering
Hadoop MapReduce-programmer er skrevet som tre Java klasser, en mapper, en reduksjonsverktøy og en driver, som kompileres, pakkes i en JAR og sendes til klyngen for å telle salg per land.

I denne opplæringen lærer du å bruke Hadoop med MapReduce-eksempler. Inndataene som brukes er SalgJan2009.csvDen inneholder salgsrelatert informasjon som produktnavn, pris, betalingsmåte, by og land kunden tilhører. Målet er å finne ut hvor mange produkter som selges i hvert land.
Første Hadoop MapReduce-program
Nå i dette MapReduce opplæring, vil vi lage vår første Java MapReduce-program:
Skjermbildet nedenfor viser rådataene for SalesJan2009, der hver linje er én transaksjon og landet står i den åttende kommaseparerte kolonnen.
Sørg for at du har Hadoop installert. Før du starter med selve prosessen, endre brukeren til 'hduser' (ID-en som ble brukt under konfigurasjonen av Hadoop – du kan bytte til bruker-ID-en som ble brukt under din egen Hadoop-konfigurasjon).
su - hduser_
Ledeteksten endres til hduser-kontoen, som vist nedenfor.
Trinn 1) Opprett prosjektkatalogen og kildefilene
Opprett en ny katalog med navnet MapReduceTutorial som vist i MapReduce-eksemplet nedenfor.
sudo mkdir MapReduceTutorial
Gi tillatelser
sudo chmod -R 777 MapReduceTutorial
Lag de tre Java kildefilene nedenfor i MapReduceTutorial. Merk at alle tre bruker den eldre org.apache.hadoop.mapred API, som fortsatt leveres med 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(); } } }
Arkivet utvides til de samme tre kildefilene, som vist her.
Sjekk filtillatelsene til alle disse filene
Hvis «lese»-tillatelser mangler, gi dem:
Trinn 2) Eksporter Hadoop-klassestien
Eksporter klassebanen som vist i Hadoop-eksemplet nedenfor.
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/*"
Den eksporterte klassebanen sendes tilbake ved ledeteksten.
Trinn 3) Kompiler Java filer
Kompiler Java filer (disse filene finnes i katalogen Final-MapReduceHandsOn). Klassefilene deres vil bli lagt i pakkekatalogen.
javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java
Denne advarselen kan trygt ignoreres – den rapporterer bare at mapped API-et er avskrevet.
Denne kompileringen vil opprette en katalog i gjeldende katalog med navnet pakkenavnet som er spesifisert i Java kildefilen (dvs. SalesCountry i vårt tilfelle) og legg alle kompilerte klassefiler i den.
Trinn 4) Opprett manifestfilen
Opprett en ny fil Manifest.txt
sudo gedit Manifest.txt
Legg til følgende linje i den:
Main-Class: SalesCountry.SalesCountryDriver
SalesCountry.SalesCountryDriver er navnet på hovedklassen. Vær oppmerksom på at du må trykke enter-tasten på slutten av denne linjen.
Trinn 5) Pakk klassene i en krukke
Lag en Jar-fil
jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class
Sjekk at jar-filen er opprettet
Trinn 6) Start Hadoop
Start Hadoop
$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh
Trinn 7) Kopier inndatafilen til HDFS
Kopier filen SalesJan2009.csv til ~/inputMapReduce
Bruk nå kommandoen nedenfor til å kopiere ~/inputMapReduce til HDFS.
$HADOOP_HOME/bin/hdfs dfs -copyFromLocal ~/inputMapReduce /
Vi kan trygt ignorere denne advarselen.
Kontroller om en fil faktisk er kopiert eller ikke.
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
Trinn 8) Kjør MapReduce-jobben
Kjør MapReduce-jobben
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
Dette vil opprette en utdatakatalog kalt mapreduce_output_sales on HDFS. Innholdet i denne katalogen vil være en fil som inneholder produktsalg per land.
Trinn 9) Les resultatene
Resultatet kan sees gjennom kommandogrensesnittet som,
$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000
Resultatene kan også sees via et nettgrensesnitt som-
Open http://localhost:50070/ i en nettleser. På Hadoop 3.x ble NameNode-nettgrensesnittet flyttet til port 9870, så bruk http://localhost:9870/ der i stedet.
Velg nå «Bla gjennom filsystemet» og naviger til /mapreduce_output_sales
Åpen del-r-00000
Forklaring av SalesMapper Class
Med jobben på gang fra ende til annen, går de neste tre delene gjennom hva hver klasse faktisk gjør.
I denne delen skal vi forstå implementeringen av SalesMapper-klassen.
1. Vi begynner med å spesifisere et pakkenavn for klassen vår. SalesCountry er navnet på pakken vår. Vær oppmerksom på at resultatet av kompileringen, SalesMapper.class, vil gå inn i en katalog med samme pakkenavn: SalesCountry.
Etterfulgt av dette importerer vi bibliotekpakker.
Øyeblikksbildet nedenfor viser en implementering av SalesMapper-klassen-
Eksempel Code Forklaring:
1. SalesMapper Klassedefinisjon-
public class SalesMapper extends MapReduceBase implements Mapper<LongWritable, Text, Text, IntWritable> {
Hver mapper-klasse må utvides fra MapReduceBase-klassen, og den må implementere Mapper-grensesnittet.
2. Definere 'kart' funksjon-
public void map(LongWritable key, Text value, OutputCollector<Text, IntWritable> output, Reporter reporter) throws IOException
Hoveddelen av Mapper-klassen er en 'map()'-metode som aksepterer fire argumenter.
Ved hvert kall til metoden 'map()' sendes et nøkkel-verdi-par ('key' og 'value' i denne koden).
Metoden 'map()' begynner med å dele inndatateksten som mottas som et argument. Den deler hver linje inn i felt.
String valueString = value.toString(); String[] SingleCountryData = valueString.split(",");
Her brukes ',' som skilletegn.
Etter dette dannes et par ved hjelp av en post ved den 7. indeksen i arrayet 'SingleCountryData' og verdien '1'.
output.collect(new Text(SingleCountryData[7]), one);
Vi velger posten på 7. indeks fordi vi trenger landsdata, og den ligger på 7. indeks i arrayet 'SingleCountryData'.
Vær oppmerksom på at inndataene våre er i formatet nedenfor (der Land er på 7. indeks, med 0 som startindeks) -
Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude
En utdata fra mapper er igjen et nøkkel-verdi-par som sendes ut ved hjelp av 'collect()'-metoden til 'OutputCollector'.
Forklaring av SalesCountryReducer Class
I denne delen skal vi forstå implementeringen av SalesCountryReducer-klassen.
1. Vi begynner med å spesifisere et navn på pakken for klassen vår. SalesCountry er navnet på pakken vår. Vær oppmerksom på at resultatet av kompileringen, SalesCountryReducer.class, vil gå inn i en katalog med samme pakkenavn: SalesCountry.
Etterfulgt av dette importerer vi bibliotekpakker.
Øyeblikksbildet nedenfor viser en implementering av SalesCountryReducer-klassen-
Code Forklaring:
1. SalesCountryReducer Klassedefinisjon-
public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {
Her er de to første datatypene, «Tekst» og «IntWritable», datatypen for input-nøkkelverdien til reduseringsverktøyet.
Utdataene fra mapperen er i form av , Denne utgangen fra mapper blir input til reduseringsverktøyet. Så, for å justere med datatypen, brukes Text og IntWritable som datatyper her.
De to siste datatypene, «Tekst» og «IntWritable», er datatypen for utdata generert av reduseringsverktøyet i form av et nøkkel-verdi-par.
Hver reduksjonsklasse må utvides fra MapReduceBase-klassen, og den må implementere Reducer-grensesnittet.
2. Definere 'redusere' funksjon-
public void reduce( Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException {
En input til reduce()-metoden er en nøkkel med en liste over flere verdier.
For eksempel, i vårt tilfelle vil det være-
, , , , , .
Dette gis til reduseringsenheten som
Så, for å akseptere argumenter av denne formen, brukes de to første datatypene, nemlig tekst og iterator Tekst er en datatype for nøkkel og iterator er en datatype for en liste over verdier for den nøkkelen.
Det neste argumentet er av typen OutputCollector som samler utgangen fra reduksjonsfasen.
reduce()-metoden begynner med å kopiere nøkkelverdien og initialisere frekvensantall til 0.
Text key = t_key;int frequencyForCountry = 0;
Deretter, ved hjelp av «while»-løkken, itererer vi gjennom listen over verdier knyttet til nøkkelen og beregner den endelige frekvensen ved å summere alle verdiene.
while (values.hasNext()) { // replace type of value with the actual type of our value IntWritable value = (IntWritable) values.next(); frequencyForCountry += value.get(); }
Nå sender vi resultatet til utdatasamleren i form av nøkkel og oppnådd frekvensantall.
Nedenfor koden gjør dette-
output.collect(key, new IntWritable(frequencyForCountry));
Forklaring av SalesCountryDriver Class
I denne delen skal vi forstå implementeringen av SalesCountryDriver-klassen.
1. Vi begynner med å spesifisere et pakkenavn for klassen vår. SalesCountry er navnet på pakken vår. Vær oppmerksom på at resultatet av kompileringen, SalesCountryDriver.class, vil gå inn i en katalog med samme pakkenavn: SalesCountry.
Her er en linje som spesifiserer pakkenavn etterfulgt av kode for å importere bibliotekpakker.
2. Definer en driverklasse som vil opprette en ny klientjobb, konfigurasjonsobjekt og annonsere Mapper- og Reducer-klasser.
Førerklassen er ansvarlig for å sette MapReduce-jobben vår til å kjøre inn HadoopI denne klassen spesifiserer vi jobbnavn, datatype for input/output og navn på mapper- og reducer-klasser.
3. I kodebiten nedenfor setter vi inn- og utdatakataloger som brukes til å konsumere henholdsvis input-datasett og produsere utdata.
arg[0] og arg[1] er kommandolinjeargumentene som sendes med en kommando gitt i MapReduce hands-on, dvs.
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
4. Trigger jobben vår
Koden nedenfor starter kjøringen av MapReduce-jobben-
try { // Run the job JobClient.runJob(job_conf); } catch (Exception e) { e.printStackTrace(); }





















