Esempio di Hadoop MapReduce: Primo Java Programma con Code

โšก Riepilogo intelligente

I programmi Hadoop MapReduce sono scritti come tre Java classi, un mapper, un reducer e un driver, che vengono compilati, impacchettati in un file jar e inviati al cluster per contare le vendite per paese.

  • ๐Ÿ”˜ dataset: Il file SalesJan2009.csv contiene una transazione per riga, con il paese indicato nell'ottava colonna, separato da virgole.
  • โ˜‘๏ธ Mappatore: SalesMapper suddivide ogni riga ed emette il paese associato al valore costante uno.
  • โœ… Reducer: SalesCountryReducer somma gli ordini provenienti da ciascun paese in un unico conteggio di frequenza.
  • ๐Ÿงช Autista: SalesCountryDriver assegna un nome al lavoro, dichiara i tipi di chiave e valore e collega tra loro entrambe le classi.
  • ๏ธ Corporatura: Compila con javac, aggiungi una voce al manifest della classe principale, quindi impacchetta tutto con jar cfm.
  • โš ๏ธ Esegui: Copia il file CSV in HDFS, invia il file jar e leggi la parte 00000 dalla directory di output.

Esempio di Hadoop MapReduce che crea un primo Java Programma

In questo tutorial imparerai a utilizzare Hadoop con gli esempi MapReduce. I dati di input utilizzati sono VenditeJan2009.csvContiene informazioni relative alle vendite, come il nome del prodotto, il prezzo, la modalitร  di pagamento, la cittร  e il paese del cliente. L'obiettivo รจ scoprire il numero di prodotti venduti in ciascun paese.

Primo programma Hadoop MapReduce

Ora in questo Tutorial su MapReduce, creeremo il nostro primo Java Programma MapReduce:

Lo screenshot qui sotto mostra i dati grezzi delle vendite di gennaio 2009, dove ogni riga rappresenta una transazione e il paese si trova nell'ottava colonna, separata da virgole.

Dati di vendita di gennaio 2009 aperti in un foglio di calcolo che mostra le colonne delle transazioni

Assicurati di aver installato Hadoop. Prima di iniziare la procedura vera e propria, cambia utente in 'hduser' (l'ID utilizzato durante la configurazione di Hadoop; puoi passare all'ID utente che usi durante la tua configurazione di Hadoop).

su - hduser_

Il prompt cambia e passa all'account hduser, come mostrato di seguito.

Richiesta del terminale dopo il passaggio all'account hduser

Passaggio 1) Creare la directory del progetto e i file sorgente

Crea una nuova directory con il nome MapReduceTutorial come mostrato nell'esempio MapReduce qui sotto.

sudo mkdir MapReduceTutorial

Comando mkdir per la creazione della directory MapReduceTutorial

Concedi i permessi

sudo chmod -R 777 MapReduceTutorial

Comando chmod che concede i permessi completi su MapReduceTutorial

Crea i tre Java file sorgente sotto all'interno di MapReduceTutorial. Nota che tutti e tre utilizzano il vecchio org.apache.hadoop.mapred API, che รจ ancora inclusa in 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();
        }
    }
}

Scarica i file qui

L'archivio si espande negli stessi tre file sorgente, come mostrato qui.

ExtracElenco degli archivi TED: SalesMapper, SalesCountryReducer e SalesCountryDriver

Controlla i permessi di tutti questi file

Elenco esteso che mostra i permessi dei file sui tre Java file sorgenti

Se mancano i permessi di lettura, concedeteli:

comando chmod che aggiunge il permesso di lettura al Java file sorgenti

Passaggio 2) Esportare il classpath di Hadoop

Esporta il classpath come mostrato nell'esempio Hadoop qui sotto.

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

Il classpath esportato viene visualizzato nuovamente nel prompt.

Prompt della shell dopo l'esportazione della variabile CLASSPATH di Hadoop

Passaggio 3) Compilare il Java file

Compila il Java file (questi file sono presenti nella directory Final-MapReduceHandsOn). I relativi file di classe verranno inseriti nella directory del pacchetto.

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

L'output di javac mostra un avviso di deprecazione durante la compilazione

Questo avviso puรฒ essere tranquillamente ignorato: indica solo che l'API di MapRed รจ obsoleta.

Questa compilazione creerร  una directory nella directory corrente denominata con il nome del pacchetto specificato in Java Inserire al suo interno il file sorgente (ad esempio SalesCountry nel nostro caso) e tutti i file di classe compilati.

Directory del pacchetto SalesCountry contenente i tre file di classe compilati

Passaggio 4) Creare il file manifest

Crea un nuovo file denominato Manifest.txt

sudo gedit Manifest.txt

Aggiungi la seguente riga:

Main-Class: SalesCountry.SalesCountryDriver

gedit Finestra che mostra la voce Main-Class nel file Manifest.txt

SalesCountry.SalesCountryDriver รจ il nome della classe principale. Nota che devi premere il tasto Invio alla fine di questa riga.

Passaggio 5) Impacchettare le classi in un barattolo

Crea un file Jar

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

Comando jar per la creazione di ProductSalePerCountry.jar dal manifesto

Verifica che il file jar sia stato creato

Elenco della directory che conferma la creazione di ProductSalePerCountry.jar

Passaggio 6) Avviare Hadoop

Avvia Hadoop

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

Passaggio 7) Copiare il file di input in HDFS

Copia il file SalesJan2009.csv nella directory ~/inputMapReduce

Ora usa il comando seguente per copiare ~/inputMapReduce in HDFS.

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

Output del comando copyFromLocal con un avviso relativo alla libreria nativa.

Possiamo tranquillamente ignorare questo avvertimento.

Verificare se un file รจ effettivamente copiato o meno.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls elenca SalesJan2009.csv all'interno della directory inputMapReduce

Passaggio 8) Eseguire il job MapReduce

Esegui il lavoro MapReduce

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

Output della console durante l'esecuzione del job SalePerCountry sul cluster

Ciรฒ creerร  una directory di output denominata mapreduce_output_sales su HDFS. Il contenuto di questa directory sarร  un file contenente le vendite di prodotti per paese.

Passaggio 9) Leggere i risultati

Il risultato puรฒ essere visualizzato tramite l'interfaccia di comando come segue:

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

Coppie Paese e numero di vendite stampate dal comando hdfs dfs -cat

I risultati possono anche essere visualizzati tramite un'interfaccia web come-

Apri http://localhost:50070/ in un browser web. Su Hadoop 3.x l'interfaccia utente web di NameNode รจ stata spostata sulla porta 9870, quindi usa http://localhost:9870/ lรฌ invece.

Pagina iniziale dell'interfaccia web di NameNode in un browser

Ora seleziona 'Esplora il filesystem' e vai a /mapreduce_output_sales

Esplora file HDFS che mostra la directory mapreduce_output_sales

Aprire la parte r-00000

File dei risultati part-r-00000 aperto nel browser HDFS

Spiegazione della classe SalesMapper

Dopo aver descritto il lavoro dall'inizio alla fine, le prossime tre sezioni illustrano nel dettaglio le funzioni di ciascuna classe.

In questa sezione, analizzeremo l'implementazione della classe SalesMapper.

1. Iniziamo specificando il nome del pacchetto per la nostra classe. SalesCountry รจ il nome del nostro pacchetto. Si noti che l'output della compilazione, SalesMapper.class, verrร  salvato in una directory con lo stesso nome del pacchetto: SalesCountry.

Successivamente importiamo i pacchetti di librerie.

L'immagine sottostante mostra un'implementazione della classe SalesMapper.

Visualizzazione dell'editor dell'implementazione completa della classe SalesMapper

Campione Code Spiegazione:

1. Definizione della classe SalesMapper-

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

Ogni classe mapper deve estendere la classe MapReduceBase e deve implementare l'interfaccia Mapper.

2. Definizione della funzione 'mappa'-

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

La parte principale della classe Mapper รจ un metodo 'map()' che accetta quattro argomenti.

Ad ogni chiamata al metodo 'map()', viene passata una coppia chiave-valore ('key' e 'value' in questo codice).

Il metodo 'map()' inizia suddividendo il testo di input ricevuto come argomento. Suddivide ogni riga in campi.

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

Qui, la virgola (,) viene utilizzata come delimitatore.

Dopodichรฉ, viene formata una coppia utilizzando un record all'indice 7 dell'array 'SingleCountryData' e un valore '1'.

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

Selezioniamo il record all'indice 7 perchรฉ abbiamo bisogno dei dati relativi al Paese, che si trovano all'indice 7 nell'array 'SingleCountryData'.

Si prega di notare che i nostri dati di input sono nel formato seguente (dove Paese si trova all'indice 7, con 0 come indice iniziale)-

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

L'output del mapper รจ nuovamente una coppia chiave-valore che viene emessa utilizzando il metodo 'collect()' di 'OutputCollector'.

Spiegazione della classe SalesCountryReducer

In questa sezione, analizzeremo l'implementazione della classe SalesCountryReducer.

1. Iniziamo specificando il nome del pacchetto per la nostra classe. SalesCountry รจ il nome del nostro pacchetto. Si noti che l'output della compilazione, SalesCountryReducer.class, verrร  salvato in una directory con lo stesso nome del pacchetto: SalesCountry.

Successivamente importiamo i pacchetti di librerie.

L'immagine seguente mostra un'implementazione della classe SalesCountryReducer.

Visualizzazione dell'editor dell'implementazione completa della classe SalesCountryReducer

Code Spiegazione:

1. Definizione della classe SalesCountryReducer-

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

Qui, i primi due tipi di dati, 'Text' e 'IntWritable', rappresentano il tipo di dati della coppia chiave-valore in ingresso al reducer.

L'output del mapper รจ nella forma di , L'output del mapper diventa l'input del reducer. Pertanto, per allinearsi al suo tipo di dati, vengono utilizzati Text e IntWritable come tipi di dati.

Gli ultimi due tipi di dati, 'Text' e 'IntWritable', rappresentano il tipo di dati di output generato dal reducer sotto forma di coppia chiave-valore.

Ogni classe reducer deve estendere la classe MapReduceBase e deve implementare l'interfaccia Reducer.

2. Definizione della funzione "riduci"

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

L'input del metodo reduce() รจ una chiave contenente un elenco di piรน valori.

Ad esempio, nel nostro caso, sarร -

, , , , , .

Questo viene dato al riduttore come

Pertanto, per accettare argomenti di questa forma, vengono utilizzati innanzitutto due tipi di dati, ovvero Testo e Iteratore. . Il testo รจ un tipo di dati di chiave e iteratore รจ un tipo di dati per un elenco di valori per quella chiave.

L'argomento successivo รจ di tipo OutputCollector che raccoglie l'uscita della fase di riduzione.

Il metodo reduce() inizia copiando il valore della chiave e inizializzando il conteggio della frequenza a 0.

Text key = t_key;
int frequencyForCountry = 0;

Successivamente, utilizzando un ciclo 'while', scorriamo l'elenco dei valori associati alla chiave e calcoliamo la frequenza finale sommando tutti i valori.

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

Ora inviamo il risultato al collettore di output sotto forma di chiave e conteggio della frequenza ottenuta.

Il codice seguente fa questo-

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

Spiegazione della classe SalesCountryDriver

In questa sezione, analizzeremo l'implementazione della classe SalesCountryDriver.

1. Iniziamo specificando il nome del pacchetto per la nostra classe. SalesCountry รจ il nome del nostro pacchetto. Si noti che l'output della compilazione, SalesCountryDriver.class, verrร  salvato in una directory con lo stesso nome del pacchetto: SalesCountry.

Ecco una riga che specifica il nome del pacchetto seguito dal codice per importare i pacchetti della libreria.

Dichiarazione del pacchetto e istruzioni di importazione di SalesCountryDriver

2. Definire una classe driver che creerร  un nuovo lavoro client, un oggetto di configurazione e pubblicizzerร  le classi Mapper e Reducer.

La classe driver รจ responsabile dell'impostazione del nostro lavoro MapReduce per l'esecuzione HadoopIn questa classe, specifichiamo il nome del job, il tipo di dati di input/output e i nomi delle classi mapper e reducer.

Il codice del driver imposta il nome del lavoro, le classi chiave e valore e i formati

3. Nello snippet di codice seguente, impostiamo le directory di input e output che vengono utilizzate rispettivamente per consumare il set di dati di input e produrre output.

arg[0] e arg[1] sono gli argomenti della riga di comando passati con un comando dato in MapReduce hands-on, vale a dire,

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

Codice del driver che passa i percorsi di input e output dalla riga di comando

4. Attiva il nostro lavoro

Il codice seguente avvia l'esecuzione del job MapReduce.

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

DOMANDE FREQUENTI

Tutte e tre le classi importano org.apache.hadoop.mapred, l'API MapReduce originale. Hadoop 3.x la include e la esegue ancora, ma il nuovo codice viene solitamente scritto utilizzando org.apache.hadoop.mapreduce, che sostituisce JobConf e JobClient con Configuration e Job.

I modelli addestrati sulla cronologia dei lavori precedenti prevedono i tempi di esecuzione, suggeriscono le dimensioni degli split e il numero di reducer, e individuano le anomalie prima che un'esecuzione termini. Segnalano inoltre i lavori i cui contatori di record persi o di attivitร  fallite si discostano dall'intervallo usuale.

Copilot gestisce bene le parti ripetitive: firme delle classi, generici, importazioni e chiamate di configurazione del driver. La logica di business, come ad esempio quale colonna contiene il paese, deve comunque essere verificata dallo sviluppatore rispetto allo schema reale.

Hadoop si rifiuta di scrivere in una directory di output esistente, quindi i risultati finali non vengono mai sovrascritti. Elimina la directory mapreduce_output_sales in modo ricorsivo con il comando hdfs dfs -rm -r oppure specifica un percorso di output diverso alla successiva esecuzione.

Lo strumento jar ignora una riga finale del manifest priva di terminatore di riga, quindi Main-Class viene silenziosamente omessa. Il jar viene quindi compilato senza errori, ma fallisce in fase di esecuzione perchรฉ non viene registrata alcuna classe principale.

Sรฌ. L'esempio specifica la versione 2.2.0 in ogni nome di file jar. In un'altra versione, sostituiscila con quella, oppure usa semplicemente l'output di hadoop classpath, che elenca tutti i file jar necessari alla distribuzione installata.

No. Un'installazione pseudo-distribuita a nodo singolo lo esegue senza modifiche, perchรฉ i comandi presuppongono solo che HDFS e YARN siano avviati. Lo stesso file jar viene inviato a un cluster multi-nodo reale senza alcuna modifica al codice.

Modifica l'indice dell'array nel mapper da 7 a 5, perchรฉ City รจ la sesta colonna di SalesJan2009.csv. Ricompila, ricrea il jar ed esegui il job su una nuova directory di output.

Riassumi questo post con: