Primjer Hadoop MapReduce: Prvi Java Program sa Code

โšก Pametni saลพetak

Hadoop MapReduce programi su napisani kao tri Java klase, mapper, reducer i driver, koji se kompajliraju, pakiraju u jar datoteku i ลกalju klasteru za brojanje prodaje po zemlji.

  • ๐Ÿ”˜ Skup podataka: Datoteka SalesJan2009.csv sadrลพi jednu transakciju po retku, a drลพava je u osmom stupcu odvojenom zarezom.
  • โ˜‘๏ธ Maper: SalesMapper dijeli svaki redak i emitira drลพavu uparenu s konstantnom vrijednoลกฤ‡u jedan.
  • โœ… Redukcija: SalesCountryReducer zbraja one koji stiลพu za svaku zemlju u jedan broj frekvencija.
  • ๐Ÿงช Vozaฤ: SalesCountryDriver imenuje posao, deklarira tipove kljuฤeva i vrijednosti te povezuje obje klase.
  • ๐Ÿ› ๏ธ Graฤ‘a: Kompajliraj s javac-om, dodaj unos u manifest glavne klase, a zatim sve zapakiraj s jar cfm-om.
  • โš ๏ธ Trฤanje: Kopirajte CSV u HDFS, poลกaljite JAR datoteku i proฤitajte dio-00000 iz izlaznog direktorija.

Primjer kreiranja prvog Hadoop MapReduce-a Java program

U ovom vodiฤu nauฤit ฤ‡ete koristiti Hadoop s MapReduce primjerima. Koriลกteni ulazni podaci su ProdajaSijeฤanj2009.csvSadrลพi podatke vezane uz prodaju kao ลกto su naziv proizvoda, cijena, naฤin plaฤ‡anja, grad i drลพava klijenta. Cilj je saznati broj prodanih proizvoda u svakoj drลพavi.

Prvi Hadoop MapReduce program

Sada u ovome Vodiฤ za MapReduce, stvorit ฤ‡emo svoj prvi Java Program MapReduce:

Snimka zaslona u nastavku prikazuje neobraฤ‘ene podatke o prodaji za sijeฤanj 2009., gdje svaki redak predstavlja jednu transakciju, a drลพava se nalazi u osmom stupcu odvojenom zarezom.

Podaci o prodaji za sijeฤanj 2009. otvoreni su u proraฤunskoj tablici koja prikazuje stupce transakcija

Provjerite imate li instaliran Hadoop. Prije nego ลกto zapoฤnete sa stvarnim postupkom, promijenite korisnika u 'hduser' (id koji se koristio prilikom konfiguriranja Hadoop-a - moลพete se prebaciti na korisniฤki ID koji se koristio tijekom vlastite konfiguracije Hadoop-a).

su - hduser_

Uputa se mijenja u hduser raฤun, kao ลกto je prikazano dolje.

Terminalski upit nakon prebacivanja na hduser raฤun

Korak 1) Izradite direktorij projekta i izvorne datoteke

Izradite novi direktorij s nazivom MapReduceTutorial kao ลกto je prikazano u donjem primjeru MapReduce.

sudo mkdir MapReduceTutorial

Naredba mkdir stvara direktorij MapReduceTutorial

Dajte dozvole

sudo chmod -R 777 MapReduceTutorial

chmod naredba koja daje puna dopuลกtenja na MapRedutoryju

Stvorite tri Java izvorne datoteke u nastavku unutar MapReduceTutoriala. Imajte na umu da sva tri koriste stariju org.apache.hadoop.mapred API, koji se joลก uvijek isporuฤuje s Hadoopom 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();
        }
    }
}

Preuzmite datoteke ovdje

Arhiva se proลกiruje u iste tri izvorne datoteke, kao ลกto je ovdje prikazano.

Extracpopis ted arhive SalesMapper, SalesCountryReducer i SalesCountryDriver

Provjerite dopuลกtenja za sve te datoteke

Dugi popis koji prikazuje dozvole za datoteke na tri Java izvorne datoteke

Ako nedostaju dozvole za 'ฤitanje', dodijelite ih:

chmod naredba dodaje dopuลกtenje za ฤitanje Java izvorne datoteke

Korak 2) Izvoz Hadoop classpath-a

Izvezite putanju klase kao ลกto je prikazano u donjem Hadoop primjeru.

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

