Hadoop MapReduce Join e contatore con l'esempio

โšก Riepilogo intelligente

Le join di MapReduce combinano due grandi dataset in base a una chiave condivisa, all'interno del mapper o del reducer, mentre i contatori di MapReduce raccolgono statistiche sul job in modo che i record errati possano essere misurati anzichรฉ ipotizzati.

  • ๐Ÿ”˜ Nozioni di base per l'adesione: Il set di dati piรน piccolo dei due viene distribuito a ogni nodo dati e utilizzato come database di riferimento.
  • โ˜‘๏ธ Unione lato mappa: Richiede che ciascun input sia partizionato, suddiviso in parti uguali e ordinato in base alla chiave di unione prima dell'esecuzione della funzione map.
  • โœ… Unione lato riduzione: Non รจ necessario alcun partizionamento, perchรฉ ogni tupla che condivide una chiave di join finisce nello stesso reducer.
  • ๐Ÿงช Esempio pratico: I file DeptName.txt e DeptStrength.txt vengono copiati in HDFS e uniti tramite Dept_ID da un file jar precompilato.
  • ๏ธ Tipologie di banconi: Cinque gruppi di contatori integrati vengono forniti con ogni lavoro e i contatori definiti dall'utente vengono dichiarati come un Java enum.
  • โš ๏ธ Utilizzo al banco: L'incremento di un contatore per ogni record mancante o non valido trasforma i problemi di qualitร  dei dati in un numero da inserire nel report di lavoro.

Tutorial su join e contatori in Hadoop MapReduce con un esempio pratico

Che cos'รจ un "Join" in MapReduce?

L'operazione di join di MapReduce viene utilizzata per combinare due grandi insiemi di dati. Tuttavia, questo processo richiede la scrittura di molto codice per eseguire l'operazione di join vera e propria. L'unione di due insiemi di dati inizia confrontando le dimensioni di ciascun insieme. Se un insieme di dati รจ piรน piccolo rispetto all'altro, l'insieme di dati piรน piccolo viene distribuito a tutti i nodi dati del cluster.

Una volta che ti unisci MapReduce Nel caso di un dataset distribuito, il Mapper o il Reducer utilizzano il dataset piรน piccolo per cercare i record corrispondenti nel dataset piรน grande e quindi combinano tali record per formare i record di output.

Tipi di unione

A seconda del punto in cui viene effettivamente eseguita l'unione, le unioni in Hadoop si classificano in due tipi.

  1. Unione lato mappa โ€” Quando l'unione viene eseguita dal mapper, si parla di unione lato mappa. In questo tipo, l'unione viene eseguita prima che i dati vengano effettivamente utilizzati dalla funzione map. รˆ obbligatorio che l'input per ogni mappa sia in forma di partizione e sia ordinato. Inoltre, deve esserci un numero uguale di partizioni e queste devono essere ordinate in base alla chiave di unione.
  2. giunzione lato ridotto โ€” Quando l'unione viene eseguita dal reducer, si parla di join lato reducer. In questo tipo di join non รจ necessario che il dataset sia strutturato (o partizionato). Qui, l'elaborazione lato map genera la chiave di join e le tuple corrispondenti di entrambe le tabelle. Come conseguenza di questa elaborazione, tutte le tuple con la stessa chiave di join vengono inviate allo stesso reducer, che unisce i record con la stessa chiave.

Un flusso complessivo del processo di join in Hadoop รจ illustrato nel diagramma seguente.

Diagramma di flusso del processo che confronta un join lato mappa con un join lato riduzione in Hadoop
Tipi di join in Hadoop MapReduce

Chiarite le due varianti, la sezione successiva illustra un join lato reduce su due piccoli file di reparto.

Come unire due set di dati: esempio di MapReduce

Sono presenti due set di dati in due file diversi (mostrati di seguito). La chiave Dept_ID รจ comune a entrambi i file. L'obiettivo รจ utilizzare MapReduce Join per combinare questi file.

Primo file di input che elenca gli ID dei dipartimenti insieme ai nomi dei dipartimenti

file 1
Secondo file di input che elenca gli ID dei dipartimenti insieme ai valori di forza dei dipartimenti

file 2

Ingresso: Il set di dati di input รจ costituito da un file di testo, denominato DeptName.txt e DeptStrength.txt.

Scarica i file di input da qui

Assicurati di averlo Hadoop installato. Prima di iniziare con il processo effettivo dell'esempio MapReduce Join, cambia l'utente in 'hduser' (l'ID utilizzato durante la configurazione di Hadoop; puoi passare all'ID utente utilizzato durante la configurazione di Hadoop).

su - hduser_

Il prompt cambia e punta all'account Hadoop, come mostrato di seguito.

Terminale dopo essere passato all'account hduser con il comando su

Passo 1) Copia il file zip nella posizione che preferisci

L'archivio MapReduceJoin scaricato รจ stato posizionato nella directory di lavoro selezionata.

Passo 2) Decomprimere il file zip

sudo tar -xvf MapReduceJoin.tar.gz

L'extracI nomi dei file scorrono mentre tar decomprime l'archivio.

Elenco della console dei file extracted da MapReduceJoin.tar.gz

Passo 3) Vai alla directory MapReduceJoin/

cd MapReduceJoin/

Prompt della shell dopo essersi spostati nella directory MapReduceJoin

Passo 4) Avvia Hadoop

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

Entrambi gli script stampano i daemon che avviano.

Messaggi di avvio dagli script del demone HDFS e YARN

Passo 5) DeptStrength.txt e DeptName.txt sono i file di input utilizzati per questo programma di esempio MapReduce Join.

Questi file devono essere copiati in HDFS utilizzando il comando seguente-

$HADOOP_HOME/bin/hdfs dfs -copyFromLocal DeptStrength.txt DeptName.txt /

Entrambi i file di testo di input sono stati copiati nella directory radice HDFS.

Passo 6) Esegui il programma utilizzando il comando seguente-

$HADOOP_HOME/bin/hadoop jar MapReduceJoin.jar MapReduceJoin/JoinDriver/DeptStrength.txt /DeptName.txt /output_mapreducejoin

Il comando viene prima ripetuto, e successivamente il processo segnala il suo stato di avanzamento sulla console.

Avvio da riga di comando del file jar MapReduceJoin incluso nel pacchetto.

Uscita console traccontrollare l'avanzamento del job di unione MapReduce

Passo 7) Dopo l'esecuzione, il file di output (denominato 'part-00000') verrร  memorizzato nella directory /output_mapreducejoin su HDFS

I risultati possono essere visualizzati utilizzando l'interfaccia della riga di comando

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

Registri del dipartimento uniti stampati da HDFS tramite il comando cat

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

Pagina di destinazione dell'interfaccia web di Hadoop utilizzata per accedere al browser del file system.

Ora seleziona 'Esplora il filesystem' e naviga fino a /output_mapreducejoin

Esplorazione della vista del file system HDFS nella directory output_mapreducejoin

Aprire la parte r-00000

Selezione del file di output part-r-00000 all'interno della visualizzazione del browser

Vengono visualizzati i risultati

Nome del dipartimento di appartenenza e numero di righe del dipartimento visualizzate nel browser

NOTA: Tieni presente che prima di eseguire questo programma per la prossima volta, dovrai eliminare la directory di output /output_mapreducejoin

$HADOOP_HOME/bin/hdfs dfs -rm -r /output_mapreducejoin

L'alternativa รจ utilizzare un nome diverso per la directory di output.

Le join ti mostrano l'aspetto dei dati. I contatori, che tratteremo in seguito, ti mostrano come si รจ comportato il processo che li ha generati.

Cos'รจ il contatore in MapReduce?

Un contatore in MapReduce รจ un meccanismo utilizzato per raccogliere e misurare informazioni statistiche sui job e sugli eventi di MapReduce. I contatori mantengono il track di varie statistiche di lavoro in MapReduce come il numero di operazioni eseguite e lo stato di avanzamento dell'operazione. I contatori vengono utilizzati per la diagnosi dei problemi in MapReduce.

I contatori Hadoop sono simili all'inserimento di un messaggio di registro nel codice di una mappa o alla riduzione. Queste informazioni potrebbero essere utili per la diagnosi di un problema nell'elaborazione del lavoro MapReduce.

In genere, questi contatori in Hadoop sono definiti in un programma (map o reduce) e vengono incrementati durante l'esecuzione quando si verifica un particolare evento o condizione (specifica per quel contatore). Un'ottima applicazione dei contatori Hadoop รจ track record validi e non validi da un dataset di input.

Tipi di contatori MapReduce

Esistono fondamentalmente 2 tipi di contatori MapReduce

  1. Contatori integrati di Hadoop: Esistono alcuni contatori Hadoop integrati per lavoro. Di seguito sono riportati i gruppi di contatori integrati:
    • Contatori attivitร  MapReduce โ€” Raccoglie informazioni specifiche sull'attivitร  (ad esempio, il numero di record di input) durante la sua esecuzione.
    • Contatori del file system โ€” Raccoglie informazioni come il numero di byte letti o scritti da un'attivitร .
    • Contatori FileInputFormat โ€” Raccoglie le informazioni relative a un certo numero di byte letti tramite FileInputFormat.
    • Contatori FileOutputFormat โ€” Raccoglie informazioni su un certo numero di byte scritti tramite FileOutputFormat.
    • Contatori di posti di lavoro โ€” Questi contatori registrano statistiche relative all'intero lavoro, come ad esempio il numero di attivitร  avviate per un lavoro.
  2. Contatori definiti dall'utente: Oltre ai contatori integrati, un utente puรฒ definire i propri contatori utilizzando funzionalitร  simili fornite dai linguaggi di programmazione. Ad esempio, in Java, un 'enum' viene utilizzato per definire contatori definiti dall'utente.

๐Ÿ’ก Nota sulla versione: I contatori dei posti di lavoro erano gestiti dal JobTracker sotto MRv1. Su YARN tale ruolo appartiene all'ApplicationMaster di MapReduce, quindi i nomi dei contatori rimangono, ma il componente che li segnala รจ cambiato.

Un lavoro non puรฒ dichiarare un numero illimitato di contatori. mapreduce.job.counters.max L'impostazione limita il totale per lavoro a 120 per impostazione predefinita e un lavoro che dichiara di piรน fallisce con un LimitExceededException, quindi i contatori sono pensati per una manciata di segnali aggregati piuttosto che per conteggi per singola chiave.

Esempio di contatori

Esempio di MapClass con contatori per contare il numero di valori mancanti e non validi. File di dati di input utilizzato in questo tutorial: Il nostro set di dati di input รจ un file CSV, SalesJan2009.csv

public static class MapClass
            extends MapReduceBase
            implements Mapper<LongWritable, Text, Text, Text>
{
    static enum SalesCounters { MISSING, INVALID };
    public void map ( LongWritable key, Text value,
                 OutputCollector<Text, Text> output,
                 Reporter reporter) throws IOException
    {
        
        //Input string is split using ',' and stored in 'fields' array
        String fields[] = value.toString().split(",", -20);
        //Value at 4th index is country. It is stored in 'country' variable
        String country = fields[4];
        
        //Value at 8th index is sales data. It is stored in 'sales' variable
        String sales = fields[8];
      
        if (country.length() == 0) {
            reporter.incrCounter(SalesCounters.MISSING, 1);
        } else if (sales.startsWith("\"")) {
            reporter.incrCounter(SalesCounters.INVALID, 1);
        } else {
            output.collect(new Text(country), new Text(sales + ",1"));
        }
    }
}

Il frammento di codice qui sopra mostra un esempio di implementazione dei contatori in Hadoop MapReduce.

Qui, Contatori di vendita รจ un contatore definito utilizzando 'enumViene utilizzato per contare i record di input MANCANTI e NON VALIDI.

Nel frammento di codice, se 'nazioneSe il campo ha lunghezza zero, il suo valore รจ mancante e quindi il contatore corrispondente SalesCounters.MISSING viene incrementato.

Successivamente, se 'venditeSe il campo inizia con ", il record viene considerato NON VALIDO. Ciรฒ รจ indicato dall'incremento del contatore SalesCounters.INVALID.

๐Ÿ’ก Nota sull'API: il frammento sopra utilizza l'originale org.apache.hadoop.mapred API, dove MapReduceBase, il Mapper interfaccia, OutputCollector and Reporter appaiono separatamente. Il codice attuale รจ scritto contro org.apache.hadoop.mapreduce, dove un singolo Context sostituisce il collettore e il reporter e un contatore viene incrementato con context.getCounter(SalesCounters.MISSING).increment(1)Il concetto di contatore รจ identico in entrambi.

DOMANDE FREQUENTI

Scegli la variante mapper quando uno dei lati รจ sufficientemente piccolo da poter essere memorizzato in memoria su ogni nodo, poichรฉ in questo modo si evita completamente la fase di shuffle. Scegli la variante reducer quando entrambi i lati sono grandi o non ordinati e accetta il costo aggiuntivo della rete.

I modelli apprendono dalla cronologia dei lavori per prevedere i tempi di esecuzione, consigliare le dimensioni degli split e il numero di reducer e rilevare le anomalie nei valori dei contatori. Segnalano inoltre i lavori i cui contatori di record persi o di attivitร  fallite si discostano dall'intervallo normale per quella pipeline.

Copilot genera scheletri plausibili per mapper e reducer, ma mescola liberamente il vecchio pacchetto mapred con il piรน recente pacchetto mapreduce in un'unica classe, il che non consente la compilazione. Correggi le importazioni e le firme dei metodi prima di fidarti della logica.

Si tratta del meccanismo che invia il file piรน piccolo a ogni nodo prima dell'avvio delle attivitร . Ogni mapper carica quindi quella copia in una mappa hash e cerca localmente le corrispondenze, il che rende possibile un join lato mapper.

Vengono stampati nel riepilogo della console al termine del lavoro, visualizzati nelle pagine web della cronologia dei lavori e del gestore delle risorse e leggibili a livello di programmazione dall'oggetto lavoro, in modo che un driver possa verificarli e interrompere un'esecuzione non riuscita.

Un singolo oggetto Context. Esso gestisce il lavoro che OutputCollector e Reporter si dividevano tra loro, quindi l'output viene scritto e i contatori vengono incrementati tramite lo stesso handle passato al metodo map.

Per la maggior parte dei lavori giornalistici, no. Unione HiveQL Si riduce allo stesso schema di rimescolamento e unione in poche righe. I job scritti manualmente valgono la pena quando la logica di unione non si adatta a una clausola SQL.

Hadoop si rifiuta di scrivere in una directory di output giร  esistente, proteggendo cosรฌ i risultati finali dalla sovrascrittura. Elimina prima la directory in modo ricorsivo oppure specifica un percorso di output diverso alla successiva esecuzione.

Riassumi questo post con: