Hadoop MapReduceの例:まず Java とのプログラム Code

⚡ スマートサマリー

Hadoop MapReduceプログラムは3つの形式で記述されます。 Java クラス、マッパー、リデューサー、ドライバーがコンパイルされ、jarファイルにパッケージ化されてクラスターに送信され、国別の売上をカウントします。

  • 🔘 データセット: SalesJan2009.csvファイルには、1行に1件の取引データが格納されており、8番目の列にカンマ区切りで国名が記載されています。
  • ☑️ マッパー: SalesMapperは各行を分割し、国名と定数値1をペアにして出力します。
  • 減速機: SalesCountryReducerは、各国ごとに到着した件数を合計して、単一の頻度カウントにまとめます。
  • 🧪 ドライバ: SalesCountryDriver はジョブに名前を付け、キーと値の型を宣言し、両方のクラスを接続します。
  • 🛠️ ビルド: javacでコンパイルし、Main-Classマニフェストエントリを追加してから、jar cfmですべてをパッケージ化します。
  • ⚠️ 実行: CSVファイルをHDFSにコピーし、jarファイルを送信して、出力ディレクトリからpart-00000を読み取ります。

Hadoop MapReduce の例で最初のものを作成する Java プログラム

このチュートリアルでは、MapReduce の例で Hadoop の使用方法を学習します。 使用される入力データは、 SalesJan2009.csvこのデータには、商品名、価格、支払い方法、顧客の居住都市と国など、販売関連情報が含まれています。目的は、各国で販売された商品の数を把握することです。

最初の Hadoop MapReduce プログラム

今これで MapReduce チュートリアル、私たちは最初の Java MapReduce プログラム:

以下のスクリーンショットは、SalesJan2009の生データを示しています。各行は1つの取引を表し、国は8番目のカンマ区切りの列に記載されています。

SalesJan2009の売上データがスプレッドシートで開かれ、取引列が表示されています。

Hadoopがインストールされていることを確認してください。実際の処理を開始する前に、ユーザーを「hduser」(Hadoopの設定時に使用したID。ご自身のHadoop設定時に使用したユーザーIDに切り替えることもできます)に変更してください。

su - hduser_

プロンプトが以下のようにhduserアカウントに変わります。

hduserアカウントに切り替えた後のターミナルプロンプト

ステップ1)プロジェクトディレクトリとソースファイルを作成する

以下の MapReduce の例に示すように、MapReduceTutorial という名前の新しいディレクトリを作成します。

sudo mkdir MapReduceTutorial

mkdirコマンドでMapReduceTutorialディレクトリを作成します

権限を与える

sudo chmod -R 777 MapReduceTutorial

MapReduceTutorial に完全な権限を付与する chmod コマンド

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つのソースファイルに展開されます。

Extractedアーカイブ一覧 SalesMapper、SalesCountryReducer、SalesCountryDriver

これらすべてのファイルのファイル権限を確認してください

3つのファイルのアクセス許可を示す長いリスト Java ソースファイル

「読み取り」権限が不足している場合は、以下の権限を付与してください。

chmod コマンドで読み取り権限を追加します Java ソースファイル

