Exemplu Hadoop MapReduce: Primul Java Program cu Code

⚡ Rezumat inteligent

Programele Hadoop MapReduce sunt scrise în trei Java clase, un mapper, un reducer și un driver, care sunt compilate, împachetate într-un fișier jar și trimise clusterului pentru a număra vânzările pe țară.

  • 🔘 Set de date: Fișierul SalesJan2009.csv conține o tranzacție pe linie, țara fiind indicată în a opta coloană separată prin virgulă.
  • ☑️ Cartograf: SalesMapper împarte fiecare linie și emite țara asociată cu valoarea constantă.
  • reductorului: SalesCountryReducer însumează frecvențele sosite pentru fiecare țară într-un singur număr de frecvență.
  • 🧪 Conducător auto: Funcția SalesCountryDriver denumește jobul, declară tipurile de cheie și valoare și conectează ambele clase.
  • 🛠️ Construi: Compilează cu javac, adaugă o intrare în manifestul Main-Class, apoi împachetează totul cu jar cfm.
  • ⚠️ Rulați: Copiați fișierul CSV în HDFS, trimiteți fișierul jar și citiți codul part-00000 din directorul de ieșire.

Exemplu Hadoop MapReduce care creează o primă Java program

În acest tutorial, veți învăța să utilizați Hadoop cu exemple MapReduce. Datele de intrare utilizate sunt VânzăriJan2009.csvConține informații legate de vânzări, cum ar fi numele produsului, prețul, metoda de plată, orașul și țara clientului. Scopul este de a afla numărul de produse vândute în fiecare țară.

Primul program Hadoop MapReduce

Acum în asta Tutorial MapReduce, vom crea primul nostru Java Programul MapReduce:

Captura de ecran de mai jos prezintă datele brute SalesJan2009, unde fiecare linie reprezintă o tranzacție, iar țara se află în a opta coloană separată prin virgulă.

Datele de vânzări SalesJan2009 s-au deschis într-o foaie de calcul care afișează coloanele de tranzacții

Asigurați-vă că aveți Hadoop instalat. Înainte de a începe procesul propriu-zis, schimbați numele utilizatorului în „hduser” (ID-ul folosit la configurarea Hadoop - puteți comuta la ID-ul de utilizator folosit în timpul propriei configurări Hadoop).

su - hduser_

Solicitarea se schimbă în contul hduser, așa cum se arată mai jos.

Solicitare terminal după trecerea la contul hduser

Pasul 1) Creați directorul proiectului și fișierele sursă

Creați un director nou cu numele MapReduceTutorial, așa cum se arată în exemplul MapReduce de mai jos.

sudo mkdir MapReduceTutorial

Comanda mkdir creează directorul MapReduceTutorial

Acordați permisiuni

sudo chmod -R 777 MapReduceTutorial

Comanda chmod care acordă permisiuni complete asupra MapReduceTutorial

Creați cele trei Java fișierele sursă de mai jos, în MapReduceTutorial. Rețineți că toate trei folosesc versiunea mai veche org.apache.hadoop.mapred API, care este încă inclus în 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();
        }
    }
}

Descărcați fișiere aici

Arhiva se extinde în aceleași trei fișiere sursă, așa cum se arată aici.

Extracarhiva ted care listează SalesMapper, SalesCountryReducer și SalesCountryDriver

Verificați permisiunile pentru toate aceste fișiere

Listă lungă care arată permisiunile fișierelor pe cele trei Java fișiere sursă

Dacă lipsesc permisiunile de „citire”, atunci acordați-le:

comanda chmod care adaugă permisiunea de citire la Java fișiere sursă

Pasul 2) Exportați calea de clasă Hadoop

Exportați calea de clasă așa cum se arată în exemplul Hadoop de mai jos.

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

Calea de clasă exportată este readusă la prompt.

Prompt shell după exportarea variabilei Hadoop CLASSPATH

Pasul 3) Compilați Java fișiere

Compilați Java fișiere (aceste fișiere sunt prezente în directorul Final-MapReduceHandsOn). Fișierele lor de clasă vor fi plasate în directorul pachetului.

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

ieșirea javac afișează un avertisment de depreciere în timpul compilării

Acest avertisment poate fi ignorat în siguranță — raportează doar că API-ul mapred este depreciat.

Această compilare va crea un director în directorul curent numit cu numele pachetului specificat în Java fișierul sursă (de exemplu, SalesCountry în cazul nostru) și puneți toate fișierele de clasă compilate în el.

Directorul pachetului SalesCountry care conține cele trei fișiere de clasă compilate

Pasul 4) Creați fișierul manifest

Creați un fișier nou Manifest.txt

sudo gedit Manifest.txt

Adăugați următoarea linie la aceasta:

Main-Class: SalesCountry.SalesCountryDriver

gedit fereastră care afișează intrarea Main-Class în Manifest.txt

SalesCountry.SalesCountryDriver este numele clasei principale. Rețineți că trebuie să apăsați tasta Enter la sfârșitul acestei linii.

Pasul 5) Împachetați clasele într-un borcan

Creați un fișier Jar

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

comanda jar construiește ProductSalePerCountry.jar din manifest

Verificați dacă fișierul jar este creat

Listă în director care confirmă crearea fișierului ProductSalePerCountry.jar

Pasul 6) Porniți Hadoop

Porniți Hadoop

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

Pasul 7) Copiați fișierul de intrare în HDFS

Copiați fișierul SalesJan2009.csv în ~/inputMapReduce

Acum folosiți comanda de mai jos pentru a copia ~/inputMapReduce în HDFS.

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

Ieșirea comenzii copyFromLocal cu un avertisment pentru biblioteca nativă

Putem ignora în siguranță acest avertisment.

Verificați dacă un fișier este copiat sau nu.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls listează SalesJan2009.csv în directorul inputMapReduce

Pasul 8) Executați jobul MapReduce

Rulați jobul MapReduce

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

Ieșirea consolei în timp ce jobul SalePerCountry rulează pe cluster

Aceasta va crea un director de ieșire numit mapreduce_output_sales on HDFS. Conținutul acestui director va fi un fișier care conține vânzările de produse pe țară.

Pasul 9) Citiți rezultatele

Rezultatul poate fi văzut prin interfața de comandă ca fiind,

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

Perechile „țară” și „număr de vânzări” afișate de comanda hdfs dfs -cat

Rezultatele pot fi văzute și printr-o interfață web ca-

Operatii Deschise http://localhost:50070/ într-un browser web. Pe Hadoop 3.x, interfața web NameNode a fost mutată în port 9870, deci folosește http://localhost:9870/ acolo în schimb.

Pagina principală a interfeței web NameNode într-un browser

Acum selectați „Răsfoiți sistemul de fișiere” și navigați la /mapreduce_output_sales

Browser de fișiere HDFS care afișează directorul mapreduce_output_sales

Deschideți piesa-r-00000

Fișierul rezultat part-r-00000 a fost deschis în browserul HDFS

Explicația clasei SalesMapper

Având în vedere că jobul se desfășoară de la un capăt la altul, următoarele trei secțiuni prezintă ce face de fapt fiecare clasă.

În această secțiune, vom înțelege implementarea clasei SalesMapper.

1. Începem prin a specifica un nume de pachet pentru clasa noastră. SalesCountry este numele pachetului nostru. Rețineți că rezultatul compilării, SalesMapper.class, va intra într-un director numit cu acest nume de pachet: SalesCountry.

După aceasta, importăm pachete de bibliotecă.

Mai jos este prezentată o implementare a clasei SalesMapper -

Vizualizare editor a implementării complete a clasei SalesMapper

Eşantion Code Explicaţie:

1. Definiția clasei SalesMapper-

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

Fiecare clasă mapper trebuie să fie extinsă din clasa MapReduceBase și trebuie să implementeze interfața Mapper.

2. Definirea funcției „hartă”-

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

Partea principală a clasei Mapper este metoda „map()” care acceptă patru argumente.

La fiecare apel al metodei „map()”, se transmite o pereche cheie-valoare („key” și „value” în acest cod).

Metoda 'map()' începe prin divizarea textului de intrare primit ca argument. Aceasta împarte fiecare linie în câmpuri.

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

Aici, se folosește „,” ca delimitator.

După aceasta, se formează o pereche folosind o înregistrare la al 7-lea index al matricei „SingleCountryData” și o valoare „1”.

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

Alegem înregistrarea la al 7-lea index deoarece avem nevoie de date despre țară, iar acestea se află la al 7-lea index din matricea „SingleCountryData”.

Vă rugăm să rețineți că datele noastre de intrare sunt în formatul următor (unde Țara se află la al 7-lea indice, cu 0 ca indice de pornire) -

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

O ieșire a funcției mapper este din nou o pereche cheie-valoare emisă folosind metoda „collect()” a funcției „OutputCollector”.

Explicația clasei SalesCountryReducer

În această secțiune, vom înțelege implementarea clasei SalesCountryReducer.

1. Începem prin a specifica numele pachetului pentru clasa noastră. SalesCountry este numele pachetului nostru. Rețineți că rezultatul compilării, SalesCountryReducer.class, va intra într-un director numit cu acest nume de pachet: SalesCountry.

După aceasta, importăm pachete de bibliotecă.

Instantaneul de mai jos prezintă o implementare a clasei SalesCountryReducer-

Vizualizare editor a implementării complete a clasei SalesCountryReducer

Code Explicaţie:

1. Definiția clasei SalesCountryReducer-

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

Aici, primele două tipuri de date, „Text” și „IntWritable”, reprezintă tipul de date al cheii-valoare de intrare a reductorului.

Rezultatul mapperului este sub forma , Această ieșire a mapper-ului devine intrare pentru reducer. Așadar, pentru a se alinia cu tipul său de date, se folosesc aici tipuri de date Text și IntWritable.

Ultimele două tipuri de date, „Text” și „IntWritable”, sunt tipurile de date de ieșire generate de reductor sub forma unei perechi cheie-valoare.

Fiecare clasă Reducer trebuie să fie extinsă din clasa MapReduceBase și trebuie să implementeze interfața Reducer.

2. Definirea funcției „reduce”-

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

O intrare pentru metoda reduce() este o cheie cu o listă de valori multiple.

De exemplu, în cazul nostru, va fi...

, , , , , .

Acest lucru este dat reductorului ca

Deci, pentru a accepta argumente de această formă, se utilizează primele două tipuri de date, și anume Text și Iterator. Textul este un tip de date cheie și iterator. este un tip de date pentru o listă de valori pentru cheia respectivă.

Următorul argument este de tip OutputCollector care colectează ieșirea fazei reducătoare.

Metoda reduce() începe prin copierea valorii cheii și inițializarea numărului de frecvență la 0.

Text key = t_key;
int frequencyForCountry = 0;

Apoi, folosind bucla „while”, iterăm prin lista de valori asociate cheii și calculăm frecvența finală prin însumarea tuturor valorilor.

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

Acum, trimitem rezultatul către colectorul de ieșire sub forma unei chei și a numărului de frecvență obținut.

Codul de mai jos face asta -

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

Explicația clasei SalesCountryDriver

În această secțiune, vom înțelege implementarea clasei SalesCountryDriver.

1. Începem prin a specifica un nume de pachet pentru clasa noastră. SalesCountry este numele pachetului nostru. Rețineți că rezultatul compilării, SalesCountryDriver.class, va intra într-un director numit cu acest nume de pachet: SalesCountry.

Aici este o linie care specifică numele pachetului urmat de cod pentru a importa pachete de bibliotecă.

Declarații de pachete și declarații de import ale SalesCountryDriver

2. Definiți o clasă de driver care va crea un nou loc de muncă client, un obiect de configurare și va face publicitate claselor Mapper și Reducer.

Clasa de șoferi este responsabilă pentru setarea jobului nostru MapReduce pentru a rula HadoopÎn această clasă, specificăm numele jobului, tipul de date de intrare/ieșire și numele claselor mapper și reducer.

Codul driverului care setează numele jobului, clasele cheie și valori și formatele

3. În fragmentul de cod de mai jos, setăm directoarele de intrare și de ieșire care sunt folosite pentru a consuma setul de date de intrare și, respectiv, a produce ieșiri.

arg[0] și arg[1] sunt argumentele din linia de comandă transmise odată cu o comandă dată în MapReduce hands-on, adică,

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

Codul driverului care transmite căile de intrare și ieșire din linia de comandă

4. Declanșează-ne jobul

Codul de mai jos pornește execuția jobului MapReduce-

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

Întrebări frecvente

Toate cele trei clase importă org.apache.hadoop.mapred, API-ul original MapReduce. Hadoop 3.x încă îl livrează și îl rulează, dar noile lucrări sunt scrise în mod normal pentru org.apache.hadoop.mapreduce, care înlocuiește JobConf și JobClient cu Configuration și Job.

Modelele antrenate pe baza istoricului anterior al joburilor prezic timpul de execuție, sugerează dimensiuni de divizare și număr de reductoare și detectează asimetriile înainte de finalizarea unei execuții. De asemenea, acestea semnalează joburile ale căror contoare de înregistrări scurse sau sarcini eșuate depășesc intervalul obișnuit.

Copilot finalizează bine părțile repetitive: semnăturile de clasă, genericele, importurile și apelurile de configurare a driverelor. Logica de business, cum ar fi coloana care conține țara, trebuie totuși verificată de către un dezvoltator în raport cu schema reală.

Hadoop refuză să scrie într-un director de ieșire existent, astfel încât rezultatele finale să nu fie niciodată suprascrise. Ștergeți recursiv mapreduce_output_sales cu hdfs dfs -rm -r sau transmiteți o cale de ieșire diferită la următoarea rulare.

Instrumentul jar ignoră o linie finală de manifest care nu are un terminator de linie, așadar Main-Class este eliminată silențios. Jarul se construiește apoi fără erori, dar eșuează la momentul execuției deoarece nu este înregistrată nicio clasă principală.

Da. Exemplul introduce 2.2.0 în fiecare nume de fișier jar. Într-o altă versiune, înlocuiți acea versiune sau pur și simplu utilizați ieșirea din hadoop classpath, care afișează fiecare fișier jar de care are nevoie distribuția instalată.

Nu. O instalare pseudo-distribuită cu un singur nod o rulează neschimbată, deoarece comenzile presupun doar că HDFS și YARN sunt pornite. Același fișier jar se trimite către un cluster real cu mai multe noduri fără nicio modificare de cod.

Schimbați indexul matricei în mapper de la 7 la 5, deoarece City este a șasea coloană a fișierului SalesJan2009.csv. Recompilați, reconstruiți fișierul jar și rulați jobul într-un director de ieșire nou.

Rezumați această postare cu: