Hadoop MapReduce-exempel: Först Java Program med Code

⚡ Smart sammanfattning

Hadoop MapReduce-program skrivs som tre Java klasser, en mapper, en reducer och en drivrutin, som kompileras, paketeras i en burk och skickas till klustret för att räkna försäljning per land.

  • 🔘 dataset: SalesJan2009.csv innehåller en transaktion per rad, med landet i den åttonde kommaseparerade kolumnen.
  • ☑️ Kartläggare: SalesMapper delar upp varje rad och skickar ut landet parat med det konstanta värdet ett.
  • Reducer: SalesCountryReducer summerar de som anländer för varje land till ett enda frekvensantal.
  • 🧪 Förare: SalesCountryDriver namnger jobbet, deklarerar nyckel- och värdetyper och kopplar båda klasserna tillsammans.
  • 🛠️ Kroppsbyggnad: Kompilera med javac, lägg till en Main-Class manifestpost och paketera sedan allt med jar cfm.
  • ⚠️ Springa: Kopiera CSV-filen till HDFS, skicka in jar-filen och läs part-00000 från utdatakatalogen.

Hadoop MapReduce-exempel skapar en första Java program

I den här handledningen lär du dig att använda Hadoop med MapReduce-exempel. Indata som används är FörsäljningJan2009.csvDen innehåller försäljningsrelaterad information såsom produktnamn, pris, betalningsmetod, stad och land för kunden. Målet är att ta reda på antalet sålda produkter i varje land.

Första Hadoop MapReduce-programmet

Nu i detta Handledning för MapReduce, kommer vi att skapa vår första Java MapReduce-program:

Skärmdumpen nedan visar rådata från SalesJan2009, där varje rad är en transaktion och landet finns i den åttonde kommaseparerade kolumnen.

Försäljningsdata för SalesJan2009 öppnades i ett kalkylblad som visar transaktionskolumner

Se till att du har Hadoop installerat. Innan du börjar med själva processen, ändra användaren till 'hduser' (det ID som användes vid konfigureringen av Hadoop — du kan byta till det användar-ID som användes under din egen Hadoop-konfiguration).

su - hduser_

Prompten ändras till hduser-kontot, som visas nedan.

Terminalprompten efter att ha bytt till hduser-kontot

Steg 1) Skapa projektkatalogen och källfilerna

Skapa en ny katalog med namnet MapReduceTutorial som visas i MapReduce-exemplet nedan.

sudo mkdir MapReduceTutorial

mkdir-kommandot skapar MapReduceTutorial-katalogen

Ge behörigheter

sudo chmod -R 777 MapReduceTutorial

chmod-kommando som ger fullständiga rättigheter på MapReduceTutorial

Skapa de tre Java källfilerna nedan i MapReduceTutorial. Observera att alla tre använder den äldre org.apache.hadoop.mapred API, som fortfarande levereras 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();
        }
    }
}

Ladda ner filer här

Arkivet expanderar till samma tre källfiler, som visas här.

Extracted-arkivlista SalesMapper, SalesCountryReducer och SalesCountryDriver

Kontrollera filbehörigheterna för alla dessa filer

Lång lista som visar filbehörigheter på de tre Java källfiler

Om läsbehörigheter saknas, bevilja dem:

chmod-kommandot lägger till läsbehörighet till Java källfiler

Steg 2) Exportera Hadoop-klassvägen

Exportera klassvägen som visas i Hadoop-exemplet nedan.

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 exporterade klassvägen visas tillbaka vid prompten.

Shell-prompt efter export av Hadoop CLASSPATH-variabeln

Steg 3) Kompilera Java filer

Kompilera Java filer (dessa filer finns i katalogen Final-MapReduceHandsOn). Deras klassfiler kommer att placeras i paketkatalogen.

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

javac-utdata visar en varning om föråldring under kompilering

Den här varningen kan säkert ignoreras – den rapporterar bara att mappred API är föråldrat.

Den här kompileringen skapar en katalog i den aktuella katalogen med namnet paketnamnet som anges i Java källfilen (dvs. SalesCountry i vårt fall) och lägg alla kompilerade klassfiler i den.

SalesCountry-paketkatalogen som innehåller de tre kompilerade klassfilerna

Steg 4) Skapa manifestfilen

Skapa en ny fil Manifest.txt

sudo gedit Manifest.txt

Lägg till följande rad i den:

Main-Class: SalesCountry.SalesCountryDriver

gedit fönster som visar Main-Class-posten i Manifest.txt

SalesCountry.SalesCountryDriver är namnet på huvudklassen. Observera att du måste trycka på Enter-tangenten i slutet av den här raden.

Steg 5) Packa klasserna i en burk