Izvezena putanja klase se vraฤ‡a u promptu.

Shell prompt nakon izvoza Hadoop CLASSPATH varijable

Korak 3) Sastavite Java slika

Sastavite Java datoteke (ove se datoteke nalaze u direktoriju Final-MapReduceHandsOn). Njihove datoteke klasa bit ฤ‡e smjeลกtene u direktorij paketa.

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

javac izlaz koji prikazuje upozorenje o zastarivanju tijekom kompilacije

Ovo upozorenje se moลพe sigurno zanemariti โ€” ono samo izvjeลกtava da je mapirani API zastario.

Ova kompilacija ฤ‡e stvoriti direktorij u trenutnom direktoriju nazvan s nazivom paketa navedenim u Java izvornu datoteku (tj. SalesCountry u naลกem sluฤaju) i u nju staviti sve kompilirane datoteke klase.

Direktorij paketa SalesCountry koji sadrลพi tri kompilirane datoteke klase

Korak 4) Izradite datoteku manifesta

Stvorite novu datoteku Manifest.txt

sudo gedit Manifest.txt

Dodajte mu sljedeฤ‡i redak:

Main-Class: SalesCountry.SalesCountryDriver

gedit prozor koji prikazuje unos Main-Class u Manifest.txt datoteci

SalesCountry.SalesCountryDriver je naziv glavne klase. Imajte na umu da morate pritisnuti tipku Enter na kraju ovog retka.

Korak 5) Zapakirajte klase u staklenku

Stvorite Jar datoteku

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

jar naredba izgradnja ProductSalePerCountry.jar iz manifesta

Provjerite je li jar datoteka stvorena

Popis direktorija koji potvrฤ‘uje da je ProductSalePerCountry.jar kreiran

Korak 6) Pokrenite Hadoop

Pokrenite Hadoop

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

Korak 7) Kopirajte ulaznu datoteku u HDFS

Kopirajte datoteku SalesJan2009.csv u ~/inputMapReduce

Sada upotrijebite naredbu u nastavku za kopiranje ~/inputMapReduce u HDFS.

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

Izlaz naredbe copyFromLocal s upozorenjem izvorne biblioteke

Moลพemo slobodno zanemariti ovo upozorenje.

Provjerite je li datoteka stvarno kopirana ili ne.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls ispisuje SalesJan2009.csv unutar direktorija inputMapReduce

Korak 8) Pokrenite zadatak MapReduce

Pokrenite zadatak MapReduce

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

Izlaz konzole dok se zadatak SalePerCountry izvrลกava na klasteru

Ovo ฤ‡e stvoriti izlazni direktorij pod nazivom mapreduce_output_sales on HDFS. Sadrลพaj ovog imenika bit ฤ‡e datoteka koja ฤ‡e sadrลพavati prodaju proizvoda po zemlji.

Korak 9) Proฤitajte rezultate

Rezultat se moลพe vidjeti kroz suฤelje naredbi kao,

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

Parovi drลพava i broj prodaje ispisani naredbom hdfs dfs -cat

Rezultati se takoฤ‘er mogu vidjeti putem web suฤelja kao-

Otvoren http://localhost:50070/ u web pregledniku. Na Hadoopu 3.x, web suฤelje NameNodea premjeลกteno je na port 9870, pa koristite http://localhost:9870/ tamo umjesto toga.

Poฤetna stranica web suฤelja NameNode u pregledniku

Sada odaberite 'Pregledaj datoteฤni sustav' i idite na /mapreduce_output_sales

Preglednik HDFS datoteka koji prikazuje direktorij mapreduce_output_sales

Otvori dio-r-00000

Datoteka rezultata part-r-00000 otvorena je u HDFS pregledniku

Objaลกnjenje klase SalesMapper

S obzirom na to da se posao izvodi od poฤetka do kraja, sljedeฤ‡a tri odjeljka opisuju ลกto svaka klasa zapravo radi.

U ovom odjeljku ฤ‡emo razumjeti implementaciju klase SalesMapper.

1. Poฤinjemo odreฤ‘ivanjem naziva paketa za naลกu klasu. SalesCountry je naziv naลกeg paketa. Imajte na umu da ฤ‡e izlaz kompilacije, SalesMapper.class, iฤ‡i u direktorij nazvan ovim nazivom paketa: SalesCountry.

Nakon toga uvozimo knjiลพniฤne pakete.

Donja snimka prikazuje implementaciju klase SalesMapper-

Prikaz urednika za cijelu implementaciju klase SalesMapper

Uzorak Code Objaลกnjenje:

1. Definicija klase SalesMapper-

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

Svaka klasa mappera mora biti proลกirena iz klase MapReduceBase i mora implementirati Mapper suฤelje.

2. Definiranje funkcije 'map'-

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

Glavni dio klase Mapper je metoda 'map()' koja prihvaฤ‡a ฤetiri argumenta.

Pri svakom pozivu metode 'map()', prosljeฤ‘uje se par kljuฤ-vrijednost ('kljuฤ' i 'vrijednost' u ovom kodu).

Metoda 'map()' zapoฤinje dijeljenjem ulaznog teksta koji se prima kao argument. Svaki redak dijeli na polja.

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

Ovdje se ',' koristi kao razdjelnik.

Nakon toga, par se formira koriลกtenjem zapisa na 7. indeksu polja 'SingleCountryData' i vrijednosti '1'.

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

Zapis na 7. indeksu odabiremo jer nam trebaju podaci o drลพavi, a on se nalazi na 7. indeksu u nizu 'SingleCountryData'.

Imajte na umu da su naลกi ulazni podaci u sljedeฤ‡em formatu (gdje je Drลพava na 7. indeksu, s 0 kao poฤetnim indeksom) -

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

Izlaz mappera je opet par kljuฤ-vrijednost koji se emitira pomoฤ‡u metode 'collect()' klase 'OutputCollector'.

Objaลกnjenje klase SalesCountryReducer

U ovom odjeljku ฤ‡emo razumjeti implementaciju klase SalesCountryReducer.

1. Poฤinjemo odreฤ‘ivanjem naziva paketa za naลกu klasu. SalesCountry je naziv naลกeg paketa. Imajte na umu da ฤ‡e izlaz kompilacije, SalesCountryReducer.class, iฤ‡i u direktorij nazvan ovim nazivom paketa: SalesCountry.

Nakon toga uvozimo knjiลพniฤne pakete.

Donja snimka prikazuje implementaciju klase SalesCountryReducer.

Prikaz urednika za kompletnu implementaciju klase SalesCountryReducer

Code Objaลกnjenje:

1. Definicija klase SalesCountryReducer-

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

Ovdje su prva dva tipa podataka, 'Text' i 'IntWritable', tip podataka ulaznog kljuฤa-vrijednosti za reducer.

Izlaz mappera je u obliku , Ovaj izlaz mappera postaje ulaz za reducer. Dakle, radi usklaฤ‘ivanja s njegovim tipom podataka, ovdje se kao tip podataka koriste Text i IntWritable.

Posljednja dva tipa podataka, 'Text' i 'IntWritable', su tipovi podataka izlaza koje generira reducer u obliku para kljuฤ-vrijednost.

Svaka klasa reducera mora biti proลกirena iz klase MapReduceBase i mora implementirati Reducer suฤelje.

2. Definiranje funkcije reduciranja-

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

Ulaz u metodu reduce() je kljuฤ s popisom viลกe vrijednosti.

Na primjer, u naลกem sluฤaju to ฤ‡e biti-

, , , , , .

To se daje reduktoru kao

Dakle, za prihvaฤ‡anje argumenata ovog oblika koriste se prva dva tipa podataka, i to tekst i iterator. Tekst je tip podataka kljuฤa i iteratora. je tip podataka za popis vrijednosti za taj kljuฤ.

Sljedeฤ‡i argument je tipa OutputCollector koji prikuplja izlaz faze redukcije.

Metoda reduce() zapoฤinje kopiranjem vrijednosti kljuฤa i inicijalizacijom brojaฤa frekvencija na 0.

Text key = t_key;
int frequencyForCountry = 0;

Zatim, koristeฤ‡i petlju 'while', iteriramo kroz popis vrijednosti povezanih s kljuฤem i izraฤunavamo konaฤnu frekvenciju zbrajanjem svih vrijednosti.

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

Sada rezultat ลกaljemo na izlazni kolektor u obliku kljuฤa i dobivenog brojaฤa frekvencija.

Donji kod radi ovo-

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

Objaลกnjenje klase SalesCountryDriver

U ovom odjeljku ฤ‡emo razumjeti implementaciju klase SalesCountryDriver.

1. Poฤinjemo odreฤ‘ivanjem naziva paketa za naลกu klasu. SalesCountry je naziv naลกeg paketa. Imajte na umu da ฤ‡e izlaz kompilacije, SalesCountryDriver.class, iฤ‡i u direktorij nazvan ovim nazivom paketa: SalesCountry.

Ovdje je redak koji navodi naziv paketa nakon kojeg slijedi kod za uvoz paketa knjiลพnice.

Deklaracija paketa i izjave o uvozu SalesCountryDrivera

2. Definirajte klasu pokretaฤa koja ฤ‡e kreirati novi posao klijenta, konfiguracijski objekt i oglaลกavati klase Mapper i Reducer.

Klasa upravljaฤkog programa odgovorna je za postavljanje naลกeg MapReduce posla za izvoฤ‘enje HadoopU ovoj klasi odreฤ‘ujemo naziv posla, tip podataka ulaza/izlaza i nazive klasa mappera i reducera.

Kod upravljaฤkog programa koji postavlja naziv posla, klase kljuฤeva i vrijednosti te formate

3. U donjem isjeฤku koda postavljamo ulazne i izlazne direktorije koji se koriste za potroลกnju ulaznog skupa podataka i proizvodnju izlaza.

arg[0] i arg[1] su argumenti naredbenog retka koji se prosljeฤ‘uju s naredbom zadanom u praktiฤnom MapReduceu, tj.

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

Kod upravljaฤkog programa koji prosljeฤ‘uje ulazne i izlazne putanje iz naredbenog retka

4. Pokrenite naลก posao

Donji kod pokreฤ‡e izvrลกavanje MapReduce posla -

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

Pitanja i odgovori

Sve tri klase uvoze org.apache.hadoop.mapred, originalni MapReduce API. Hadoop 3.x ga i dalje isporuฤuje i pokreฤ‡e, ali novi rad se obiฤno piลกe na org.apache.hadoop.mapreduce, koji zamjenjuje JobConf i JobClient s Configuration i Job.

Modeli obuฤeni na povijesti proลกlih poslova predviฤ‘aju vrijeme izvoฤ‘enja, predlaลพu veliฤine podjele i broj reduktora te uoฤavaju odstupanja prije zavrลกetka izvoฤ‘enja. Takoฤ‘er oznaฤavaju poslove ฤiji brojaฤi prolivenih zapisa ili neuspjelih zadataka izlaze izvan uobiฤajenog raspona.

Copilot dobro dovrลกava repetitivne dijelove: potpise klasa, generike, uvoze i pozive konfiguracije upravljaฤkog programa. Poslovnu logiku, poput stupca koji sadrลพi drลพavu, programer i dalje mora provjeriti u odnosu na stvarnu shemu.

Hadoop odbija pisati u postojeฤ‡i izlazni direktorij tako da se gotovi rezultati nikada ne prepisuju. Rekurzivno izbriลกite mapreduce_output_sales s hdfs dfs -rm -r ili proslijedite drugi izlazni put pri sljedeฤ‡em pokretanju.

Alat jar ignorira zavrลกni redak manifesta koji nema terminator retka, pa se Main-Class tiho odbacuje. Jar se zatim gradi bez greลกke, ali ne uspijeva tijekom izvoฤ‘enja jer nije zabiljeลพena glavna klasa.

Da. Primjer navodi 2.2.0 u svakom nazivu jar datoteke. U drugom izdanju zamijenite tu verziju ili jednostavno upotrijebite izlaz hadoop classpath-a, koji ispisuje svaku jar datoteku koju instalirana distribucija treba.

Ne. Pseudodistribuirana instalacija s jednim ฤvorom pokreฤ‡e ga nepromijenjenog, jer naredbe pretpostavljaju samo da su HDFS i YARN pokrenuti. Isti JAR se ลกalje na pravi klaster s viลกe ฤvorova bez ikakve promjene koda.

Promijenite indeks polja u maperu sa 7 na 5, jer je City ลกesti stupac datoteke SalesJan2009.csv. Ponovno kompajlirajte, ponovno izgradite JAR datoteku i pokrenite zadatak na novom izlaznom direktoriju.

Saลพmite ovu objavu uz: