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

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

MapReduce のコストはマップとリデュースの間のシャッフルに集中している

TT Labで続きを見る

一言でいうと

MapReduceは、マップ → (分割・ソート・シャッフル) → リデュースという枠組みで、ユーザーが書くのは、前後の2つの関数だけです。中間のソートとシャッフルはフレームワークが行い、ジョブのコストの大部分がそこで生じます。コンバイナーはその中間を減らす仕組みで、カウンターはその中間をのぞき込む窓です。

なぜこのような形が必要だったのか

数百台に散らばったテラバイトのデータで、単語の数を数えなければならないとします。データを1台のサーバーに集めると、ネットワークが先に死にます。HDFS設計ドキュメントが掲げる原則が、「計算をデータの近くに移すほうが、データを移すより安い」である理由です。すると、計算をブロックが置かれたサーバーごとに分けて実行する必要がありますが、今度は、分けて実行した結果を同じ単語同士で再び集める必要があります。サーバーが死ねば、その分を再実行しなければならず、遅いサーバーは待たなければなりません。

この問題をジョブごとに新しく解くと、コードの大部分が分散処理の後始末になります。MapReduceは、その後始末をフレームワークに切り出しました。MapReduceチュートリアルは、フレームワークがタスクの配置と監視、失敗したタスクの再実行を担い、アプリケーションは入出力の場所とmap・reduce関数だけを渡すと整理しています。代償は、すべての計算をキーと値のペアで表現しなければならないという制約です。その制約のおかげで、「同じキーは同じリデューサーに行く」という約束1つで、集約を解決します。

どう動くのか

流れをペアの形で書くと、次のようになります。チュートリアルに載っているとおりです。

(입력) <k1, v1> -> map -> <k2, v2> -> combine -> <k2, v2> -> reduce -> <k3, v3> (출력)

マップ側: チュートリアルによると、マップの数は、通常、入力の全体サイズ、つまり入力ファイルのブロック数が左右します。ブロック1つが、たいてい入力スプリット(split)1つになり、スプリットごとにマップタスクが1つ起動します。マップが出力するペアは、すぐにはディスクに行かず、メモリバッファに溜まります。mapred-default.xmlで、このバッファサイズmapreduce.task.io.sort.mbのデフォルトは100MBで(タスクのヒープが180MBしかないラボのイメージは、32MBに減らしてあります)、mapreduce.map.sort.spill.percentのデフォルト0.80に届くと、後ろのスレッドが内容をソートしてディスクに書き出します(spill)。マップが終わると、書き出した断片を1つにマージします。このファイルは、リデューサーの数だけのパーティションに分かれており、各パーティションの中は、キーでソートされています。

どのキーがどのパーティションに行くかは、パーティショナー(Partitioner)が決めます。デフォルトはHashPartitionerで、キーのハッシュをリデューサー数で割った余りです。同じキーは、どのマップから出てきても、同じ番号のリデューサーに行きます。集約が成り立つ根拠が、この1行です。

リデュース側: チュートリアルは、リデューサーがシャッフル・ソート・リデュースの3つの段階を経ると説明しています。シャッフルで、すべてのマップ出力のうち自分のパーティションだけをHTTPで取得し、取得している間にマージソートを行って、同じキーの値を1列に並べます。そのあと、キーごとにreduce(키, 값 목록)(プレースホルダーは順にキー、値のリストです)を1回ずつ呼びます。リデューサー数を0にすると、マップ出力がそのまま結果になり、ソートもしません。

リデューサー数mapreduce.job.reducesのデフォルトは1です。設定しなければ、すべてのマップ出力が1つのリデューサーに集中し、結果ファイルもpart-r-000001つです。リデューサーを3つに増やせば、結果ファイルが3つになります。リデューサーはキーの順に呼ばれるので、単語の数え上げのように、受け取ったキーをそのまま使うジョブなら、各ファイルの中はキー順です。ただし、チュートリアルが明言するように、フレームワークがリデューサーの出力を再びソートすることはなく、3つのファイルをつなげると、全体の順序は混ざります。ハッシュで分けたからです。

全体を1列にソートするには、パーティショナーを変える必要があります。TeraSortの例が、その方法を示しています。入力からキーのサンプルを取ってN-1個の境界を決め、境界の区間ごとにリデューサーを1つずつ割り当てます。すると、リデューサーiの出力はすべてリデューサーi+1より小さく、ファイルを番号順につなげると、全体がソートされています。TeraValidateは、それを検査するジョブです。

コンバイナーとカウンター

単語の数え上げのマップは、(the, 1)を何万回も出力します。これをそのままシャッフルすると、ネットワークを1が何万回も渡ります。コンバイナーは、マップ側で事前に(the, 38211)に減らして送るローカル集計です。チュートリアルは、コンバイナーがスピルのときに動き、リデュース側のマージ中にも動くことがあり、バッファより大きなレコードがコンバイナーを通るかは決まっていないと書いています。まとめると、コンバイナーが何回動くかは保証されません。そのため、1回動いても3回動いても結果が同じになる演算だけを、コンバイナーとして使えます。合計と最大値は使え、平均は使えません。平均の平均は平均ではないからです。

コンバイナーが役目を果たしているかは、カウンターで見ます。ジョブが終わると、フレームワークが数えた値が出力されます。マップ入力レコードとマップ出力レコード、コンバイナー入力と出力レコード、リデュース入力グループ数、スピルしたレコード数などです。コンバイナー出力が入力よりはるかに小さければ、シャッフルがその分減ったという意味です。単語の数え上げなら、リデュース入力グループ数が、そのまま結果に出る異なる単語の数です。カウンターは、ジョブが何をしたかを、ログを漁らずに数値で証明する最も安い方法です。

YARNの上で動くということ

Hadoop 2以降、MapReduceは、自分ではリソースを管理しません。YARNアーキテクチャドキュメントが説明するように、クラスター全体を調停するResourceManagerと、アプリごとに1つ起動するApplicationMasterに役割が分かれ、MapReduceジョブのApplicationMasterがMRAppMasterです。ジョブを1つ投入すると、コンテナが、最低でもAM 1つ、マップタスクの数だけ、リデューサーの数だけ起動します。

ジョブが終わると、AMは履歴を残します。デフォルトのステージングディレクトリが/tmp/hadoop-yarn/stagingで、中間の完了ディレクトリは、その下のhistory/done_intermediateです。ジョブ履歴サーバーを起動しなくても、ここに.jhist(イベント履歴)と.summary、ジョブ設定のXMLが、ユーザーごとに溜まります。履歴サーバーは、これを移して見せる役割にすぎず、記録自体はAMが行います。

現場での姿

1つ目は、ジョブがクラスターではなくローカルで動くことです。mapreduce.framework.nameのデフォルトはlocalです。yarnに変更していないクライアントで投入したジョブは、ResourceManagerに現れず、投入したマシンのJVM 1つで静かに動きます。「なぜResourceManagerの画面にジョブがないのか」に対する、最もよくある答えです。ラボのイメージは、この値をyarnに変えてあります。

2つ目は、AMが入る場所がないことです。yarn.app.mapreduce.am.resource.mbのデフォルトは1536MBです。NodeManagerがコンテナに渡す分がそれより小さいクラスターでは、AMコンテナが入るノードがなく、ジョブは開始すらできません。このラボ環境のように、割り当てが1024MBの場所では、AMのサイズを別に下げておかないと、ジョブが起動しません。ラボのイメージは、AMを384MB、マップとリデュースを256MBに設定して、AM 1つとタスク2つが同時に入るようにしてあります。

3つ目は、リデューサー1つが終わらないことです。デフォルト値の1のままにしているか、キー1つに値が集中しています。前者は設定で直し、後者はキーの設計で直します。

4つ目は、コンバイナーのせいで数値が間違うことです。平均や中央値を、リデューサーのコードのままコンバイナーとして設定すると、コンバイナーが動く回数によって結果が変わります。合ったり間違ったりするバグなので、見つけにくいです。

実務で本当に大切なこと

次のラボですること

YARNをオンにしてNodeManagerが登録されたかを確認したあと、本5冊をHDFSにアップロードして、サンプルjarの単語カウントのジョブを実行し、ジョブIDを書き留めます。ジョブ履歴サーバーなしで、HDFSのdone_intermediateに残ったジョブ履歴ファイルを探して開きますが、ラボのイメージは、このファイルを人間が読めるJSONの行で書くようにしてあります。ジョブが終わるときの全カウンターから、マップ入力・出力とコンバイナー入力・出力、リデュース入力グループ数と出力レコードを取り出して、コンバイナーがシャッフルをどれだけ減らしたかを計算します。リデューサーを3つに増やして結果ファイルがどう分かれるかを見て、最後にteragen・terasort・teravalidateで全体ソートを実行し、検証まで通します。