Skapa en Jar-fil

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

jar-kommandot skapar ProductSalePerCountry.jar från manifestet

Kontrollera att jar-filen är skapad

Kataloglista som bekräftar att ProductSalePerCountry.jar skapades

Steg 6) Starta Hadoop

Starta Hadoop

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

Steg 7) Kopiera indatafilen till HDFS

Kopiera filen SalesJan2009.csv till ~/inputMapReduce

Använd nu kommandot nedan för att kopiera ~/inputMapReduce till HDFS.

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

copyFromLocal-kommandoutdata med en varning för native-biblioteket

Vi kan lugnt ignorera denna varning.

Kontrollera om en fil faktiskt har kopierats eller inte.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

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

Steg 8) Kör MapReduce-jobbet

Kör MapReduce-jobbet

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

Konsolutdata medan SalePerCountry-jobbet körs på klustret

Detta kommer att skapa en utdatakatalog med namnet mapreduce_output_sales on HDFS. Innehållet i denna katalog kommer att vara en fil som innehåller produktförsäljning per land.

Steg 9) Läs resultaten

Resultatet kan ses via kommandogränssnittet som,

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

Land- och försäljningsantalspar utskrivna av kommandot hdfs dfs -cat

Resultaten kan också ses via ett webbgränssnitt som-

Öppet http://localhost:50070/ i en webbläsare. På Hadoop 3.x flyttades NameNode-webbgränssnittet till port 9870, så använd http://localhost:9870/ där istället.

NameNode webbgränssnitts hemsida i en webbläsare

Välj nu "Bläddra i filsystemet" och navigera till /mapreduce_output_sales

HDFS-filläsare som visar katalogen mapreduce_output_sales

Öppna del-r-00000

resultatfilen part-r-00000 öppnad i HDFS-webbläsaren

Förklaring av SalesMapper Class

Med jobbet på gång från början till slut går de kommande tre avsnitten igenom vad varje klass faktiskt gör.

I det här avsnittet kommer vi att förstå implementeringen av SalesMapper-klassen.

1. Vi börjar med att ange ett paketnamn för vår klass. SalesCountry är namnet på vårt paket. Observera att utdata från kompileringen, SalesMapper.class, kommer att hamna i en katalog med detta paketnamn: SalesCountry.

Därefter importerar vi bibliotekspaket.

Nedanstående ögonblicksbild visar en implementering av SalesMapper-klassen-

Redigerarvy av den fullständiga implementeringen av SalesMapper-klassen

Prov Code Förklaring:

1. SalesMapper Klassdefinition-

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

Varje mapper-klass måste utökas från MapReduceBase-klassen och den måste implementera Mapper-gränssnittet.

2. Definiera 'karta' funktion-

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

Huvuddelen av Mapper-klassen är en 'map()'-metod som accepterar fyra argument.

Vid varje anrop till metoden 'map()' skickas ett nyckel-värde-par ('key' och 'value' i denna kod).

Metoden 'map()' börjar med att dela upp inmatningstexten som tas emot som ett argument. Den delar upp varje rad i fält.

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

Här används ',' som avgränsare.

Efter detta bildas ett par med hjälp av en post vid det sjunde indexet i arrayen 'SingleCountryData' och värdet '1'.

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

Vi väljer posten vid 7:e index eftersom vi behöver landsdata och den finns vid 7:e index i arrayen 'SingleCountryData'.

Observera att våra indata är i formatet nedan (där Land är på 7:e index, med 0 som startindex) -

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

Utdata från mapper är återigen ett nyckel-värde-par som genereras med hjälp av 'collect()'-metoden i 'OutputCollector'.

Förklaring av SalesCountryReducer Class

I det här avsnittet kommer vi att förstå implementeringen av SalesCountryReducer-klassen.

1. Vi börjar med att ange ett namn på paketet för vår klass. SalesCountry är namnet på vårt paket. Observera att utdata från kompileringen, SalesCountryReducer.class, kommer att hamna i en katalog med detta paketnamn: SalesCountry.

Därefter importerar vi bibliotekspaket.

Nedanstående ögonblicksbild visar en implementering av SalesCountryReducer-klassen-

Redigerarvy av den fullständiga implementeringen av SalesCountryReducer-klassen

Code Förklaring:

1. SalesCountryReducer Class Definition-

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

Här är de två första datatyperna, 'Text' och 'IntWritable', datatypen för indatanyckelvärdet till reduceraren.

Utdata från mappern är i form av , Denna utdata från mappern blir indata till reduceraren. Så, för att anpassa den till dess datatyp, används Text och IntWritable som datatyper här.

De två sista datatyperna, 'Text' och 'IntWritable', är datatypen för utdata som genereras av reduceraren i form av ett nyckel-värde-par.

Varje reducer-klass måste utökas från MapReduceBase-klassen och den måste implementera Reducer-gränssnittet.

2. Definiera "reducera" funktion-

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

Indata till metoden reduce() är en nyckel med en lista med flera värden.

Till exempel, i vårt fall kommer det att vara-

, , , , , .

Detta ges till reducerare som

Så, för att acceptera argument av denna form används de två första datatyperna, nämligen Text och Iterator Text är en datatyp för nyckel och iterator är en datatyp för en lista med värden för den nyckeln.

Nästa argument är av typen OutputCollector som samlar in utgången från reducerfasen.

Metoden reduce() börjar med att kopiera nyckelvärdet och initiera frekvensantalet till 0.

Text key = t_key;
int frequencyForCountry = 0;

Sedan, med hjälp av 'while'-loopen, itererar vi igenom listan med värden som är associerade med nyckeln och beräknar den slutliga frekvensen genom att summera alla värden.

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

Nu skickar vi resultatet till utdatainsamlaren i form av nyckel och erhållet frekvensantal.

Nedanstående kod gör detta-

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

Förklaring av SalesCountryDriver Class

I det här avsnittet kommer vi att förstå implementeringen av SalesCountryDriver-klassen.

1. Vi börjar med att ange ett paketnamn för vår klass. SalesCountry är namnet på vårt paket. Observera att utdata från kompileringen, SalesCountryDriver.class, kommer att hamna i en katalog med detta paketnamn: SalesCountry.

Här är en rad som anger paketnamn följt av kod för att importera bibliotekspaket.

Paketdeklaration och importutlåtanden för SalesCountryDriver

2. Definiera en drivrutinsklass som skapar ett nytt klientjobb, konfigurationsobjekt och annonserar Mapper- och Reducer-klasser.

Förarklassen ansvarar för att ställa in vårt MapReduce-jobb att köra in HadoopI den här klassen anger vi jobbnamn, datatyp för indata/utdata och namn på mapper- och reducer-klasser.

Drivrutinskod som ställer in jobbnamn, nyckel- och värdeklasser och format

3. I kodavsnittet nedan ställer vi in ​​indata- och utdatakataloger som används för att konsumera indatauppsättning respektive producera utdata.

arg[0] och arg[1] är kommandoradsargumenten som skickas med ett kommando som ges i MapReduce hands-on, dvs.

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

Drivrutinskod som skickar in- och utmatningsvägarna från kommandoraden

4. Trigga vårt jobb

Nedanstående kod startar körningen av MapReduce-jobbet-

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

Vanliga frågor

Alla tre klasserna importerar org.apache.hadoop.mapred, det ursprungliga MapReduce API:et. Hadoop 3.x levereras och körs fortfarande, men nytt arbete skrivs normalt mot org.apache.hadoop.mapreduce, som ersätter JobConf och JobClient med Configuration och Job.

Modeller som tränats på tidigare jobbhistorik förutspår körtid, föreslår delningsstorlekar och antal reducerare och identifierar skevheter innan en körning avslutas. De flaggar också jobb vars räknare för spillda poster eller misslyckade uppgifter avviker utanför det vanliga intervallet.

Copilot slutför de repetitiva delarna väl: klassignaturer, generiska program, importer och drivrutinskonfigurationsanrop. Affärslogiken, som vilken kolumn som innehåller landet, måste fortfarande kontrolleras mot det verkliga schemat av en utvecklare.

Hadoop vägrar att skriva till en befintlig utdatakatalog så att färdiga resultat aldrig skrivs över. Ta bort mapreduce_output_sales rekursivt med hdfs dfs -rm -r, eller skicka en annan utdatasökväg vid nästa körning.

Jar-verktyget ignorerar en sista manifestrad som inte har någon radavslutning, så Main-Class tas bort i tysthet. Jar-filen byggs sedan utan fel men misslyckas vid körning eftersom ingen huvudklass registreras.

Ja. Exemplet fäster 2.2.0 i varje jar-namn. I en annan utgåva, ersätt den versionen, eller använd helt enkelt utdata från hadoop-klassvägen, som skriver ut varje jar-fil som den installerade distributionen behöver.

Nej. En pseudodistribuerad installation med en enda nod kör den oförändrad, eftersom kommandona bara antar att HDFS och YARN startas. Samma jar skickas till ett verkligt kluster med flera noder utan någon kodändring.

Ändra arrayindexet i mappern från 7 till 5, eftersom City är den sjätte kolumnen i SalesJan2009.csv. Kompilera om, bygg om burken och kör jobbet mot en ny utdatakatalog.

Sammanfatta detta inlägg med: