Hadoop MapReduceの例:まず Java とのプログラム Code
⚡ スマートサマリー
Hadoop MapReduceプログラムは3つの形式で記述されます。 Java クラス、マッパー、リデューサー、ドライバーがコンパイルされ、jarファイルにパッケージ化されてクラスターに送信され、国別の売上をカウントします。
このチュートリアルでは、MapReduce の例で Hadoop の使用方法を学習します。 使用される入力データは、 SalesJan2009.csvこのデータには、商品名、価格、支払い方法、顧客の居住都市と国など、販売関連情報が含まれています。目的は、各国で販売された商品の数を把握することです。
最初の Hadoop MapReduce プログラム
今これで MapReduce チュートリアル、私たちは最初の Java MapReduce プログラム:
以下のスクリーンショットは、SalesJan2009の生データを示しています。各行は1つの取引を表し、国は8番目のカンマ区切りの列に記載されています。
Hadoopがインストールされていることを確認してください。実際の処理を開始する前に、ユーザーを「hduser」(Hadoopの設定時に使用したID。ご自身のHadoop設定時に使用したユーザーIDに切り替えることもできます)に変更してください。
su - hduser_
プロンプトが以下のようにhduserアカウントに変わります。
ステップ1)プロジェクトディレクトリとソースファイルを作成する
以下の MapReduce の例に示すように、MapReduceTutorial という名前の新しいディレクトリを作成します。
sudo mkdir MapReduceTutorial
権限を与える
sudo chmod -R 777 MapReduceTutorial
3つを作成する Java 以下のソースファイルはMapReduceTutorial内にあります。3つとも古いバージョンを使用していることに注意してください。 org.apache.hadoop.mapred このAPIは、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); } }
Sales CountryReducer.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)); } }
Sales CountryDriver.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(); } } }
アーカイブは、ここに示されているように、同じ3つのソースファイルに展開されます。
これらすべてのファイルのファイル権限を確認してください
「読み取り」権限が不足している場合は、以下の権限を付与してください。
ステップ2)Hadoopクラスパスをエクスポートする
以下の Hadoop の例に示すように、クラスパスをエクスポートします。
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/*"
エクスポートされたクラスパスがプロンプトに表示されます。
ステップ3)コンパイル Java ファイル
コンパイル Java ファイル(これらのファイルは Final-MapReduceHandsOn ディレクトリに存在します)。クラスファイルはパッケージディレクトリに配置されます。
javac -d . SalesMapper.java SalesCountryReducer.java SalesCountryDriver.java
この警告は無視しても問題ありません。これは単にmapred APIが非推奨になったことを報告しているだけです。
このコンパイルでは、現在のディレクトリに、指定されたパッケージ名で名前が付けられたディレクトリが作成されます。 Java ソースファイル(この場合はSalesCountry)を作成し、コンパイル済みのクラスファイルをすべてその中に配置します。
ステップ4)マニフェストファイルを作成する
Manifest.txtという新しいファイルを作成します。
sudo gedit Manifest.txt
そこに次の行を追加してください。
Main-Class: SalesCountry.SalesCountryDriver
SalesCountry.SalesCountryDriver はメインクラスの名前です。この行の最後に Enter キーを押す必要があることに注意してください。
ステップ5)クラスをジャーにパッケージ化する
Jarファイルを作成する
jar cfm ProductSalePerCountry.jar Manifest.txt SalesCountry/*.class
jarファイルが作成されたことを確認する
ステップ6)Hadoopを起動する
Hadoopを開始する
$HADOOP_HOME/sbin/start-dfs.sh
$HADOOP_HOME/sbin/start-yarn.sh
ステップ7)入力ファイルをHDFSにコピーする
SalesJan2009.csv ファイルを ~/inputMapReduce にコピーします。
次に、以下のコマンドを使用して ~/inputMapReduce を HDFS にコピーします。
$HADOOP_HOME/bin/hdfs dfs -copyFromLocal ~/inputMapReduce /
この警告は無視しても問題ありません。
実際にファイルがコピーされたかどうかを確認します。
$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce
ステップ8)MapReduceジョブを実行する
MapReduce ジョブを実行する
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
これにより、mapreduce_output_sales という名前の出力ディレクトリが作成されます。 HDFS。 このディレクトリの内容は、国別の製品売上高を含むファイルになります。
ステップ9)結果を読む
結果はコマンドインターフェースを通して以下のように確認できます。
$HADOOP_HOME/bin/hdfs dfs -cat /mapreduce_output_sales/part-00000
結果は Web インターフェイス経由でも確認できます。
店は開いています http://localhost:50070/ Web ブラウザで。Hadoop 3.x では、NameNode の Web UI はポートに移動しました。 9870、使用する http://localhost:9870/ 代わりにそこにあります。
次に「ファイルシステムの参照」を選択し、/mapreduce_output_sales に移動します。
パートr-00000を開く
SalesMapperクラスの説明
ジョブが最初から最後まで実行されることを踏まえ、次の3つのセクションでは、各クラスが実際にどのような役割を果たすのかを詳しく説明します。
このセクションでは、SalesMapperクラスの実装について理解を深めます。
1. まず、クラスのパッケージ名を指定します。SalesCountry がパッケージ名です。コンパイルの出力ファイル SalesMapper.class は、このパッケージ名「SalesCountry」と同じ名前のディレクトリに保存されますのでご注意ください。
続いて、ライブラリ パッケージをインポートします。
以下のスナップショットは、SalesMapperクラスの実装を示しています。
サンプル Code 説明:
1. SalesMapper クラスの定義 -
public class SalesMapper extends MapReduceBase implements Mapper<LongWritable, Text, Text, IntWritable> {
すべてのマッパークラスはMapReduceBaseクラスを継承し、Mapperインターフェースを実装する必要があります。
2.「マップ」関数の定義 -
public void map(LongWritable key, Text value, OutputCollector<Text, IntWritable> output, Reporter reporter) throws IOException
Mapperクラスの主要部分は、4つの引数を受け取る「map()」メソッドです。
'map()' メソッドが呼び出されるたびに、キーと値のペア (このコードでは 'key' と 'value') が渡されます。
'map()' メソッドは、引数として受け取った入力テキストを分割することから始まります。各行をフィールドに分割します。
String valueString = value.toString(); String[] SingleCountryData = valueString.split(",");
ここでは、「,」が区切り文字として使用されています。
その後、配列「SingleCountryData」の7番目のインデックスにあるレコードと値「1」を使用してペアが作成されます。
output.collect(new Text(SingleCountryData[7]), one);
国データが必要なので、7番目のインデックスにあるレコードを選択します。このデータは配列「SingleCountryData」の7番目のインデックスにあります。
入力データは以下の形式になっていることにご注意ください(国名は7番目のインデックスにあり、0が開始インデックスです)。
Transaction_date,Product,Price,Payment_Type,Name,City,State,Country,Account_Created,Last_Login,Latitude,Longitude
マッパーの出力は、OutputCollectorのcollect()メソッドを使用して出力されるキーと値のペアです。
Sales CountryReducerクラスの説明
このセクションでは、SalesCountryReducer クラスの実装について理解を深めます。
1. まず、クラスのパッケージ名を指定します。パッケージ名は SalesCountry です。コンパイルの出力ファイル SalesCountryReducer.class は、このパッケージ名 SalesCountry という名前のディレクトリに格納されますのでご注意ください。
続いて、ライブラリ パッケージをインポートします。
以下のスナップショットは、SalesCountryReducer クラスの実装を示しています。
Code 説明:
1. Sales CountryReducer クラスの定義 -
public class SalesCountryReducer extends MapReduceBase implements Reducer<Text, IntWritable, Text, IntWritable> {
ここで、最初の 2 つのデータ型である「Text」と「IntWritable」は、リデューサーへの入力キーと値のデータ型です。
マッパーの出力は、 、このマッパーの出力はリデューサーへの入力となります。そのため、データ型に合わせて、ここではTextとIntWritableがデータ型として使用されます。
最後の2つのデータ型、「Text」と「IntWritable」は、リデューサーによって生成される出力のデータ型であり、キーと値のペアの形式をとります。
すべてのリデューサークラスはMapReduceBaseクラスを継承し、Reducerインターフェースを実装する必要があります。
2. 「reduce」関数の定義 -
public void reduce( Text t_key, Iterator<IntWritable> values, OutputCollector<Text,IntWritable> output, Reporter reporter) throws IOException {
reduce() メソッドへの入力は、複数の値のリストを持つキーです。
たとえば、私たちの場合は次のようになります-
、 、 、 、 、 。
これはリデューサーに次のように渡されます
この形式の引数を受け入れるには、まずテキストとイテレータという2つのデータ型が使用されます。テキストはキーとイテレータのデータ型ですこれは、そのキーに対応する値のリストを表すデータ型です。
次の引数は OutputCollector 型ですこれは、リデューサーフェーズの出力を収集します。
reduce() メソッドは、まずキーと値をコピーし、頻度カウントを 0 に初期化することから始まります。
Text key = t_key;int frequencyForCountry = 0;
次に、「while」ループを使用して、キーに関連付けられた値のリストを繰り返し処理し、すべての値を合計して最終的な頻度を計算します。
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));
Sales CountryDriverクラスの説明
このセクションでは、SalesCountryDriver クラスの実装について理解を深めます。
1. まず、クラスのパッケージ名を指定します。パッケージ名は SalesCountry です。コンパイル結果の出力ファイル SalesCountryDriver.class は、このパッケージ名「SalesCountry」と同じ名前のディレクトリに保存されますのでご注意ください。
ここでは、パッケージ名を指定する行と、その後にライブラリ パッケージをインポートするコードが続きます。
2. 新しいクライアント ジョブ、構成オブジェクトを作成し、Mapper クラスと Reducer クラスをアドバタイズするドライバー クラスを定義します。
ドライバー クラスは、MapReduce ジョブを実行するように設定する責任があります。 Hadoopのこのクラスでは、ジョブ名、入出力のデータ型、マッパークラスとリデューサークラスの名前を指定します。
3. 以下のコード スニペットでは、入力データセットの使用と出力の生成にそれぞれ使用される入力ディレクトリと出力ディレクトリを設定します。
arg[0]とarg[1]は、MapReduceの実践で与えられたコマンドとともに渡されるコマンドライン引数です。
$HADOOP_HOME/bin/hadoop jar ProductSalePerCountry.jar /inputMapReduce /mapreduce_output_sales
4. ジョブをトリガーする
以下のコードはMapReduceジョブの実行を開始します。
try { // Run the job JobClient.runJob(job_conf); } catch (Exception e) { e.printStackTrace(); }





















