Hadoop MapReduce'i näide: esimene Java Programm koos Code

⚡ Nutikas kokkuvõte

Hadoop MapReduce programmid on kirjutatud kolmena Java klassid, kaardistaja, reduktor ja draiver, mis kompileeritakse, pakitakse purki ja esitatakse klastrile, et lugeda müüki riigiti.

  • 🔘 Andmekogum: Failis SalesJan2009.csv on üks tehing rea kohta, kusjuures riik on kaheksandas komadega eraldatud veerus.
  • ☑️ Kaardistaja: SalesMapper jagab iga rea ​​ja väljastab riigi, mis on seotud konstantse väärtusega üks.
  • Reduktor: SalesCountryReducer summeerib iga riigi kohta saabuvad väärtused üheks sagedusloenduriks.
  • 🧪 Juht: SalesCountryDriver paneb tööle nime, deklareerib võtme- ja väärtusetüübid ning ühendab mõlemad klassid omavahel.
  • 🛠️ Ehitada: Kompileeri javaciga, lisa Main-klassi manifesti kirje ja seejärel paki kõik jar cfm-iga.
  • ⚠️ Run: Kopeeri CSV HDFS-i, saada JAR-fail üles ja loe väljundkataloogist part-00000.

Hadoop MapReduce'i näide esimese loomisest Java programm

Selles õpetuses õpite kasutama Hadoopi koos MapReduce'i näidetega. Kasutatud sisendandmed on MüükJaan2009.csvSee sisaldab müügiga seotud teavet, nagu toote nimi, hind, makseviis, kliendi linn ja riik. Eesmärk on välja selgitada igas riigis müüdud toodete arv.

Esimene programm Hadoop MapReduce

Nüüd selles MapReduce'i õpetus, loome oma esimese Java MapReduce programm:

Allolev ekraanipilt näitab SalesJan2009 toorandmeid, kus iga rida tähistab ühte tehingut ja riik asub kaheksandas komaga eraldatud veerus.

SalesJan2009 müügiandmed avatud arvutustabelis, mis näitab tehingute veerge

Veenduge, et teil oleks Hadoop installitud. Enne tegeliku protsessi alustamist muutke kasutajaks 'hduser' (ID, mida kasutati Hadoopi konfigureerimisel – saate lülituda oma Hadoopi konfigureerimise ajal kasutatud kasutajatunnusele).

su - hduser_

Viip muutub hduser kontoks, nagu allpool näidatud.

Terminali viip pärast hduser-kontole üleminekut

1. samm) Looge projektikataloog ja lähtekoodifailid

Loo uus kataloog nimega MapReduceTutorial, nagu on näidatud allolevas MapReduce'i näites.

sudo mkdir MapReduceTutorial

mkdir käsk loob MapReduceTutorial kataloogi

Andke load

sudo chmod -R 777 MapReduceTutorial

chmod käsk annab MapReduceTutorial'ile täielikud õigused

Loo kolm Java MapReduceTutoriali allolevad lähtekoodifailid. Pange tähele, et kõik kolm kasutavad vanemat versiooni. org.apache.hadoop.mapred API, mis on endiselt kaasas Hadoop 3.x-ga.

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();
        }
    }
}

Laadige failid alla siit

Arhiiv laieneb samadeks kolmeks lähtefailiks, nagu siin näidatud.

ExtracTedi arhiiviloend SalesMapper, SalesCountryReducer ja SalesCountryDriver

Kontrollige kõigi nende failide failiõigusi

Pikk nimekiri, mis näitab kolme failiõigusi Java lähtefailid

Kui lugemisõigused puuduvad, siis andke need:

chmod käsk lisab lugemisõiguse Java lähtefailid

2. samm) Ekspordi Hadoopi klassirada

Ekspordi klassirada, nagu on näidatud allolevas Hadoopi näites.

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

Eksporditud klassitee kuvatakse käsuviibal tagasi.

Shelli viip pärast Hadoop CLASSPATH muutuja eksportimist

3. samm) Koostage Java failid

Koostage Java failid (need failid asuvad kataloogis Final-MapReduceHandsOn). Nende klassifailid pannakse paketikataloogi.

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

javaci väljund, mis näitab kompileerimise ajal aegumise hoiatust

Seda hoiatust võib ohutult ignoreerida – see annab ainult teada, et kaardistatud API on aegunud.

See kompileerimine loob praegusesse kataloogi kataloogi, mille nimi on määratud käsus Java lähtefail (st meie puhul SalesCountry) ja pane sinna kõik kompileeritud klassifailid.

SalesCountry paketikataloog, mis sisaldab kolme kompileeritud klassifaili

4. samm) Looge manifestifail

Loo uus fail Manifest.txt

sudo gedit Manifest.txt

Lisage sellele järgmine rida:

Main-Class: SalesCountry.SalesCountryDriver

gedit aken, mis kuvab Manifest.txt failis Main-Class kirjet

SalesCountry.SalesCountryDriver on põhiklassi nimi. Pane tähele, et selle rea lõpus tuleb vajutada sisestusklahvi (Enter).

5. samm) Pakkige klassid purki

Looge Jar-fail

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

jar-käsu abil ProductSalePerCountry.jar faili loomine manifestist

Kontrollige, kas jar-fail on loodud

Kataloogiloend, mis kinnitab ProductSalePerCountry.jar faili loomist

6. samm) Käivitage Hadoop

Käivitage Hadoop

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

7. samm) Kopeeri sisendfail HDFS-i

Kopeeri fail SalesJan2009.csv kausta ~/inputMapReduce

Nüüd kasuta allolevat käsku, et kopeerida ~/inputMapReduce HDFS-i.

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

copyFromLocal käsu väljund koos natiivse teegi hoiatusega

Võime seda hoiatust julgelt ignoreerida.

Kontrollige, kas fail on tegelikult kopeeritud või mitte.

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls kuvab faili SalesJan2009.csv faili inputMapReduce kataloogis

8. samm) Käivita MapReduce'i töö

Käivitage MapReduce'i töö

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

Konsooli väljund, kui klastris töötab SalePerCountry töö

See loob väljundkataloogi nimega mapreduce_output_sales on HDFS. Selle kataloogi sisu on fail, mis sisaldab toodete müüki riigiti.

9. samm) Lugege tulemusi

Tulemust saab käsurea kaudu näha järgmiselt:

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

Riigi ja müügiarvu paarid, mille prindib käsk hdfs dfs -cat

Tulemusi saab näha ka veebiliidese kaudu,

avatud http://localhost:50070/ veebibrauseris. Hadoop 3.x-s kolis NameNode veebiliides porti 9870, seega kasutage http://localhost:9870/ hoopis seal.

NameNode veebiliidese avaleht brauseris

Nüüd vali „Sirvi failisüsteemi” ja navigeeri kausta /mapreduce_output_sales

HDFS-failibrauser, mis kuvab kataloogi mapreduce_output_sales

Ava osa-r-00000

osa-r-00000 tulemusfail avatud HDFS-brauseris

SalesMapperi klassi selgitus

Töö käigus otsast lõpuni vaatavad järgmised kolm osa läbi, mida iga klass tegelikult teeb.

Selles osas vaatleme SalesMapperi klassi rakendamist.

1. Alustame oma klassi paketi nime määramisega. Müügiriik on meie paketi nimi. Pane tähele, et kompileerimise väljund SalesMapper.class läheb selle paketi nimega kataloogi: Müügiriik.

Seejärel impordime teegipakette.

Allolev hetktõmmis näitab SalesMapperi klassi rakendamist -

SalesMapperi klassi täieliku implementatsiooni redaktorivaade

Proov Code Selgitus:

1. SalesMapperi klassi definitsioon-

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

Iga kaardistaja klass peab olema laiendatud MapReduceBase klassist ja see peab implementeerima kaardistaja liidese.

2. Kaardi funktsiooni defineerimine-

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

Mapperi klassi põhiosa on meetod 'map()', mis aktsepteerib nelja argumenti.

Iga 'map()' meetodi kutsumise korral edastatakse võtme-väärtuse paar (selles koodis 'võti' ja 'väärtus').

Meetod 'map()' alustab argumendina saadud sisendteksti jagamisega. See jagab iga rea ​​väljadeks.

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

Siin kasutatakse eraldajana ','.

Pärast seda moodustatakse paar, kasutades massiivi 'SingleCountryData' 7. indeksi kirjet ja väärtust '1'.

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

Valime 7. indeksi kirje, kuna vajame riigi andmeid ja need asuvad massiivi 'SingleCountryData' 7. indeksis.

Pange tähele, et meie sisendandmed on allolevas vormingus (kus riik on 7. indeksil ja 0 on algusindeks):

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

Kaardistaja väljund on jällegi võtme-väärtuse paar, mis väljastatakse 'OutputCollector' meetodi 'collect()' abil.

SalesCountryReduceri klassi selgitus

Selles osas mõistame SalesCountryReduceri klassi rakendamist.

1. Alustame oma klassi paketi nime määramisega. Müügiriik on meie paketi nimi. Pane tähele, et kompileerimise väljund SalesCountryReducer.class läheb selle paketi nimega kataloogi: Müügiriik.

Seejärel impordime teegipakette.

Allolev hetktõmmis näitab SalesCountryReduceri klassi rakendamist -

SalesCountryReduceri klassi täieliku implementatsiooni redaktorivaade

Code Selgitus:

1. SalesCountryReduceri klassi definitsioon-

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

Siin on kaks esimest andmetüüpi, „Text” ja „IntWritable”, reduktori sisendvõtme väärtuse andmetüübid.

Kaardistaja väljund on kujul , See kaardistaja väljund saab reduktori sisendiks. Seega, et see vastaks selle andmetüübile, kasutatakse siin andmetüüpidena Text ja IntWritable.

Kaks viimast andmetüüpi, „Text” ja „IntWritable”, on reduktori genereeritud väljundi andmetüübid võtme-väärtuse paari kujul.

Iga reduktoriklass peab olema laiendatud MapReduceBase klassist ja see peab implementeerima reduktoriliidese.

2. Vähendamisfunktsiooni määratlemine-

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

Meetodi reduce() sisendiks on võti, mis sisaldab mitmest väärtusest koosnevat loendit.

Näiteks meie puhul on see

, , , , , .

See antakse reduktorile kujul

Seega sellise vormi argumentide vastuvõtmiseks kasutatakse kahte esimest andmetüüpi, nimelt teksti ja iteraatorit. Tekst on võtme ja iteraatori andmetüüp. on selle võtme väärtuste loendi andmetüüp.

Järgmine argument on tüüpi OutputCollector mis kogub reduktori faasi väljundit.

Meetod reduce() alustab võtme väärtuse kopeerimise ja sagedusloenduri initsialiseerimisega väärtuseks 0.

Text key = t_key;
int frequencyForCountry = 0;

Seejärel, kasutades 'while' tsüklit, itereerime läbi võtmega seotud väärtuste loendi ja arvutame lõppsageduse, summeerides kõik väärtused.

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

Nüüd edastame tulemuse väljundkollektorisse võtme ja saadud sagedusloenduri kujul.

Allolev kood teeb seda -

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

SalesCountryDriver klassi selgitus

Selles osas vaatleme SalesCountryDriver klassi rakendamist.

1. Alustame oma klassi paketi nime määramisega. MüügiRiik on meie paketi nimi. Pane tähele, et kompileerimise väljund SalesCountryDriver.class läheb selle paketi nimega kataloogi: MüügiRiik.

Siin on rida, mis määrab paketi nime, millele järgneb kood teegipakettide importimiseks.

SalesCountryDriveri paketideklaratsioon ja impordiväljavõtted

2. Määratlege draiveriklass, mis loob uue klienditöö, konfiguratsiooniobjekti ja reklaamib Mapper ja Reducer klasse.

Juhiklass vastutab meie MapReduce'i töö käivitamise eest hadoopSelles klassis määrame töö nime, sisendi/väljundi andmetüübi ning kaardistaja ja reduktori klasside nimed.

Draiveri koodi seadistamine töö nime, võtme- ja väärtusklasside ning vormingute jaoks

3. Allpool olevas koodilõigul määrame sisend- ja väljundkataloogid, mida kasutatakse vastavalt sisendandmestiku tarbimiseks ja väljundi tootmiseks.

arg[0] ja arg[1] on käsurea argumendid, mis edastatakse MapReduce'i praktilises kasutusjuhendis antud käsuga, st

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

Draiveri kood, mis edastab käsurealt sisend- ja väljundteid

4. Käivitage meie töö

Allolev kood alustab MapReduce'i töö käivitamist

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

KKK

Kõik kolm klassi impordivad org.apache.hadoop.mapred'i, mis on algne MapReduce'i API. Hadoop 3.x on endiselt saadaval ja käitab seda, kuid uus töö kirjutatakse tavaliselt org.apache.hadoop.mapreduce'i vastu, mis asendab JobConfi ja JobClienti Configurationi ja Jobiga.

Varasemate tööde ajaloo põhjal treenitud mudelid ennustavad käitusaega, pakuvad välja jaotussuurusi ja vähenduste arvu ning tuvastavad enne käituse lõppu nihke. Samuti märgistavad need tööd, mille lekkinud kirjete või ebaõnnestunud ülesannete loendurid triivivad tavapärasest vahemikust väljapoole.

Copilot teeb korduvad osad hästi ära: klassi signatuurid, geneerilised koodid, import ja draiveri konfiguratsioonikõned. Äriloogikat, näiteks seda, milline veerg sisaldab riiki, peab arendaja ikkagi tegeliku skeemi suhtes kontrollima.

Hadoop keeldub kirjutamast olemasolevasse väljundkataloogi, et lõpptulemusi kunagi üle ei kirjutataks. Kustuta mapreduce_output_sales rekursiivselt käsuga hdfs dfs -rm -r või edasta järgmisel käivitamisel erinev väljundrada.

Jar-tööriist ignoreerib viimast manifesti rida, millel puudub rea lõpp, seega Main-Class kustutatakse vaikselt. Seejärel ehitatakse Jar-fail veatult, kuid ebaõnnestub käitusajal, kuna main-klassi pole salvestatud.

Jah. Näide lisab igale jar-faili nimele versiooni 2.2.0. Teises versioonis asenda see versioon või kasuta lihtsalt hadoop classpathi väljundit, mis prindib kõik jar-failid, mida installitud distributsioon vajab.

Ei. Ühesõlmelise pseudohajutatud installi korral töötab see muutmata kujul, sest käsud eeldavad ainult HDFS-i ja YARN-i käivitamist. Sama JAR-fail esitatakse päris mitmesõlmelisele klastrile ilma koodimuudatusteta.

Muutke kaardistaja massiivi indeks 7-lt 5-le, kuna City on SalesJan2009.csv faili kuues veerg. Kompileerige JAR-fail uuesti, looge see uuesti ja käivitage töö uue väljundkausta vastu.

Võta see postitus kokku järgmiselt: