Apache Hadoop — Stand up and run HDFS and YARN in one pod
Run word count and TeraSort on YARN and read them through counters
Goal
You turn on YARN and count the words of five books with MapReduce. After the job finishes, you open the job history (.jhist) left in HDFS to read the counters, and confirm with TeraSort how the output is divided when you change the number of reducers and how the shuffle guarantees sorting.
Why it matters
MapReduce is rarely written newly these days, but its model (the map emits (key, value), the framework gathers by key and sorts and hands over (the shuffle), and the reduce combines per key) has been inherited by the shuffle of every Spark, Flink and SQL engine. Here you can see most plainly why a shuffle is expensive, why a combiner (pre-combining on the map side) reduces the shuffle, and why the number of reducers becomes the number of output files. Counters are a record of what the job did, left as numbers: how many lines the map read and how many it emitted, how many lines the combiner reduced them to, how many keys the reduce received. It is the first place to look when the result looks strange, and it remains in the history file even after the job has finished. The YARN in this Pod keeps containers small to fit the Pod memory (2Gi) (AM 384MB, map and reduce 256MB). So tasks run one or two at a time in turn: it is slow, but the results and counters are the same.
Steps
- Turn on YARN with
lab-hadoop start yarnand save the output ofyarn node -list -allto /root/hdp/mr/nodes.txt. - Upload the five files
/data/books/book-*.txtto HDFS /user/root/mr/in/. - Count /user/root/mr/in with the example jar's
wordcountand write the output to /user/root/mr/wc. - Write the ID (
job_…) of the step 3 job to /root/hdp/mr/jobid.txt. - From that job's history file (.jhist), pull out the TaskCounter values
MAP_INPUT_RECORDS,MAP_OUTPUT_RECORDS,COMBINE_INPUT_RECORDS,COMBINE_OUTPUT_RECORDS,REDUCE_INPUT_GROUPSandREDUCE_OUTPUT_RECORDSand write them to /root/hdp/mr/counters.json. - Run the same wordcount with three reducers (
-D mapreduce.job.reduces=3) and write to /user/root/mr/wc3. - Make 100,000 lines with
terageninto /user/root/mr/tera-in, sort them withterasortinto /user/root/mr/tera-out, and then write the validation result withteravalidateto /user/root/mr/tera-val. - In /root/hdp/mr/report.md, write three sections:
## 맵과 리듀스,## 카운터and## 정렬(use exactly these Korean headings in this order; they mean "Map and reduce", "Counters" and "Sorting"). Put theMAP_OUTPUT_RECORDSandCOMBINE_OUTPUT_RECORDSfrom step 5 in the second section.
Notes
- Example jar:
/opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.5.0.jar. Run it asyarn jar <jar> wordcount <입력> <출력>(the placeholders are the jar, the input and the output). If the output directory already exists, the job does not start. - Job history: in
/tmp/hadoop-yarn/staging/history/done_intermediate/root/there are a<잡ID>-…-<잡 이름>-…-<상태>-<대기열>-….jhist(job ID, job name, state and queue) and the.summaryand_conf.xml. After its first two lines (format and schema), the .jhist has events as one JSON line each, and the counters are in thetotalCountersof theJOB_FINISHEDevent. - One job takes about a minute. Keeping YARN on uses nearly 1GB more memory, so do not do other heavy work at the same time.
- Common mistakes: rerunning without deleting the output directory; and writing the application ID (
application_…) instead of the job ID (the numbers are the same but the names differ). - Official documentation: MapReduce Tutorial · Apache Hadoop YARN · mapred-default.xml · YARN Commands
Turn on YARN
Turn on the ResourceManager and NodeManager with lab-hadoop start yarn, and save the output of yarn node -list -all to /root/hdp/mr/nodes.txt.
The NodeManager must tell the ResourceManager how much memory and how many cores it can offer before it can accept jobs. The node should show as RUNNING. Check with lab-hadoop status that the four JVMs (NameNode, DataNode, ResourceManager, NodeManager) are up.
Upload the input
Upload /data/books/book-1.txt to book-5.txt to HDFS /user/root/mr/in/.
Each file in the input directory produces an input split, and one map task attaches to one split. Five files are five maps.
Run wordcount
Count the words with yarn jar /opt/hadoop/share/hadoop/mapreduce/hadoop-mapreduce-examples-3.5.0.jar wordcount /user/root/mr/in /user/root/mr/wc.
When the job finishes, _SUCCESS and part-r-00000 appear in the output directory. The map emits (word, 1) for each word, the combiner adds first on the map side, and the reducer sums per word. The grader compares your output with the values counted directly from the original.
Find the job history
Write the ID of the step 3 job (job_<숫자>_<숫자>, where each placeholder is a number) as one line to /root/hdp/mr/jobid.txt.
When you run the job, 'Running job: job_…' is printed on the console. If you missed it, look in the listing of HDFS /tmp/hadoop-yarn/staging/history/done_intermediate/root/ for a .jhist whose name contains word+count: even without a history server, the AM leaves it here.
Read the counters
From the .jhist of the step 4 job, pull out from the org.apache.hadoop.mapreduce.TaskCounter group MAP_INPUT_RECORDS, MAP_OUTPUT_RECORDS, COMBINE_INPUT_RECORDS, COMBINE_OUTPUT_RECORDS, REDUCE_INPUT_GROUPS and REDUCE_OUTPUT_RECORDS and write them to /root/hdp/mr/counters.json as {"이름": 정수} (each counter name maps to an integer).
MAP_OUTPUT_RECORDS is the number of words (including duplicates), and COMBINE_OUTPUT_RECORDS is the sum over the maps of their distinct words: the number of lines that actually passed to the shuffle. REDUCE_INPUT_GROUPS should equal the total number of distinct words. The ratio of the two numbers is the shuffle the combiner saved.
Three reducers mean three outputs
Run a job with three reducers using yarn jar <jar> wordcount -D mapreduce.job.reduces=3 /user/root/mr/in /user/root/mr/wc3 (the placeholder is the jar path).
The map output is assigned by the key's hash to one of the three reducers (HashPartitioner). So a word appears in exactly one file, and each file is sorted inside itself but the files are not sorted relative to each other.
TeraSort: the shuffle creates the sort
Make 100,000 lines (100 bytes each) with teragen 100000 /user/root/mr/tera-in, sort them with terasort /user/root/mr/tera-in /user/root/mr/tera-out, and then validate with teravalidate /user/root/mr/tera-out /user/root/mr/tera-val (all from the same example jar).
TeraSort samples the input, cuts the key range into as many pieces as there are reducers (total order partitioning), and each reducer sorts its own range and writes it out. Then the output files concatenated in order become the total sort. TeraValidate writes 'misorder' if there is a place that is out of order and only a checksum if there is none.
Record the shuffle in numbers
In /root/hdp/mr/report.md, write three sections: ## 맵과 리듀스, ## 카운터 and ## 정렬 (use exactly these Korean headings in this order; they mean "Map and reduce", "Counters" and "Sorting"). Put the values of MAP_OUTPUT_RECORDS and COMBINE_OUTPUT_RECORDS from step 5 in the second section as numbers.
In the first section, write how many maps and reduces there were and how the output was divided; in the second, by what fraction the combiner reduced the shuffle; and in the third, how TeraSort produced the total order.