Apache Hadoop — Stand up and run HDFS and YARN in one pod
The cost of MapReduce is concentrated in the shuffle between map and reduce
In one line
MapReduce is the framework map → (partition, sort, shuffle) → reduce, and what the user writes is only the two functions at the ends. The sorting and shuffling in the middle are done by the framework, and most of a job's cost comes from there. A combiner is a device that shrinks that middle, and counters are the window for looking into it.
Why this shape was needed
Suppose you have to count words in terabytes of data scattered over hundreds of machines. If you gather the data on one server, the network dies first. That is why the principle set out by the HDFS design document is "moving computation near the data is cheaper than moving the data". Then you have to divide the computation and run it on each server where the blocks lie, and this time you have to gather the divided results again by the same words. If a server dies, its share has to be rerun, and a slow server has to be waited for.
If you solve this problem anew for every job, most of the code becomes distributed-processing housekeeping. MapReduce moved that housekeeping into the framework. The MapReduce tutorial summarizes that the framework takes care of placing and monitoring tasks and of rerunning failed tasks, and the application supplies only the input and output locations and the map and reduce functions. The price is the constraint that every computation must be expressed as key-value pairs. Thanks to that constraint, one promise, "the same key goes to the same reducer", solves the gathering.
How it works
Written in the shape of pairs, the flow is this, exactly as in the tutorial.
(입력) <k1, v1> -> map -> <k2, v2> -> combine -> <k2, v2> -> reduce -> <k3, v3> (출력)
The map side. According to the tutorial, the number of maps is usually governed by the total size of the input, that is, the number of blocks of the input files. One block usually becomes one input split, and one map task starts per split. The pairs a map emits do not go straight to disk but pile up in a memory buffer. In mapred-default.xml, the default of this buffer size mapreduce.task.io.sort.mb is 100MB (the lab image, whose task heap is only 180MB, reduced it to 32MB), and when it reaches the default 0.80 of mapreduce.map.sort.spill.percent, a background thread sorts the contents and spills them to disk. When the map finishes, the spilled pieces are merged into one. This file is divided into as many segments as there are reducers, and inside each segment it is sorted by key.
Which key goes to which segment is decided by the Partitioner. The default is HashPartitioner, the remainder of the key's hash divided by the number of reducers. The same key, whichever map it came from, goes to the reducer with the same number. This one line is the basis on which the gathering holds.
The reduce side. The tutorial explains that a reducer goes through three phases: shuffle, sort and reduce. In the shuffle, it fetches only its own segment from all the map outputs over HTTP, and while fetching, it does a merge sort to line up the values of the same key. Then it calls reduce(키, 값 목록) (key, list of values) once per key. If you set the number of reducers to 0, the map output becomes the result as is, with no sorting.
The default of the number of reducers mapreduce.job.reduces is 1. If not set, the whole map output is crowded into one reducer and the result file is one part-r-00000. If you raise the reducers to three, the result files become three. Reducers are called in key order, so for a job that uses the received keys as they are, like word counting, each file is in key order inside. However, as the tutorial states firmly, the framework does not re-sort the reducers' output, and if you concatenate the three files, the overall order is mixed. This is because it was divided by hash.
To sort the whole thing into one line, you have to change the partitioner. The TeraSort example shows how. It samples keys from the input to decide N-1 boundaries and assigns one reducer to each boundary interval. Then the output of reducer i is all smaller than that of reducer i+1, and if you concatenate the files in numeric order, the whole is sorted. TeraValidate is the job that checks this.
Combiners and counters
The map of word counting emits (the, 1) tens of thousands of times. If you shuffle that as is, a 1 crosses the network tens of thousands of times. A combiner is a local aggregation that reduces it in advance on the map side and sends (the, 38211). The tutorial says the combiner runs at spill time, may also run during the reduce-side merge, and that whether a record larger than the buffer goes through the combiner is not defined. To summarize, how many times the combiner runs is not guaranteed. So only operations whose result is the same whether they run once or three times can be used as a combiner. Sums and maximums are fine, and averages are not: the average of averages is not the average.
Whether the combiner does its job is seen with counters. When the job finishes, the values counted by the framework are printed. They are things like map input records and map output records, combiner input and output records, reduce input groups and spilled records. If the combiner output is far smaller than its input, it means the shuffle shrank by that much. For word counting, the number of reduce input groups is the number of distinct words in the result. Counters are the cheapest way to prove in numbers what a job did without digging through logs.
Running on YARN
Since Hadoop 2, MapReduce does not manage resources by itself. As the YARN architecture document explains, the roles are divided into the ResourceManager, which arbitrates the whole cluster, and an ApplicationMaster started one per application, and the ApplicationMaster of a MapReduce job is MRAppMaster. When you submit one job, containers start: at least one AM, as many as the number of map tasks, and as many as the number of reducers.
When the job finishes, the AM leaves a history. The default staging directory is /tmp/hadoop-yarn/staging and the intermediate done directory is history/done_intermediate beneath it. Even without starting a job history server, .jhist (the event history), .summary and the job configuration XML pile up here per user. The history server is only in the role of moving and displaying them; the recording itself is done by the AM.
What it looks like in the field
First, the job runs locally, not on the cluster. The default of mapreduce.framework.name is local. A job submitted from a client that has not been changed to yarn does not appear in the ResourceManager and quietly runs in one JVM on the submitting machine. It is the most common answer to "why is there no job on the ResourceManager screen". The lab image has changed this value to yarn.
Second, there is no room for the AM. The default of yarn.app.mapreduce.am.resource.mb is 1536MB. On a cluster where the share a NodeManager gives to containers is smaller than that, there is no node where the AM container can fit, so the job cannot even start. In a place where the share is 1024MB, like this lab environment, you must lower the AM size separately for the job to start. The lab image set the AM to 384MB and the map and reduce to 256MB, so that one AM and two tasks fit at the same time.
Third, one reducer never finishes. Either the default of 1 was left as it is, or the values are concentrated on one key. The former is fixed with configuration, and the latter with key design.
Fourth, the numbers go wrong because of the combiner. If you apply an average or a median as a combiner, using the reducer code as it is, the result differs depending on how many times the combiner runs. It is a bug that is right sometimes and wrong at other times, so it is hard to find.
What really matters in practice
- The same key goes to the same reducer. The basis of the gathering is a line of the partitioner.
- The shuffle is the cost. Most of tuning is reducing the bytes that cross between map and reduce.
- You do not know how many times a combiner runs. Use only operations that give the same result even when applied several times.
- The default number of reducers is 1. Both the number of result files and the parallelism depend on this value.
- The framework.name default is local. If a job is not on the ResourceManager, look at this first.
- Prove it with counters. Input, output, combiner and group counts show in numbers what the job did.
What you will do in the next lab
You turn on YARN and confirm that the NodeManager is registered, then upload five books to HDFS, run the example jar's word count job and write down the job ID. Without a history server, you find and open the job history file left in done_intermediate in HDFS (the lab image has set this file to be written as human-readable JSON lines). From the full counters at the end of the job, you pull out map input and output, combiner input and output, reduce input groups and output records, and calculate how much the combiner reduced the shuffle. You raise the reducers to three and see how the result files are divided, and finally run a total sort with teragen, terasort and teravalidate and pass it through validation.