ステップ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/*"

エクスポートされたクラスパスがプロンプトに表示されます。

Hadoop CLASSPATH変数をエクスポートした後のシェルプロンプト

ステップ3)コンパイル Java ファイル

コンパイル Java ファイル(これらのファイルは Final-MapReduceHandsOn ディレクトリに存在します)。クラスファイルはパッケージディレクトリに配置されます。

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

javacの出力でコンパイル中に非推奨警告が表示される

この警告は無視しても問題ありません。これは単にmapred APIが非推奨になったことを報告しているだけです。

このコンパイルでは、現在のディレクトリに、指定されたパッケージ名で名前が付けられたディレクトリが作成されます。 Java ソースファイル(この場合はSalesCountry)を作成し、コンパイル済みのクラスファイルをすべてその中に配置します。

SalesCountryパッケージディレクトリには、コンパイル済みの3つのクラスファイルが格納されています。

ステップ4)マニフェストファイルを作成する

Manifest.txtという新しいファイルを作成します。

sudo gedit Manifest.txt

そこに次の行を追加してください。

Main-Class: SalesCountry.SalesCountryDriver

gedit Manifest.txt の Main-Class エントリを表示するウィンドウ

SalesCountry.SalesCountryDriver はメインクラスの名前です。この行の最後に Enter キーを押す必要があることに注意してください。

ステップ5)クラスをジャーにパッケージ化する

Jarファイルを作成する

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

jarコマンドは、マニフェストからProductSalePerCountry.jarを構築します。

jarファイルが作成されたことを確認する

ProductSalePerCountry.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 /

ネイティブライブラリの警告を含む copyFromLocal コマンドの出力

この警告は無視しても問題ありません。

実際にファイルがコピーされたかどうかを確認します。

$HADOOP_HOME/bin/hdfs dfs -ls /inputMapReduce

hdfs dfs -ls コマンドで、inputMapReduce ディレクトリ内の SalesJan2009.csv を一覧表示します。

ステップ8)MapReduceジョブを実行する

MapReduce ジョブを実行する

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

クラスター上でSalePerCountryジョブが実行されている間のコンソール出力

これにより、mapreduce_output_sales という名前の出力ディレクトリが作成されます。 HDFS。 このディレクトリの内容は、国別の製品売上高を含むファイルになります。

ステップ9)結果を読む

結果はコマンドインターフェースを通して以下のように確認できます。

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

hdfs dfs -cat コマンドで表示される国と販売件数のペア

結果は Web インターフェイス経由でも確認できます。

店は開いています http://localhost:50070/ Web ブラウザで。Hadoop 3.x では、NameNode の Web UI はポートに移動しました。 9870、使用する http://localhost:9870/ 代わりにそこにあります。

NameNodeのWebインターフェースのホームページをブラウザで表示

次に「ファイルシステムの参照」を選択し、/mapreduce_output_sales に移動します。

mapreduce_output_salesディレクトリを表示するHDFSファイルブラウザ

パートr-00000を開く

part-r-00000 結果ファイルが HDFS ブラウザで開かれました

SalesMapperクラスの説明

ジョブが最初から最後まで実行されることを踏まえ、次の3つのセクションでは、各クラスが実際にどのような役割を果たすのかを詳しく説明します。

このセクションでは、SalesMapperクラスの実装について理解を深めます。

1. まず、クラスのパッケージ名を指定します。SalesCountry がパッケージ名です。コンパイルの出力ファイル SalesMapper.class は、このパッケージ名「SalesCountry」と同じ名前のディレクトリに保存されますのでご注意ください。

続いて、ライブラリ パッケージをインポートします。

以下のスナップショットは、SalesMapperクラスの実装を示しています。

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 クラスの実装を示しています。

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」と同じ名前のディレクトリに保存されますのでご注意ください。

ここでは、パッケージ名を指定する行と、その後にライブラリ パッケージをインポートするコードが続きます。

SalesCountryDriverのパッケージ宣言とインポートステートメント

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

よくあるご質問

これら3つのクラスはすべて、オリジナルのMapReduce APIであるorg.apache.hadoop.mapredをインポートします。Hadoop 3.xには引き続き同梱され、このAPIが実行されますが、新しいコードは通常org.apache.hadoop.mapreduceに対して記述され、JobConfとJobClientはConfigurationとJobに置き換えられます。

過去のジョブ履歴に基づいて学習されたモデルは、実行時間を予測し、分割サイズとリデューサー数を提案し、実行終了前にスキューを検出します。また、スピルレコード数や失敗タスク数が通常の範囲から外れたジョブも警告します。

Copilotは、クラスシグネチャ、ジェネリクス、インポート、ドライバ構成呼び出しといった反復的な部分をうまく処理します。しかし、国をどの列に格納するかといったビジネスロジックについては、開発者が実際のスキーマと照合して確認する必要があります。

Hadoop は、完了した結果が上書きされないように、既存の出力ディレクトリへの書き込みを拒否します。`hdfs dfs -rm -r` コマンドを使用して `mapreduce_output_sales` ディレクトリを再帰的に削除するか、次回の実行時に別の出力パスを指定してください。

jarツールは、行末記号のない最後のマニフェスト行を無視するため、Main-Classは警告なしに削除されます。その結果、jarファイルはエラーなくビルドされますが、メインクラスが記録されていないため、実行時に失敗します。

はい。この例では、すべてのjarファイル名に2.2.0を指定しています。別のリリースを使用する場合は、そのバージョンに置き換えるか、インストールされているディストリビューションに必要なすべてのjarファイルを表示するhadoop classpathの出力を使用してください。

いいえ。単一ノードの擬似分散環境では、コマンドはHDFSとYARNが起動していることを前提としているため、変更なしで実行されます。同じjarファイルを実際のマルチノードクラスタに投入しても、コードの変更は一切不要です。

マッパー内の配列インデックスを7から5に変更してください。CityはSalesJan2009.csvの6列目です。再コンパイルしてjarファイルを再構築し、新しい出力ディレクトリに対してジョブを実行してください。