TT Lab
はじめる
学ぶ 学習パス コース

Apache Hadoop — 1つのPodにHDFSとYARNを立てて運用する

Hadoop Streaming では標準入出力の一行がそのまま契約になる

TT Labで続きを見る

一言でいうと

Hadoop Streamingは、マッパーとリデューサーを、標準入力を読み、標準出力に書く任意のプログラムに差し替えるツールです。契約は単純です。1行が1レコード、最初のタブの前がキー、リデューサーにはキーでソートされた行が来ます。単純な分、キーが変わる境界をスクリプトが自分で見つける必要があり、失敗とカウンターも、終了コードと標準エラー出力で伝える必要があります。

なぜ必要なのか

MapReduceの元の契約は、Javaのインターフェースです。Mapperを継承し、Writable型を合わせ、jarにまとめる必要があります。アクセスログでステータスコード別のリクエスト数を数える作業に、それだけの形式は重すぎます。しかも、データを扱う人の大半は、Pythonやシェルで、すでに1台のマシンで動くスクリプトを持っています。

Hadoop Streamingドキュメントは、この隙間を埋めます。実行ファイルやスクリプトなら何でも、マッパーとリデューサーとして使えて、ドキュメントの最初の例は、マッパーに/bin/cat、リデューサーに/usr/bin/wcを使います。1台でcat log | map.py | sort | reduce.pyとして動いていたパイプラインが、ほぼそのまま数百台に広がります。途中のsortを、MapReduceのシャッフルが代わりに担うわけです。

どう動くのか

ドキュメントによると、マップタスクは、開始時に指定された実行ファイルを別のプロセスとして起動します。入力スプリットを行に変換して、そのプロセスの標準入力に流し込み、標準出力から出てくる行を集めて、キーと値のペアに変換します。デフォルトの規則は、最初のタブ文字の前がキー、後ろが値で、タブがなければ、行全体がキーで、値は空です。リデューサー側も同じです。フレームワークがソート・マージしたペアを、再び키\t값(プレースホルダーは順にキー、値です)の行に戻して、リデューサープロセスの標準入力に入れます。

この契約で最も重要な事実は、リデューサーが受け取るのが、キーごとのまとまりではなく、ソートされた行の流れだという点です。Javaのリデューサーはreduce(키, 값 목록)(プレースホルダーは順にキー、値のリストです)をキーごとに1回ずつ呼ばれて受け取りますが、ストリーミングのリデューサーは、行を1つずつ読みながら、「キーが変わったか」を自分で判断する必要があります。MapReduceチュートリアルが言うとおり、マップ出力はソートされたあと、リデューサーごとに分けられて届くので、同じキーの行は必ず連続して来ます。リデューサーは、その保証1つに頼って、前の行のキーを覚えておき、キーが変わる瞬間に、合計を出力して初期化します。

import sys

current, total = None, 0
for line in sys.stdin:
    key, _, value = line.rstrip("\n").partition("\t")
    if key != current:
        if current is not None:
            print(f"{current}\t{total}")
        current, total = key, 0
    total += int(value)
if current is not None:
    print(f"{current}\t{total}")

最後のキーを、ループの後でもう一度出力する行を忘れるのが、最もよくあるミスです。そうすると、結果からソート上の一番最後のキー1つが、静かに消えます。

マッパーとリデューサーのスクリプトは、ワーカーノードにある必要があります。-filesで渡すと、タスクの作業フォルダーに、同じ名前のシンボリックリンクができます。ドキュメントは、-filesや-Dのような汎用オプションをストリーミングオプションより前に置かなければ、コマンドが失敗すると警告しています。

シャッフルを扱うノブ

コンバイナー: -combinerで、実行ファイルをもう1つ渡せます。ステータスコードごとに1を出力するマッパーなら、リデューサーと同じ合計スクリプトをコンバイナーとしてかけて、シャッフルに渡る行数を大きく減らせます。コンバイナーが何回動くかは保証されないので、何回適用しても結果が同じ合計系だけをかけます。

キーフィールドと分割: キーが2026-09-19.404のように複数のフィールドでできているとき、ソートは全体のキーで行い、分割は前のフィールドだけでしたいことがあります。ドキュメントのKeyFieldBasedPartitionerの例が、それです。stream.num.map.output.key.fieldsでキーが何フィールドか、map.output.key.field.separatorでフィールドの区切り文字を、mapreduce.partition.keypartitioner.options=-k1,1のように、分割に使うフィールドを指定します。すると、同じ日付の行はすべて同じリデューサーに行き、その中では、全体のキー順にソートされて届きます。日付ごとの結果ファイルが欲しいときや、1つのリデューサーの中で日付ごとの小計を出したいときに使います。

リデューサー数とマップだけのジョブ: デフォルトのリデューサー数は、mapred-default.xmlのmapreduce.job.reducesのデフォルト値1です。ドキュメントによると、0にすると、リデューサーなしで、マップ出力がそのまま結果になります。行を絞り込んだり、形式を変えるだけの仕事なら、シャッフル自体をなくすのが最も速いです。

失敗とカウンターの伝え方

ストリーミングのスクリプトは、Java APIがないので、フレームワークと2つの経路だけで話します。1つは終了コードです。ドキュメントによると、デフォルトでは、0以外の終了コードで終わったタスクは失敗として扱われます。失敗したタスクは再試行され、mapreduce.map.maxattemptsのデフォルトが4なので、同じマップが4回失敗するとジョブ全体が失敗します。壊れた行1つで例外を投げるマッパーは、その行が入ったスプリットで4回死んで、ジョブを道連れにします。

もう1つは標準エラー出力です。reporter:counter:<그룹>,<카운터>,<양>の形式の行をstderrに書くとカウンターが上がり、reporter:status:<메시지>はステータスの文言を変えます(プレースホルダーは順にグループ、カウンター、量、メッセージです)。そのため、壊れた行は、例外で死ぬのではなく、スキップしながらreporter:counter:logs,malformed,1を書くほうがよいです。ジョブは最後まで動き、何行が壊れていたかは、カウンターに数値として残ります。標準出力に混ぜて書くと、それが結果レコードになってしまうので、必ずstderrに送ります。

設定値も、環境変数として読めます。ドキュメントによると、設定名のドットがアンダースコアに変わり、たとえばmapreduce.job.idはmapreduce_job_idとして見えます。

現場での姿

1つ目は、ローカルでは合うのに、クラスターで最後のキーが抜けることです。リデューサーが、ループの後で最後のまとまりを出力していない場合です。1台でsortを使ってテストするとき、小さなデータでは、たまたま表に出ないこともあります。

2つ目は、ジョブが4回再試行して死ぬことです。ログの1行のエンコードが壊れていたり、フィールド数が足りない行があったりします。Pythonのマッパーがその行で例外を投げると、同じスプリットの再試行も、同じ行で死にます。再試行は、一時的な障害のための仕組みであって、悪いデータのための仕組みではありません。

3つ目は、カウンターが結果ファイルに出力されることです。reporter:counterをprintで標準出力に書いたからです。結果にreporter:counter:...で始まるキーができます。

4つ目は、スクリプトがノードにないというエラーです。-filesを忘れたか、ストリーミングオプションの後に置いたかです。実行権限と、1行目のインタープリターの指定も、併せて確認します。

実務で本当に大切なこと

次のラボですること

アクセスログからステータスコードごとに1を出力するPythonのマッパーと、ソートされた入力を読み継いで合計を出すリデューサーを作り、ストリーミングジョブとして実行します。同じ合計スクリプトをコンバイナーとしてかけて、マップ出力とリデュース入力のレコード数をカウンターで比べ、日付とパスをタブでつないだ2つのフィールドのキーを、KeyFieldBasedPartitionerで最初のフィールド(日付)だけを見て、リデューサー2つに分けます。壊れた行は標準エラー出力でユーザーカウンターを上げて数え、特定のパスでわざと死ぬマッパーをマップ試行1回で実行して、ジョブがすぐに失敗する様子を見たあと、同じ名前のジョブを正常なマッパーでもう一度実行して成功させます。