Apache Hadoop — 1つのPodにHDFSとYARNを立てて運用する
単語カウントと TeraSort を YARN に載せ、カウンタで読む
目標
YARNを起動して、本5冊の単語をMapReduceで数えます。ジョブが終わったあとにHDFSに残ったジョブ履歴(.jhist)を開いてカウンターを読み、リデューサー数を変えたとき出力がどう分かれるか、シャッフルがソートをどう保証するかを、TeraSortで確認します。
なぜ重要なのか
MapReduceは、最近新しく書くことは少なくなりましたが、そのモデル、つまり、マップが(キー、値)を出力し、フレームワークがキーで集めてソートして渡し(シャッフル)、リデュースがキーごとにまとめるという仕組みは、Spark・Flink・SQLエンジンのシャッフルが、すべて受け継いでいます。シャッフルがなぜ高価なのか、コンバイナー(マップ側での事前集約)がなぜシャッフルを減らすのか、リデューサー数がなぜ出力ファイル数になるのかを、ここで最も生々しく見られます。 カウンターは、ジョブが何をしたかを数値として残した記録です。マップが何行を読んで何行を出力したか、コンバイナーが何行に減らしたか、リデュースがキーをいくつ受け取ったか。結果がおかしいときに最初に見る場所で、ジョブが終わったあとも、履歴ファイルに残っています。 このPodのYARNは、Podのメモリ(2Gi)に合わせて、コンテナを小さく設定してあります(AM 384MB、マップ・リデュース256MB)。そのため、タスクが1–2個ずつ順番に動きます。遅くても、結果とカウンターは同じです。
ステップ
lab-hadoop start yarnでYARNを起動し、yarn node -list -allの出力を、/root/hdp/mr/nodes.txtに保存してください。/data/books/book-*.txtの5冊を、HDFSにアップロードしてください(アップロード先: /user/root/mr/in/)。- サンプルjarの
wordcountで、/user/root/mr/inを数えて、結果を、/user/root/mr/wcに書き込んでください。 - ステップ3のジョブのID(
job_…)を、/root/hdp/mr/jobid.txtに書いてください。 - そのジョブの履歴ファイル(.jhist)から、TaskCounterの
MAP_INPUT_RECORDS・MAP_OUTPUT_RECORDS・COMBINE_INPUT_RECORDS・COMBINE_OUTPUT_RECORDS・REDUCE_INPUT_GROUPS・REDUCE_OUTPUT_RECORDSを取り出して、/root/hdp/mr/counters.jsonに書いてください。 - 同じwordcountを、リデューサー3個(
-D mapreduce.job.reduces=3)で実行し、結果を、/user/root/mr/wc3に書き込んでください。 teragenで10万行を作成し(出力先: /user/root/mr/tera-in)、terasortでソートして(出力先: /user/root/mr/tera-out)、teravalidateで検証結果を、/user/root/mr/tera-valに書き込んでください。- /root/hdp/mr/report.mdに、
## 맵과 리듀스・## 카운터・## 정렬の3つの節を書いてください(見出しは韓国語で、順に「マップとリデュース」「カウンター」「ソート」を意味します)。2つ目の節に、ステップ5のMAP_OUTPUT_RECORDSとCOMBINE_OUTPUT_RECORDSを入れてください。
参考
- サンプルjar:
/opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.5.0.jar。実行はyarn jar <jar> wordcount <입력> <출력>です(プレースホルダーは順に入力、出力です)。出力ディレクトリがすでにあると、ジョブは開始しません。 - ジョブ履歴:
/tmp/hadoop-yarn/staging/history/done_intermediate/root/に、<잡ID>-…-<잡 이름>-…-<상태>-<대기열>-….jhistと.summary・_conf.xmlがあります(プレースホルダーは順にジョブID、ジョブ名、状態、キューです)。.jhistは最初の2行(形式・スキーマ)のあとに、イベントが1行ずつJSONで入っていて、JOB_FINISHEDイベントのtotalCountersにカウンターがあります。 - ジョブ1つが1分前後かかります。YARNを起動しておくと、メモリを1GB近く余分に使うので、ほかの重い作業は一緒にしないでください。
- よくあるミス: 出力ディレクトリを削除せずに再実行すること、ジョブIDの代わりにアプリケーションID(
application_…)を書くこと(数字は同じで、名前が違います)。 - 公式ドキュメント: MapReduce Tutorial・Apache Hadoop YARN・mapred-default.xml・YARN Commands
YARNを起動する
lab-hadoop start yarnでResourceManagerとNodeManagerを起動し、yarn node -list -allの出力を、/root/hdp/mr/nodes.txtに保存してください。
NodeManagerが、自分が提供できるメモリ・コアをResourceManagerに知らせて初めて、ジョブを受け取れます。ノードがRUNNINGと表示される必要があります。lab-hadoop statusで、JVM 4つ(NameNode・DataNode・ResourceManager・NodeManager)が起動していることを確認してください。
入力をアップロードする
/data/books/book-1.txt–book-5.txtを、HDFSにアップロードしてください(アップロード先: /user/root/mr/in/)。
入力ディレクトリの中のファイルごとに入力スプリット(split)ができ、スプリット1つにマップタスク1つが付きます。ファイル5つは、マップ5つです。
wordcountを実行する
yarn jar /opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.5.0.jar wordcount /user/root/mr/in /user/root/mr/wcで、単語を数えてください。
ジョブが終わると、出力ディレクトリに_SUCCESSとpart-r-00000ができます。マップは単語ごとに(単語、1)を出力し、コンバイナーがマップ側で先に足し、リデューサーが単語ごとにまとめます。採点ツールは、自分の出力を、元データから直接数えた値と比べます。
ジョブ履歴を探す
ステップ3のジョブのID(job_<숫자>_<숫자>)を、/root/hdp/mr/jobid.txtに1行で書いてください(プレースホルダーは順に数字、数字です)。
ジョブを実行するとき、コンソールに「Running job: job_…」が出力されます。見逃したなら、HDFSの/tmp/hadoop-yarn/staging/history/done_intermediate/root/の一覧から、名前にword+countが入った.jhistを探してください。履歴サーバーなしでも、AMがここに残します。
カウンターを読む
ステップ4のジョブの.jhistから、org.apache.hadoop.mapreduce.TaskCounterグループのMAP_INPUT_RECORDS・MAP_OUTPUT_RECORDS・COMBINE_INPUT_RECORDS・COMBINE_OUTPUT_RECORDS・REDUCE_INPUT_GROUPS・REDUCE_OUTPUT_RECORDSを取り出して、/root/hdp/mr/counters.jsonに{"이름": 정수}の形式で書き込んでください(プレースホルダーは順に名前、整数です)。
MAP_OUTPUT_RECORDSは単語数(重複を含む)で、COMBINE_OUTPUT_RECORDSはマップごとのユニークな単語数の合計です。シャッフルで実際に渡された行数です。REDUCE_INPUT_GROUPSは、全体のユニークな単語数と同じになるはずです。2つの数値の比率が、コンバイナーが節約したシャッフルです。
リデューサーが3つなら出力も3つ
yarn jar <jar> wordcount -D mapreduce.job.reduces=3 /user/root/mr/in /user/root/mr/wc3で、リデューサー3個のジョブを実行してください。
マップ出力は、キーのハッシュでリデューサー3つのうち1つに割り当てられます(HashPartitioner)。そのため、1つの単語は、ちょうど1つのファイルにだけ出てきて、各ファイルはその中ではソートされていますが、ファイル同士はソートされていません。
TeraSort(シャッフルがソートを作る)
teragen 100000 /user/root/mr/tera-inで10万行(100バイトずつ)を作成し、terasort /user/root/mr/tera-in /user/root/mr/tera-outでソートしたあと、teravalidate /user/root/mr/tera-out /user/root/mr/tera-valで検証してください(すべて同じサンプルjarです)。
TeraSortは、入力をサンプリングして、キーの範囲をリデューサー数だけに切り(全順序パーティショニング)、各リデューサーが自分の範囲をソートして出力します。すると、出力ファイルを順につなげたものが、全体のソートになります。TeraValidateは、ずれた場所があれば「misorder」を、なければチェックサムだけを書きます。
シャッフルを数値で残す
/root/hdp/mr/report.mdに、## 맵과 리듀스・## 카운터・## 정렬の3つの節を書いてください(見出しは韓国語で、順に「マップとリデュース」「カウンター」「ソート」を意味します)。2つ目の節に、ステップ5のMAP_OUTPUT_RECORDSとCOMBINE_OUTPUT_RECORDSの値を、数値で入れてください。
最初の節には、マップ・リデュースがいくつあって、出力がどう分かれたか、2つ目の節には、コンバイナーがシャッフルを何分の1に減らしたか、3つ目の節には、TeraSortが全体の順序をどう作ったかを書いてください。