TT Lab
Get started
Learn Learning paths Courses

Apache Hadoop — Stand up and run HDFS and YARN in one pod

Run word count and TeraSort on YARN and read them through counters

Continue in TT Lab

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

  1. Turn on YARN with lab-hadoop start yarn and save the output of yarn node -list -all to /root/hdp/mr/nodes.txt.
  2. Upload the five files /data/books/book-*.txt to HDFS /user/root/mr/in/.
  3. Count /user/root/mr/in with the example jar's wordcount and write the output to /user/root/mr/wc.
  4. Write the ID (job_…) of the step 3 job to /root/hdp/mr/jobid.txt.
  5. 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_GROUPS and REDUCE_OUTPUT_RECORDS and write them to /root/hdp/mr/counters.json.
  6. Run the same wordcount with three reducers (-D mapreduce.job.reduces=3) and write to /user/root/mr/wc3.
  7. Make 100,000 lines with teragen into /user/root/mr/tera-in, sort them with terasort into /user/root/mr/tera-out, and then write the validation result with teravalidate to /user/root/mr/tera-val.
  8. 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 MAP_OUTPUT_RECORDS and COMBINE_OUTPUT_RECORDS from step 5 in the second section.

Notes

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.