Apache Hadoop — 1つのPodにHDFSとYARNを立てて運用する
アクセスログを Python の MapReduce で集計する
目標
Hadoop StreamingでPythonのマッパー・リデューサーを書いて、アクセスログ3日分をステータスコード別に数えます。コンバイナーがシャッフルをどれだけ減らすかをカウンターで測り、日付でリデューサーを選ぶ分割(KeyFieldBasedPartitioner)、壊れた行を数えるユーザーカウンター、わざと死ぬマッパーによるタスクの失敗と再試行まで扱います。
なぜ重要なのか
Streamingは、マップとリデュースを標準入出力の契約に置き換えます。マッパーは、入力の行を標準入力で受け取り、「キー値」の行を標準出力に出力し、フレームワークがキーでソートして、リデューサーの標準入力に流します。リデューサーが受け取るのは、キーごとのまとまりではなく、ソートされた行の流れなので、キーが変わる場所を自分で見つける必要があります。この契約さえ守れば、どんな言語でもMapReduceを書け、ローカルでcat | mapper | sort | reducerとして同じようにテストできます。
コンバイナーは、マップ側で事前にまとめるリデューサーです。結果が変わらないためには、演算が結合法則と交換法則を満たす必要があり(合計・最大値は使え、平均は使えません)、入力と出力の形が同じである必要があります。正しく付ければ、シャッフルが何十分の1に減ります。
デフォルトの分割は、キー全体のハッシュです。キーが「日付パス」で、1つの日付のすべてのパスを1つのリデューサーで見たいなら、キーの最初のフィールドだけで分割するように変える必要があります。そして、本番でストリーミングジョブが静かに間違う最も一般的な理由は、壊れた行です。捨てるのはかまいませんが、何行を捨てたかを数えないことはよくありません。
ステップ
- /root/hdp/streaming/mapper.pyを書いてください。1行を空白で分割して、フィールドが10個以上あり、9つ目のフィールド(0から数えて8)が3桁の数字で、
#で始まらない行なら、상태코드<TAB>1(プレースホルダーはステータスコードです)を出力し、そうでなければ捨てます。 - /root/hdp/streaming/reducer.pyを書いてください。キーでソートされた
키<TAB>수の行を受け取り、同じキーの数を足して키<TAB>합を出力します(プレースホルダーは順にキー、数、キー、合計です)。数は1でなくてもかまいません。 /data/logs/access-2026-03-01.log–03.logを、HDFSにアップロードし(アップロード先: /user/root/streaming/in/)、ジョブ名hdp-statusでストリーミングジョブを実行して、結果を、/user/root/streaming/statusに書き込んでください。- 同じジョブに
-combinerでリデューサーを付け、ジョブ名hdp-status-comb、出力先は、/user/root/streaming/status_combで実行して、そのジョブのMAP_OUTPUT_RECORDS・COMBINE_OUTPUT_RECORDS・REDUCE_INPUT_RECORDSを、/root/hdp/streaming/shuffle.jsonに{"map_output": …, "combine_output": …, "reduce_input": …}の形式で書き込んでください。 - /root/hdp/streaming/mapper2.pyが
날짜(yyyy-MM-dd)<TAB>경로<TAB>1(プレースホルダーは順に日付、パスです)を出力するようにし、ステップ1のreducer.pyをそのままリデューサーとして使って、ジョブ名hdp-daily、リデューサー2個、キー2フィールド、最初のフィールドだけで分割する設定で実行して、結果を、/user/root/streaming/dailyに書き込んでください。 - ステップ1のマッパーに、壊れた行ごとにユーザーカウンター
Lab・BadLinesを1上げる行を入れ(標準エラー出力にreporter:counter:Lab,BadLines,1)、ジョブ名hdp-badで実行して、結果を、/user/root/streaming/badに書き込んでください。 - /root/hdp/streaming/crash_mapper.py(
/event/springパスに出会うと例外で終わるマッパー)を使って、ジョブ名hdp-crash、マップ試行1回(-D mapreduce.map.maxattempts=1)で、結果を、/user/root/streaming/crashに書き込むジョブを実行して失敗させ、そのジョブIDを、/root/hdp/streaming/crash.txtに書いたあと、同じジョブ名でステップ1のマッパーを使って、結果を、/user/root/streaming/crash-okに書き込むジョブをもう一度実行して成功させてください。 - /root/hdp/streaming/report.mdに、
## 표준 입출력 계약・## 컴바이너와 분할・## 실패와 카운터の3つの節を書いてください(見出しは韓国語で、順に「標準入出力の契約」「コンバイナーと分割」「失敗とカウンター」を意味します)。2つ目の節にステップ4のmap_outputとreduce_inputを入れてください。
参考
- ストリーミングジョブ:
mapred streaming -D mapreduce.job.name=<이름> -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input <입력> -output <출력>(プレースホルダーは順に名前、入力、出力です)。-Dオプションは、ほかのオプションより前に置きます。 - ローカルテスト:
head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | python3 reducer.py。 - キー2フィールド・最初のフィールドで分割: 汎用オプション
-D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1は-filesの前に、コマンドオプション-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2は-mapperの後に置きます。順序が違うと、「Unrecognized option」で止まります。 - ジョブ履歴は、
/tmp/hadoop-yarn/staging/history/done_intermediate/root/に残ります(ジョブ名の-は、ファイル名では%2Dと書かれます)。失敗したジョブも残ります。 - よくあるミスは、リデューサーがキーごとに1回だけ呼ばれると思い込むこと(行ごとに流れ込んできます)、コンバイナーの出力の形を入力と違うものにすること、出力ディレクトリを削除せずに再実行することです。
- 公式ドキュメント: Hadoop Streaming・MapReduce Tutorial・mapred-default.xml
マッパー(1行を受け取ってキーと値を出力する)
/root/hdp/streaming/mapper.pyを書いてください。標準入力の行ごとに空白で分割し、フィールドが10個以上あり、9つ目のフィールド(0から数えて8)が3桁の数字で、#で始まらなければ、상태코드<TAB>1を標準出力に出力し(プレースホルダーはステータスコードです)、そうでなければ何も出力しません。
head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | uniq -cで、まずテストしてください。採点ツールは、自分のマッパーを、壊れた行が混ざった新しい入力で直接実行してみます。
リデューサー(ソートされた流れの中でキーが変わる場所を見つける)
/root/hdp/streaming/reducer.pyを書いてください。標準入力でキー順にソートされた키<TAB>수の行を受け取り、同じキーの数を足して、키<TAB>합を標準出力に出力します(プレースホルダーは順にキー、数、キー、合計です)。数は1でないこともあります。
リデューサーは、キーごとのリストを受け取りません。ソートされた行が流れ込んでくるだけなので、前の行とキーが変わる瞬間に、前のキーの合計を出力して、新しく始めます。最後のキーを出力するのを忘れないでください。数が1でなくてもよいように書けば、そのままコンバイナーになります。
YARNに載せる
lab-hadoop start yarnでYARNを起動し、/data/logs/access-2026-03-01.log–03.logを、HDFSにアップロードしたあと(アップロード先: /user/root/streaming/in/)、mapred streaming -D mapreduce.job.name=hdp-status -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input /user/root/streaming/in -output /user/root/streaming/statusで実行してください。
-filesで渡したスクリプトは、各タスクの作業フォルダーにコピーされます。そのため、-mapperには、パスなしでファイル名だけを書きます。採点ツールは、出力が元のログから直接数えた値と同じか、履歴に成功したhdp-statusジョブがあるかを確認します。
コンバイナーでシャッフルを減らす
ステップ3のジョブに-combiner 'python3 reducer.py'を加えて、ジョブ名hdp-status-comb、出力先は、/user/root/streaming/status_combで実行してください。そのジョブの履歴から、TaskCounterのMAP_OUTPUT_RECORDS・COMBINE_OUTPUT_RECORDS・REDUCE_INPUT_RECORDSを取り出して、/root/hdp/streaming/shuffle.jsonに{"map_output": 정수, "combine_output": 정수, "reduce_input": 정수}の形式で書き込んでください(プレースホルダーは順に整数、整数、整数です)。
合計は結合法則と交換法則を満たすので、リデューサーをそのままコンバイナーとして使っても、結果は同じです。コンバイナーのあと、シャッフルに渡った行数(=リデュース入力)が、マップ出力よりどれだけ少ないかを見てください。
日付でリデューサーを選ぶ
/root/hdp/streaming/mapper2.pyが、有効な行ごとに날짜(yyyy-MM-dd)<TAB>경로<TAB>1を出力するようにし(プレースホルダーは順に日付、パスです)、reducer.pyをリデューサーとして使って、ジョブ名hdp-daily、-D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1(汎用オプション、前側)と、-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2(コマンドオプション、後ろ側)で、結果を、/user/root/streaming/dailyに書き込んでください。
キーは最初の2つのフィールド(日付・パス)なので、ソートは2つのフィールドで行いますが、分割は最初のフィールド(日付)だけで行います。そのため、1つの日付のすべてのパスが1つのリデューサーに行き、出力ファイル1つが、日付1つを丸ごと持ちます。reducer.pyが最後のタブで分割するように書かれていれば、キーが2フィールドでも、そのまま動作します。
壊れた行をユーザーカウンターで数える
mapper.pyが、壊れた行ごとに標準エラー出力へreporter:counter:Lab,BadLines,1を書くようにして、ジョブ名hdp-badで、結果を、/user/root/streaming/badに書き込むジョブを実行してください。そのジョブのカウンターLab・BadLinesが、3つのログの壊れた行数と同じである必要があります。
ストリーミングは、標準エラー出力のreporter:counter:<그룹>,<이름>,<증가량>の行を、カウンターの更新として理解します(プレースホルダーは順にグループ、名前、増加量です)。このカウンターはジョブ履歴に残るので、ジョブが終わったあとも、「何行を捨てたか」に答えられます。
死ぬマッパー(失敗と再実行)
/root/hdp/streaming/crash_mapper.py(/event/springパスに出会うと例外で終わるマッパー)を使って、ジョブ名hdp-crash、-D mapreduce.map.maxattempts=1で、結果を、/user/root/streaming/crashに書き込むジョブを実行して失敗させ、そのジョブIDを、/root/hdp/streaming/crash.txtに書いてください。そのあと、同じジョブ名でステップ1のmapper.pyを使って、結果を、/user/root/streaming/crash-okに書き込むジョブをもう一度実行して成功させてください。
マップタスクが失敗すると、フレームワークは別の試行で再実行します(デフォルトは4回)。試行を1回に減らすと、最初の失敗がそのままジョブの失敗です。失敗したジョブも履歴が残るので、何がなぜ死んだかは、yarn logs -applicationId application_<같은 숫자>(プレースホルダーは同じ数字です)でタスクの標準エラー出力を読んで確認します。
契約・コンバイナー・失敗を書き残す
/root/hdp/streaming/report.mdに、## 표준 입출력 계약・## 컴바이너와 분할・## 실패와 카운터の3つの節を書いてください(見出しは韓国語で、順に「標準入出力の契約」「コンバイナーと分割」「失敗とカウンター」を意味します)。2つ目の節にステップ4のmap_outputとreduce_inputを、数値で入れてください。
リデューサーが受け取るものは何か、コンバイナーがシャッフルを何分の1に減らしたか、壊れた行と失敗したタスクをどうやって見つけたかを書いてください。