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.

  • 🔘 Datasett: SalesJan2009.csv inneholder én transaksjon per linje, med landet i den åttende kommaseparerte kolonnen.
  • ☑️ Kartlegger: SalesMapper deler hver linje og sender ut landet parret med den konstante verdien én.
  • Reducer: SalesCountryReducer summerer antallet som ankommer for hvert land til én frekvenstelling.
  • 🧪 Sjåfør: SalesCountryDriver navngir jobben, deklarerer nøkkel- og verdityper og kobler begge klassene sammen.
  • 🛠️ Bygge: Kompiler med javac, legg til en manifestoppføring i Main-Class, og pakk deretter alt med jar cfm.
  • ⚠️ Løpe: Kopier CSV-filen til HDFS, send inn jar-filen og les part-00000 fra utdatakatalogen.

Hadoop MapReduce-eksempel på å lage en første Java program

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.

Salgsdata for SalesJan2009 åpnet i et regneark som viser transaksjonskolonner

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.

Terminalledetekst etter bytte til hduser-kontoen

Trinn 1) Opprett prosjektkatalogen og kildefilene

Opprett en ny katalog med navnet MapReduceTutorial som vist i MapReduce-eksemplet nedenfor.

sudo mkdir MapReduceTutorial

mkdir-kommandoen som oppretter MapReduceTutorial-katalogen

Gi tillatelser

sudo chmod -R 777 MapReduceTutorial

chmod-kommando som gir fulle tillatelser på 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();
        }
    }
}

Last ned filer her

Arkivet utvides til de samme tre kildefilene, som vist her.

Extracted-arkivliste SalesMapper, SalesCountryReducer og SalesCountryDriver

Sjekk filtillatelsene til alle disse filene

Lang liste som viser filtillatelser på de tre Java kildefiler

Hvis «lese»-tillatelser mangler, gi dem:

chmod-kommandoen legger til lesetillatelse til Java kildefiler

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.

Shell-ledetekst etter eksport av Hadoop CLASSPATH-variabelen

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

javac-utdata som viser en advarsel om avskrivning under kompilering

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.

SalesCountry-pakkekatalogen som inneholder de tre kompilerte klassefilene

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

gedit vindu som viser Main-Class-oppføringen i Manifest.txt

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

jar-kommandobygging ProductSalePerCountry.jar fra manifestet

Sjekk at jar-filen er opprettet

Katalogoppføring som bekrefter at ProductSalePerCountry.jar ble 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 /

copyFromLocal-kommandoutdata med en advarsel fra native-biblioteket

Vi kan trygt ignorere denne advarselen.

Kontroller om en fil faktisk er kopiert eller ikke.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls liste SalesJan2009.csv i inputMapReduce-katalogen

Trinn 8) Kjør MapReduce-jobben

Kjør MapReduce-jobben

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

Konsollutdata mens SalePerCountry-jobben kjører på klyngen

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

Land- og salgstallpar skrevet ut av hdfs dfs -cat-kommandoen

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.

NameNode webgrensesnitt hjemmeside i en nettleser

Velg nå «Bla gjennom filsystemet» og naviger til /mapreduce_output_sales

HDFS-filleser som viser mappen mapreduce_output_sales

Åpen del-r-00000

Resultatfilen part-r-00000 åpnet i HDFS-nettleseren

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-

Redigeringsvisning av den komplette implementeringen 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-

Redigeringsvisning av den komplette implementeringen 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.

Pakkedeklarasjon og importerklæringer for SalesCountryDriver

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.

Driverkode som angir jobbnavn, nøkkel- og verdiklasser og formater

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

Driverkode som sender inn- og utgangsbanene fra kommandolinjen

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();
}

Spørsmål og svar

Alle tre klassene importerer org.apache.hadoop.mapred, det originale MapReduce API-et. Hadoop 3.x sender og kjører det fortsatt, men nytt arbeid skrives vanligvis mot org.apache.hadoop.mapreduce, som erstatter JobConf og JobClient med Configuration og Job.

Modeller trent på tidligere jobbhistorikk forutsier kjøretid, foreslår delte størrelser og antall reduseringsenheter, og oppdager skjevheter før en kjøring er fullført. De flagger også jobber der tellerne for sølte poster eller mislykkede oppgaver avviker utenfor det vanlige området.

Copilot fullfører de repeterende delene på en god måte: klassesignaturer, generiske programmer, import og driverkonfigurasjonskall. Forretningslogikken, som hvilken kolonne som inneholder landet, må fortsatt kontrolleres mot det virkelige skjemaet av en utvikler.

Hadoop nekter å skrive inn i en eksisterende utdatamappe, slik at ferdige resultater aldri blir overskrevet. Slett mapreduce_output_sales rekursivt med hdfs dfs -rm -r, eller send en annen utdatasti ved neste kjøring.

Jar-verktøyet ignorerer en endelig manifestlinje som ikke har noen linjeavslutning, så Main-Class slettes stille. Jar-filen bygger deretter uten feil, men mislykkes ved kjøretid fordi ingen hovedklasse er registrert.

Ja. Eksemplet fester 2.2.0 i hvert jar-navn. I en annen utgivelse, erstatt den versjonen, eller bruk ganske enkelt utdataene fra hadoop-klassestien, som skriver ut alle jar-koder som den installerte distribusjonen trenger.

Nei. En pseudodistribuert installasjon med én node kjører den uendret, fordi kommandoene bare antar at HDFS og YARN er startet. Den samme jar-filen sender til en ekte klynge med flere noder uten noen kodeendring.

Endre arrayindeksen i mapperen fra 7 til 5, fordi City er den sjette kolonnen i SalesJan2009.csv. Kompiler på nytt, gjenoppbygg JAR-filen og kjør jobben mot en ny utdatakatalog.

Oppsummer dette innlegget med: