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.

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.
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.
1. samm) Looge projektikataloog ja lähtekoodifailid
Loo uus kataloog nimega MapReduceTutorial, nagu on näidatud allolevas MapReduce'i näites.
sudo mkdir MapReduceTutorial
Andke load
sudo chmod -R 777 MapReduceTutorial
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(); } } }
Arhiiv laieneb samadeks kolmeks lähtefailiks, nagu siin näidatud.
Kontrollige kõigi nende failide failiõigusi
Kui lugemisõigused puuduvad, siis andke need:
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.
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
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.
4. samm) Looge manifestifail
Loo uus fail Manifest.txt
sudo gedit Manifest.txt
Lisage sellele järgmine rida:
Main-Class: SalesCountry.SalesCountryDriver
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
Kontrollige, kas jar-fail on loodud
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 /
Võime seda hoiatust julgelt ignoreerida.
Kontrollige, kas fail on tegelikult kopeeritud või mitte.
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
8. samm) Käivita MapReduce'i töö
Käivitage MapReduce'i töö
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
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
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.
Nüüd vali „Sirvi failisüsteemi” ja navigeeri kausta /mapreduce_output_sales
Ava osa-r-00000
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 -
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 -
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.
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.
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
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(); }





















