Apache Hadoop — Stand up and run HDFS and YARN in one pod
Size containers to fit a 2Gi pod and split the queues
Goal
You check with the ResourceManager's REST API what resources this cluster can offer, and see how a job that asks for a container larger than that ends. You create the etl and adhoc queues in the Capacity Scheduler and send jobs, and cover that a nonexistent queue is rejected and how to gather and read the logs of a finished job.
Why it matters
YARN is a resource manager that hands out a cluster's memory and cores in units of containers. The NodeManager announces the share it can offer (yarn.nodemanager.resource.memory-mb), and the ResourceManager's scheduler rounds requests up to the minimum allocation unit and rejects those larger than the maximum allocation. Everything that runs on YARN, whether MapReduce or Spark, follows these rules.
In operations, "my job doesn't move from ACCEPTED" and "my job dies right away" are usually problems of these numbers. If the container size is larger than the node's share, room never opens up, and if it is larger than the maximum allocation, it is rejected immediately. This Pod has 2Gi of memory, so the node share was set to 1024MB: the numbers are small, so the rules are easier to see.
Queues are a way for several teams to share one cluster. The Capacity Scheduler sets, for each queue, a guaranteed share (capacity) and an upper limit it may borrow up to (maximum-capacity). Under the same parent, the sum of guaranteed shares must be 100, and after editing the configuration file, it is reread with refreshQueues. The logs of finished containers are gathered to HDFS by the NodeManager (log aggregation), so you can read them with yarn logs even after the Pod has disappeared.
Steps
- Turn on YARN with
lab-hadoop start yarnand save the response ofcurl -s http://localhost:8088/ws/v1/cluster/metricsto /root/hdp/yarn/metrics.json. - Run the example jar's
piwith-D mapreduce.map.memory.mb=2048and let the job end, then save the output ofyarn application -status <애플리케이션 ID>(where the placeholder is the application ID) for that application to /root/hdp/yarn/toobig.txt. - Edit /opt/hadoop/etc/hadoop/capacity-scheduler.xml to set the queues under
roottodefault,etl,adhoc, the guaranteed shares to 20, 50 and 30, and the upper limit ofadhocto 50, runyarn rmadmin -refreshQueues, and then save the output ofyarn queue -status etlto /root/hdp/yarn/etl.txt. - Upload
/data/books/book-1.txtto HDFS /user/root/yarn/in/ and runwordcountwith-D mapreduce.job.queuename=etlto /user/root/yarn/wc-etl. - Try submitting the same wordcount to the nonexistent queue
nope(output /user/root/yarn/wc-nope) and save the error output to /root/hdp/yarn/badqueue.txt. - Save the response of
curl -s http://localhost:8088/ws/v1/cluster/schedulerto /root/hdp/yarn/scheduler.json. - Gather the application logs of the step 4 job with
yarn logs -applicationId <애플리케이션 ID>(the placeholder is the application ID) and save them to /root/hdp/yarn/etl-logs.txt. - In /root/hdp/yarn/report.md, write three sections:
## 컨테이너 크기,## 대기열and## 로그(use exactly these Korean headings in this order; they mean "Container size", "Queues" and "Logs"). Put thetotalMBfrom step 1 and the 2048 requested in step 2 in the first section.
Notes
- YARN in this Pod: NodeManager share 1024MB and 2 cores, minimum allocation 128MB, maximum allocation 1024MB, AM 384MB, map and reduce 256MB (
yarn-site.xmlandmapred-site.xml). - The job ID
job_<A>_<B>and the application IDapplication_<A>_<B>use the same numbers. The history of finished jobs is left in/tmp/hadoop-yarn/staging/history/done_intermediate/root/. - Because the history server was not started, for a failed job, the client may end with a connection error (10020) while fetching the final state. See the job's real result with
yarn application -statusor the history file. - Common mistakes: not making the sum of guaranteed shares 100, so refreshQueues is rejected; submitting with
root.attached to the queue name (write only the name); and callingyarn logsbefore the application finishes. - Official documentation: Apache Hadoop YARN · Capacity Scheduler · ResourceManager REST APIs · yarn-default.xml
What the cluster can offer
Turn on YARN with lab-hadoop start yarn, and save the response JSON of curl -s http://localhost:8088/ws/v1/cluster/metrics to /root/hdp/yarn/metrics.json.
The totalMB, totalVirtualCores and activeNodes of clusterMetrics are the share the NodeManager announced. Think about why totalMB is much smaller than the Pod's 2Gi: the four HDFS and YARN daemons and the client JVM take their places first.
A container larger than the maximum allocation
Run yarn jar /opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.5.0.jar pi -D mapreduce.map.memory.mb=2048 1 10 and let the job end, then save the output of yarn application -status application_<…> (the placeholder is the application ID) for that application to /root/hdp/yarn/toobig.txt.
One map asks for 2048MB while the maximum allocation is 1024MB. The AM ends the job as KILLED before even requesting a map container. The application ID is the same number as the job ID, and you can find it by looking in the history folder for the KILLED job named QuasiMonteCarlo. The Diagnostics line tells the reason.
Create three queues
In /opt/hadoop/etc/hadoop/capacity-scheduler.xml, set yarn.scheduler.capacity.root.queues to default,etl,adhoc, root.default.capacity, root.etl.capacity and root.adhoc.capacity to 20, 50 and 30, and root.adhoc.maximum-capacity to 50. After it is reread with yarn rmadmin -refreshQueues, save the output of yarn queue -status etl to /root/hdp/yarn/etl.txt.
The sum of the guaranteed shares (capacity) under the same parent must be 100 for refreshQueues to accept it. The upper limit (maximum-capacity) is the limit up to which it can borrow spare resources. Just editing the configuration file does nothing: it has to be reread.
Send a job to the etl queue
Upload /data/books/book-1.txt to HDFS /user/root/yarn/in/ and run yarn jar <예제 jar> wordcount -D mapreduce.job.queuename=etl /user/root/yarn/in /user/root/yarn/wc-etl (the placeholder is the example jar).
For the queue, write only the name (without root.). See whether root.etl is printed in the queue column of the job history file name. The grader looks at that column and the _SUCCESS of the output.
A nonexistent queue is rejected at submission
Try submitting the same wordcount with -D mapreduce.job.queuename=nope to /user/root/yarn/wc-nope, and save the error output (including standard error) to /root/hdp/yarn/badqueue.txt.
A nonexistent queue is rejected by the scheduler at the moment it receives the application. The job does not even start, so it is not left in the history either: that is why this error exists only in the submitter's output.
The queues as the scheduler sees them
Save the response JSON of curl -s http://localhost:8088/ws/v1/cluster/scheduler to /root/hdp/yarn/scheduler.json.
In the response's scheduler.schedulerInfo.queues.queue, each queue has capacity, maxCapacity and usedCapacity. The grader checks whether those values equal the configuration file you edited: it is a way of confirming that refreshQueues really took effect.
Gather and read the logs of a finished job
Run yarn logs -applicationId <ID> with the application ID (application_<…>) of the step 4 job and save the output to /root/hdp/yarn/etl-logs.txt.
When a container ends, the NodeManager gathers its logs under /tmp/logs/<사용자>/ in HDFS (log aggregation; the placeholder is the user). So the logs remain even if the node disappears. The gathering takes a few seconds after the app ends, so if you call it right away, the output may be empty. In the output, for each container, a Container: header and LogType:syslog, stdout and stderr come out in turn.
Record the resource numbers
In /root/hdp/yarn/report.md, write three sections: ## 컨테이너 크기, ## 대기열 and ## 로그 (use exactly these Korean headings in this order; they mean "Container size", "Queues" and "Logs"). Put the totalMB from step 1 and the 2048 requested in step 2 in the first section, as numbers.
Write how the node share, minimum allocation, maximum allocation and AM size mesh together, what a queue's guaranteed share and upper limit mean, and where the logs were gathered.