Przykład Hadoop MapReduce: Pierwszy Java Program z Code

⚡ Inteligentne podsumowanie

Programy Hadoop MapReduce są pisane jako trzy Java klasy, maper, reduktor i sterownik, które są kompilowane, pakowane do pliku jar i przesyłane do klastra w celu zliczenia sprzedaży według kraju.

  • 🔘 Zestaw danych: W pliku SalesJan2009.csv każda linia zawiera jedną transakcję, a kraj znajduje się w ósmej kolumnie, rozdzielonej przecinkami.
  • Twórca map: SalesMapper rozdziela każdy wiersz i emituje nazwę kraju sparowaną ze stałą wartością jeden.
  • Reduktor: SalesCountryReducer sumuje liczbę przybyć dla każdego kraju w jednym zliczeniu częstotliwości.
  • 🧪 Kierowca: SalesCountryDriver nadaje nazwę zadaniu, deklaruje typy kluczy i wartości, a następnie łączy obie klasy.
  • 🛠️. Budować: Skompiluj za pomocą javac, dodaj wpis manifestu Main-Class, a następnie spakuj wszystko za pomocą pliku jar cfm.
  • ⚠️ Biegać: Skopiuj plik CSV do HDFS, prześlij plik jar i odczytaj part-00000 z katalogu wyjściowego.

Przykład tworzenia pierwszego Hadoop MapReduce Java program

W tym samouczku nauczysz się używać Hadoop z przykładami MapReduce. Stosowane dane wejściowe to SprzedażJan2009.csvZawiera informacje związane ze sprzedażą, takie jak nazwa produktu, cena, sposób płatności, miasto i kraj klienta. Celem jest ustalenie liczby produktów sprzedanych w każdym kraju.

Pierwszy program Hadoop MapReduce

Teraz w tym Samouczek MapReduce, stworzymy nasz pierwszy Java Program MapReduce:

Poniższy zrzut ekranu pokazuje surowe dane SalesJan2009, gdzie każdy wiersz odpowiada jednej transakcji, a kraj znajduje się w ósmej kolumnie rozdzielonej przecinkami.

Dane sprzedaży ze stycznia 2009 r. otwarte w arkuszu kalkulacyjnym, pokazującym kolumny transakcji

Upewnij się, że masz zainstalowany Hadoop. Zanim rozpoczniesz właściwy proces, zmień użytkownika na „hduser” (identyfikator używany podczas konfiguracji Hadoop — możesz zmienić identyfikator użytkownika używany podczas własnej konfiguracji Hadoop).

su - hduser_

Monit zmieni się na konto hduser, jak pokazano poniżej.

Monit terminala po przełączeniu się na konto hduser

Krok 1) Utwórz katalog projektu i pliki źródłowe

Utwórz nowy katalog o nazwie MapReduceTutorial, jak pokazano w poniższym przykładzie MapReduce.

sudo mkdir MapReduceTutorial

Polecenie mkdir tworzące katalog MapReduceTutorial

Przyznaj uprawnienia

sudo chmod -R 777 MapReduceTutorial

Polecenie chmod przyznające pełne uprawnienia w MapReduceTutorial

Utwórz trzy Java Pliki źródłowe znajdują się poniżej w MapReduceTutorial. Należy pamiętać, że wszystkie trzy korzystają ze starszej wersji org.apache.hadoop.mapred API, które nadal jest zawarte w 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();
        }
    }
}

Pobierz pliki tutaj

Archiwum rozrasta się do trzech takich samych plików źródłowych, jak pokazano tutaj.

ExtracArchiwum TED zawierające SalesMapper, SalesCountryReducer i SalesCountryDriver

Sprawdź uprawnienia do wszystkich tych plików

Długa lista pokazująca uprawnienia do plików dla trzech Java pliki źródłowe

Jeżeli brakuje uprawnień „odczytu”, to przyznaj je:

polecenie chmod dodające uprawnienie do odczytu Java pliki źródłowe

Krok 2) Eksportuj ścieżkę klas Hadoop

