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.
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.
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.
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
Przyznaj uprawnienia
sudo chmod -R 777 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(); } } }
Archiwum rozrasta się do trzech takich samych plików źródłowych, jak pokazano tutaj.
Sprawdź uprawnienia do wszystkich tych plików
Jeżeli brakuje uprawnień „odczytu”, to przyznaj je:
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ń.
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
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.
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
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
Sprawdź, czy plik jar został utworzony
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 /
Możemy spokojnie zignorować to ostrzeżenie.
Sprawdź, czy plik rzeczywiście został skopiowany, czy nie.
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
Krok 8) Uruchom zadanie MapReduce
Uruchom zadanie MapReduce
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
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
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.
Teraz wybierz „Przeglądaj system plików” i przejdź do /mapreduce_output_sales
Otwórz część-r-00000
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
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
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.
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.
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
4. Uruchom naszą pracę
Poniższy kod rozpoczyna wykonywanie zadania MapReduce-
try { // Run the job JobClient.runJob(job_conf); } catch (Exception e) { e.printStackTrace(); }





















