Apache Hadoop — 1つのPodにHDFSとYARNを立てて運用する
YARN はメモリをコンテナの大きさに切って配る
一言でいうと
YARNは、クラスターのメモリとCPUをコンテナという断片に切って、アプリに分け与えます。断片のサイズには下限と上限があり、下限より小さい要求は下限に引き上げられ、上限より大きい要求は拒否されます。誰に先に渡すかはキューが決め、終了したコンテナのログは、ログ集約をオンにして初めて1か所に残ります。
なぜ必要なのか
Hadoop 1のJobTrackerは、リソース管理とジョブの監視を1人でこなしていました。クラスターが大きくなるほど、この1つのプロセスがボトルネックになり、スロットがマップ用とリデュース用に最初から分かれていたため、マップだけが動く時間には、リデュースのスロットが遊んでいました。しかも、そのクラスターでは、MapReduce以外のものを動かせませんでした。
YARNアーキテクチャドキュメントは、これを2つに分けた設計だと説明しています。リソース管理は、クラスターに1つだけのResourceManagerが、ジョブの計画と監視は、アプリごとに1つずつ起動するApplicationMasterが担います。マシンごとに動くNodeManagerは、コンテナを起動し、リソース使用量を監視して報告します。ResourceManagerの中のSchedulerは、ドキュメントの表現で「純粋な」スケジューラーです。リソースを分け与えるだけで、アプリの状態を追跡したり、失敗したタスクを再起動したりしません。それはAMの役目です。そのため、MapReduceも、Sparkも、ほかのフレームワークも、それぞれのAMさえ持ってくれば、同じクラスターを分け合えます。
どう動くのか
ジョブを投入すると、ResourceManagerが最初のコンテナを確保してAMを起動します。AMは「メモリXMB、vcoreY個のコンテナをN個」という形で要求し、スケジューラーが、NodeManagerの空きに合わせて渡します。ここでサイズの規則が関わってきます。yarn-default.xmlによると、次のとおりです。
yarn.scheduler.minimum-allocation-mbデフォルト1024: これより小さい要求は、この値に引き上げられます。これよりメモリが小さく設定されたNodeManagerは、ResourceManagerが終了させます。yarn.scheduler.maximum-allocation-mbデフォルト8192: これより大きい要求は、InvalidResourceRequestExceptionで拒否されます。- vcoreは、最小1、最大4がデフォルトです。
リソースモデルドキュメントは、要求が最小・最大に合わせられたり、設定された増分の倍数に変更されたりすることがあると付け加えています。要求したサイズと受け取ったサイズが違うことがあるという意味です。MapReduceのAMは、タスクの要求が最大割り当てを超えることを事前に察知し、コンテナを要求する前に、ジョブをKILLEDで終了させ、その理由をアプリケーションの診断(Diagnostics)行に残します。
NodeManagerがコンテナに提供する総量は、yarn.nodemanager.resource.memory-mbです。デフォルト値-1は、ハードウェアの自動検出がオンのときだけ計算され、それ以外の場合は8192MBになります。マシンに実際にそれだけのメモリがあるかは、確認しません。
このラボのPodを数値で見ると、規則がなぜ重要かがわかります。Podのメモリは2Giで、NameNode・DataNode・ResourceManager・NodeManagerの4つのデーモンのヒープ上限だけを合わせても、864MB(256・160・256・192)です。NodeManagerがコンテナに渡す分は1024MBです。最小割り当てをデフォルトの1024MBにすると、コンテナ1つがノード全体です。AMがその1つを取ると、マップが入る場所がなく、ジョブはAMだけを起動したまま、永遠に待ちます。しかも、MapReduce AMのデフォルトの要求であるyarn.app.mapreduce.am.resource.mbは1536MBなので、1024MBのノードには、そもそも入りません。小さなノードでは、最小割り当てを下げ、AMとタスクの要求も併せて下げる必要があります。ラボのイメージが、最小割り当て128MB・最大割り当て1024MB、AM 384MB、マップ・リデュース256MBに設定してある理由です。
要求とJVMヒープは別物であることも覚えておいてください。mapreduce.map.memory.mbがコンテナのサイズで、ヒープは、デフォルトの割合0.8で、その中で決まります。残りの20%は、JVM自体とネイティブメモリの取り分です。NodeManagerはデフォルトで、物理メモリと仮想メモリの両方のチェックをオンにしていて(vmem-pmem-ratio 2.1)、超過したコンテナを強制終了します(ラボのイメージは、JVMの予約アドレス空間を使用量として誤って数える仮想メモリのチェックを、オフにしてあります)。
キューは、誰が先に受け取るかを決める
デフォルトのスケジューラーはCapacitySchedulerです。Capacity Schedulerドキュメントによると、すべてのキューはrootの子で、yarn.scheduler.capacity.root.queuesに、カンマ区切りで子を書きます。各キューのcapacityはパーセントで、1つの階層の合計が100である必要があります。この取り分は保証であって、フェンスではありません。ほかのキューが空いていれば、自分の取り分より多く使え、その融通をmaximum-capacityで制限します。
<property>
<name>yarn.scheduler.capacity.root.queues</name>
<value>default,etl,adhoc</value>
</property>
<property>
<name>yarn.scheduler.capacity.root.etl.capacity</name>
<value>60</value>
</property>
設定ファイルを変更したあとは、ResourceManagerを再起動せず、yarn rmadmin -refreshQueuesで反映します。ジョブはmapreduce.job.queuename(デフォルトdefault)でキューを選び、存在しないキューで投入したジョブは、投入の段階で拒否されます。
小さなクラスターでよく引っかかる値が、もう1つあります。maximum-am-resource-percentは、クラスターのリソースのうちAMに使える割合で、デフォルトは10%です。同時にアクティブになるアプリ数を、これが決めます。1024MBのノードの10%は、AM 1つよりも小さいので、ジョブを2つ同時に投入すると、2つ目がACCEPTEDのまま止まりやすいです。
ログはオンにして初めて集まる
コンテナのログは、そのコンテナを起動したNodeManagerのローカルディスクに残ります。yarn.log-aggregation-enableのデフォルトはfalseで、このとき、ログはyarn.nodemanager.log.retain-secondsのデフォルト10800秒(3時間)後に削除されます。オンにすると、アプリが終わったあと、ログがHDFSの/tmp/logsの下に移され、yarn logs -applicationId <앱 ID>(プレースホルダーはアプリIDです)の1行で、すべてのコンテナのログをまとめて読めます。ラボのイメージは、この値をオンにしてあります。ノードが数百台なら、これなしでは、失敗したタスク1つのログを探すのに半日かかります。
クラスターが今どれだけ空いているかは、ResourceManager RESTの/ws/v1/cluster/metricsで見ます。totalMB・allocatedMB・availableMBとappsPendingが一度に出ます。
現場での姿
1つ目は、ジョブがACCEPTEDから動かないことです。原因は、たいてい3つのうちのどれかです。AMが入れるだけの大きさのノードがないか、キューのAMの取り分がいっぱいか、AMは起動したのにタスクコンテナが入る場所がないかです。availableMBと要求サイズを並べて見れば、ほとんど解決します。
2つ目は、要求を減らしたのに、使うメモリがそのままであることです。最小割り当てより小さく要求すると、最小割り当てに引き上げられます。要求256MB、最小1024MBなら、受け取るのは1024MBです。
3つ目は、コンテナがメモリ超過で死ぬことです。ヒープは要求の中に収まっていても、スレッドスタックとネイティブバッファが、残りを超えた場合です。ヒープの割合を下げるか、要求を大きくします。
4つ目は、キューを削除したら、ジョブがすべて拒否されることです。投入スクリプトがデフォルトのdefaultキューを使っているのに、そのキューを一覧から外したからです。
実務で本当に大切なこと
- 要求したサイズと受け取ったサイズは違います。下限に引き上げられ、上限を超えれば拒否されます。
- NodeManagerの取り分は、実際のメモリを知りません。デフォルトのままだと、8192MBだと嘘をつきます。
- コンテナのサイズはヒープより大きいです。デフォルトでヒープはその0.8で、残りはJVM自体の取り分です。
- キューの取り分は保証であって、フェンスではありません。フェンスはmaximum-capacityです。
- キューの変更はrefreshQueuesで反映します。再起動はしません。
- ログ集約はデフォルトでオフです。オンにしておかないと、数時間後に失敗の証拠が消えます。
次のラボですること
ResourceManagerのRESTでクラスター全体のメモリと空きを読み、マップ1つに最大割り当てより大きい2048MBを要求したジョブが、コンテナを受け取れずKILLEDで終わることを、アプリケーションの状態の診断行で確認します。capacity-scheduler.xmlにetlとadhocのキューを加えて、保証枠の合計を100に合わせ、refreshQueuesで反映したあと、etlキューで投入したジョブがそのキューで動くことと、存在しないキューで投入したジョブが投入の段階で拒否されることを見ます。スケジューラーRESTでキューの状態を読み、終わったジョブの集約されたログをyarn logsで取得して、レポートにまとめます。