Eksportuj ścieżkę klas, jak pokazano w poniższym przykładzie Hadoop.

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/*"

Wyeksportowana ścieżka klas zostanie wyświetlona w wierszu poleceń.

Monit powłoki po wyeksportowaniu zmiennej Hadoop CLASSPATH

Krok 3) Kompilacja Java pliki

Skompiluj Java Pliki (te pliki znajdują się w katalogu Final-MapReduceHandsOn). Ich pliki klas zostaną umieszczone w katalogu pakietów.

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

dane wyjściowe javac pokazujące ostrzeżenie o wycofaniu podczas kompilacji

To ostrzeżenie można bezpiecznie zignorować — informuje ono jedynie o tym, że interfejs API mapred jest przestarzały.

Ta kompilacja utworzy katalog w bieżącym katalogu o nazwie odpowiadającej nazwie pakietu określonej w Java plik źródłowy (w naszym przypadku SalesCountry) i umieścić w nim wszystkie skompilowane pliki klas.

Katalog pakietu SalesCountry zawierający trzy skompilowane pliki klas

Krok 4) Utwórz plik manifestu

Utwórz nowy plik Manifest.txt

sudo gedit Manifest.txt

Dodaj do niego następujący wiersz:

Main-Class: SalesCountry.SalesCountryDriver

gedit okno pokazujące wpis klasy głównej w Manifest.txt

SalesCountry.SalesCountryDriver to nazwa klasy głównej. Pamiętaj, że musisz nacisnąć klawisz Enter na końcu tego wiersza.

Krok 5) Spakuj klasy do pliku jar

Utwórz plik Jar

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

polecenie jar tworzące plik ProductSalePerCountry.jar z manifestu

Sprawdź, czy plik jar został utworzony

Wpis w katalogu potwierdzający utworzenie pliku ProductSalePerCountry.jar

Krok 6) Uruchom Hadoop

Uruchom Hadoopa

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

Krok 7) Skopiuj plik wejściowy do HDFS

Skopiuj plik SalesJan2009.csv do ~/inputMapReduce

Teraz użyj poniższego polecenia, aby skopiować ~/inputMapReduce do HDFS.

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

Wyjście polecenia copyFromLocal z ostrzeżeniem dotyczącym biblioteki natywnej

Możemy spokojnie zignorować to ostrzeżenie.

Sprawdź, czy plik rzeczywiście został skopiowany, czy nie.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls wyświetla SalesJan2009.csv w katalogu inputMapReduce

Krok 8) Uruchom zadanie MapReduce

Uruchom zadanie MapReduce

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

Dane wyjściowe konsoli podczas działania zadania SalePerCountry w klastrze

Spowoduje to utworzenie katalogu wyjściowego o nazwie mapreduce_output_sales on HDFS. Zawartość tego katalogu będzie plikiem zawierającym sprzedaż produktów w poszczególnych krajach.

Krok 9) Odczytaj wyniki

Wynik można zobaczyć za pomocą interfejsu poleceń w następujący sposób:

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

Pary kraju i liczby sprzedaży wydrukowane poleceniem hdfs dfs -cat

Wyniki można również zobaczyć za pośrednictwem interfejsu internetowego, ponieważ:

Otwarte http://localhost:50070/ w przeglądarce internetowej. W Hadoop 3.x interfejs użytkownika NameNode został przeniesiony do portu 9870, więc użyj http://localhost:9870/ tam zamiast.

Strona główna interfejsu internetowego NameNode w przeglądarce

Teraz wybierz „Przeglądaj system plików” i przejdź do /mapreduce_output_sales

Przeglądarka plików HDFS pokazująca katalog mapreduce_output_sales

Otwórz część-r-00000

plik wyników part-r-00000 otwarty w przeglądarce HDFS

Objaśnienie klasy SalesMapper

Po zakończeniu pracy od początku do końca w kolejnych trzech sekcjach omówiono, co właściwie robi każda klasa.

W tej sekcji zrozumiemy implementację klasy SalesMapper.

1. Zaczynamy od podania nazwy pakietu dla naszej klasy. SalesCountry to nazwa naszego pakietu. Należy pamiętać, że wynik kompilacji, SalesMapper.class, trafi do katalogu o nazwie odpowiadającej nazwie pakietu: SalesCountry.

Następnie importujemy pakiety bibliotek.

Poniższy zrzut ekranu przedstawia implementację klasy SalesMapper

Widok edytora kompletnej implementacji klasy SalesMapper

Próba Code Wyjaśnienie:

1. Definicja klasy SalesMapper-

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

Każda klasa mappera musi być rozszerzeniem klasy MapReduceBase i musi implementować interfejs Mapper.

2. Definiowanie funkcji „mapa”-

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

Główną część klasy Mapper stanowi metoda 'map()', która przyjmuje cztery argumenty.

Przy każdym wywołaniu metody 'map()' przekazywana jest para klucz-wartość ('klucz' i 'wartość' w tym kodzie).

Metoda „map()” rozpoczyna się od podzielenia tekstu wejściowego otrzymanego jako argument. Dzieli ona każdy wiersz na pola.

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

W tym przypadku znak ',' użyto jako rozgranicznika.

Następnie tworzona jest para przy użyciu rekordu o 7. indeksie tablicy „SingleCountryData” i wartości „1”.

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

Wybieramy rekord o 7. indeksie, ponieważ potrzebujemy danych krajowych, a znajdują się one pod 7. indeksem w tablicy „SingleCountryData”.

Należy pamiętać, że nasze dane wejściowe są w poniższym formacie (gdzie Kraj znajduje się na 7. indeksie, a 0 jest indeksem początkowym)-

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

Wynikiem działania mappera jest ponownie para klucz-wartość emitowana przy użyciu metody „collect()” obiektu „OutputCollector”.

Objaśnienie klasy SalesCountryReducer

W tej sekcji zrozumiemy implementację klasy SalesCountryReducer.

1. Zaczynamy od podania nazwy pakietu dla naszej klasy. SalesCountry to nazwa naszego pakietu. Należy pamiętać, że wynik kompilacji, SalesCountryReducer.class, trafi do katalogu o nazwie odpowiadającej nazwie pakietu: SalesCountry.

Następnie importujemy pakiety bibliotek.

Poniższy zrzut ekranu przedstawia implementację klasy SalesCountryReducer

Widok edytora kompletnej implementacji klasy SalesCountryReducer

Code Wyjaśnienie:

1. Definicja klasy SalesCountryReducer-

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

W tym przypadku pierwsze dwa typy danych, „Tekst” i „IntWritable”, są typem danych wejściowej pary klucz-wartość dla reduktora.

Wyjście mappera ma postać , To wyjście mappera staje się wejściem do reduktora. Zatem, aby dopasować się do jego typu danych, użyto tutaj typów danych Text i IntWritable.

Ostatnie dwa typy danych, „Tekst” i „IntWritable”, to typ danych wyjściowych generowanych przez reduktor w formie pary klucz-wartość.

Każda klasa reduktora musi być rozszerzeniem klasy MapReduceBase i musi implementować interfejs Reducer.

2. Definiowanie funkcji „zmniejsz”

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

Danymi wejściowymi metody reduce() są klucze z listą wielu wartości.

W naszym przypadku będzie to na przykład:

, , , , , .

To jest podawane reduktorowi jako

Aby zaakceptować argumenty w tej formie, najpierw używa się dwóch typów danych, tj. Tekst i Iterator . Tekst jest typem danych klucza i iteratora jest typem danych dla listy wartości dla danego klucza.

Następny argument jest typu OutputCollector który zbiera wyjście fazy reduktora.

Metoda reduce() rozpoczyna od skopiowania wartości klucza i zainicjowania licznika częstotliwości na 0.

Text key = t_key;
int frequencyForCountry = 0;

Następnie, używając pętli 'while', przeglądamy listę wartości skojarzonych z kluczem i obliczamy końcową częstotliwość, sumując wszystkie wartości.

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

Teraz przekazujemy wynik do kolektora wyjściowego w postaci klucza i uzyskanej liczby częstotliwości.

Poniższy kod to robi-

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

Objaśnienie klasy SalesCountryDriver

W tej sekcji zrozumiemy implementację klasy SalesCountryDriver.

1. Zaczynamy od podania nazwy pakietu dla naszej klasy. SalesCountry to nazwa naszego pakietu. Należy pamiętać, że wynik kompilacji, SalesCountryDriver.class, zostanie umieszczony w katalogu o nazwie odpowiadającej nazwie pakietu: SalesCountry.

Oto linia określająca nazwę pakietu, po której następuje kod importowania pakietów bibliotek.

Deklaracja pakietu i oświadczenia importowe SalesCountryDriver

2. Zdefiniuj klasę sterownika, która utworzy nowe zadanie klienta, obiekt konfiguracyjny i rozgłosi klasy Mapper i Reduktor.

Klasa sterownika jest odpowiedzialna za ustawienie naszego zadania MapReduce do uruchomienia HadoopW tej klasie określamy nazwę zadania, typ danych wejścia/wyjścia oraz nazwy klas mapperów i reducerów.

Kod sterownika ustawiający nazwę zadania, klasy kluczy i wartości oraz formaty

3. W poniższym fragmencie kodu ustawiamy katalogi wejściowe i wyjściowe, które służą odpowiednio do pobierania wejściowego zbioru danych i generowania danych wyjściowych.

arg[0] i arg[1] to argumenty wiersza poleceń przekazywane wraz z poleceniem w praktycznym użyciu MapReduce, tj.

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

Kod sterownika przekazujący ścieżki wejściowe i wyjściowe z wiersza poleceń

4. Uruchom naszą pracę

Poniższy kod rozpoczyna wykonywanie zadania MapReduce-

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

FAQ

Wszystkie trzy klasy importują org.apache.hadoop.mapred, oryginalne API MapReduce. Hadoop 3.x nadal jest dostępny i działa, ale nowe prace są zazwyczaj pisane z wykorzystaniem org.apache.hadoop.mapreduce, które zastępuje JobConf i JobClient interfejsami Configuration i Job.

Modele trenowane na podstawie historii zadań przewidują czas wykonania, sugerują rozmiary podziału i liczbę reduktorów oraz wykrywają przekłamania przed zakończeniem wykonania. Sygnalizują również zadania, których liczniki rozlanych rekordów lub nieudanych zadań wykraczają poza typowy zakres.

Copilot dobrze radzi sobie z powtarzalnymi elementami: sygnaturami klas, typami generycznymi, importami i wywołaniami konfiguracji sterownika. Logika biznesowa, taka jak kolumna zawierająca kraj, nadal musi zostać sprawdzona przez programistę pod kątem zgodności ze schematem rzeczywistym.

Hadoop odmawia zapisu do istniejącego katalogu wyjściowego, więc gotowe wyniki nigdy nie zostaną nadpisane. Usuń rekurencyjnie plik mapreduce_output_sales za pomocą polecenia hdfs dfs -rm -r lub przekaż inną ścieżkę wyjściową przy następnym uruchomieniu.

Narzędzie jar ignoruje końcowy wiersz manifestu, który nie zawiera znaku zakończenia, więc Main-Class jest automatycznie usuwany. Plik jar kompiluje się bez błędów, ale kończy się niepowodzeniem w czasie wykonywania, ponieważ nie jest rejestrowana żadna klasa główna.

Tak. Przykład przypina 2.2.0 do każdej nazwy pliku jar. W innej wersji zastąp tę wersję tą wersją lub po prostu skorzystaj z wyników z Hadoop ClassPath, które wyświetlają wszystkie pliki jar potrzebne zainstalowanej dystrybucji.

Nie. Instalacja pseudorozproszona z jednym węzłem uruchamia ją bez zmian, ponieważ polecenia zakładają jedynie uruchomienie HDFS i YARN. Ten sam plik jar jest przesyłany do rzeczywistego klastra wielowęzłowego bez żadnych zmian w kodzie.

Zmień indeks tablicy w maperze z 7 na 5, ponieważ „City” to szósta kolumna pliku SalesJan2009.csv. Przekompiluj, odbuduj plik jar i uruchom zadanie w nowym katalogu wyjściowym.

Podsumuj ten post następująco: