TT Lab
Get started
Learn Learning paths Courses

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

Aggregate access logs with MapReduce written in Python

Continue in TT Lab

Goal

Using Hadoop Streaming, you write a Python mapper and reducer and count three days of access logs by status code. You measure with counters how much the combiner reduces the shuffle, and cover partitioning that chooses the reducer by date (KeyFieldBasedPartitioner), a user counter that counts broken lines, and task failure and retries with a mapper that deliberately dies.

Why it matters

Streaming turns map and reduce into a standard input/output contract. The mapper receives input lines on standard input and emits 'keyvalue' lines on standard output, and the framework sorts them by key and streams them to the reducer's standard input. What the reducer receives is not a bundle per key but a stream of sorted lines, so it has to find where the key changes by itself. As long as you keep this contract, you can write MapReduce in any language, and you can test it locally in the same way with cat | mapper | sort | reducer. A combiner is a reducer that combines ahead of time on the map side. For the result not to change, the operation must satisfy associativity and commutativity (sums and maximums are fine, averages are not), and the shape of input and output must be the same. If you attach it correctly, the shuffle shrinks to a few tens of times smaller. The default partitioning is the hash of the whole key. If the key is 'datepath' and you want to see all the paths of one date in one reducer, you have to change it to partition by only the first column of the key. And the most common reason a streaming job in production is silently wrong is broken lines: discarding them is fine, but not counting how many lines you discarded is not.

Steps

  1. Write /root/hdp/streaming/mapper.py. Split a line on whitespace, and if the line has 10 or more columns, the ninth column (index 8, counting from 0) is a three-digit number and it does not start with #, emit 상태코드<TAB>1 (a status code, TAB, 1); otherwise discard it.
  2. Write /root/hdp/streaming/reducer.py. It receives 키<TAB>수 lines sorted by key (key, TAB, count), adds the counts of the same key and emits 키<TAB>합 (key, TAB, sum); the count need not be 1.
  3. Upload /data/logs/access-2026-03-01.log to 03.log to HDFS /user/root/streaming/in/, run the streaming job with the job name hdp-status, and write to /user/root/streaming/status.
  4. Run the same job with the reducer attached as -combiner, with the job name hdp-status-comb and the output /user/root/streaming/status_comb, and write MAP_OUTPUT_RECORDS, COMBINE_OUTPUT_RECORDS and REDUCE_INPUT_RECORDS of that job to /root/hdp/streaming/shuffle.json as {"map_output": …, "combine_output": …, "reduce_input": …}.
  5. Make /root/hdp/streaming/mapper2.py emit 날짜(yyyy-MM-dd)<TAB>경로<TAB>1 (date, TAB, path, TAB, 1), use the reducer.py from step 1 as it is as the reducer, and run it with the job name hdp-daily, 2 reducers, a two-column key and partitioning by the first column only, writing to /user/root/streaming/daily.
  6. In the step 1 mapper, add a line that raises the user counter Lab·BadLines by 1 for each broken line (write reporter:counter:Lab,BadLines,1 to standard error), and run it with the job name hdp-bad to /user/root/streaming/bad.
  7. With /root/hdp/streaming/crash_mapper.py, which ends with an exception when it meets the path /event/spring, run it with the job name hdp-crash and one map attempt (-D mapreduce.map.maxattempts=1) to /user/root/streaming/crash to make it fail, write that job ID to /root/hdp/streaming/crash.txt, and then run it again with the same name using the step 1 mapper to /user/root/streaming/crash-ok to make it succeed.
  8. In /root/hdp/streaming/report.md, write three sections: ## 표준 입출력 계약, ## 컴바이너와 분할 and ## 실패와 카운터 (use exactly these Korean headings in this order; they mean "The standard I/O contract", "Combiner and partitioning" and "Failures and counters"). Put map_output and reduce_input from step 4 in the second section.

Notes

The mapper: take a line and emit a key and a value

Write /root/hdp/streaming/mapper.py. For each line of standard input, split it on whitespace, and if it has 10 or more columns, the 9th column (8 counting from 0) is a three-digit number and the line does not start with #, write 상태코드<TAB>1 (a status code, TAB, 1) to standard output; otherwise emit nothing.

Test first with head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | uniq -c. The grader runs your mapper directly on new input mixed with broken lines.

The reducer: find where the key changes in a sorted stream

Write /root/hdp/streaming/reducer.py. It receives on standard input 키<TAB>수 lines (key, TAB, count) sorted in key order, adds the counts of the same key, and writes 키<TAB>합 (key, TAB, sum) to standard output. The count may not be 1.

The reducer does not receive a list per key. Sorted lines just flow in, so the moment the key differs from the previous line, emit the sum of the previous key and start over. Do not forget to emit the last key. If you write it so that the count need not be 1, it is directly usable as a combiner.

Put it on YARN

Turn on YARN with lab-hadoop start yarn, upload /data/logs/access-2026-03-01.log to 03.log to HDFS /user/root/streaming/in/, and then run mapred streaming -D mapreduce.job.name=hdp-status -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input /user/root/streaming/in -output /user/root/streaming/status.

Scripts given with -files are copied into each task's working folder. So for -mapper you write only the file name without a path. The grader checks whether the output equals the values counted directly from the original logs and whether there is a successful hdp-status job in the history.

Reduce the shuffle with a combiner

Add -combiner 'python3 reducer.py' to the step 3 job and run it with the job name hdp-status-comb and the output /user/root/streaming/status_comb. From that job's history, pull out TaskCounter MAP_OUTPUT_RECORDS, COMBINE_OUTPUT_RECORDS and REDUCE_INPUT_RECORDS and write them to /root/hdp/streaming/shuffle.json as {"map_output": 정수, "combine_output": 정수, "reduce_input": 정수} (all three values are integers).

A sum satisfies associativity and commutativity, so using the reducer as it is as the combiner gives the same result. See how much smaller the number of lines passed to the shuffle after the combiner (= the reduce input) is than the map output.

Choose the reducer by date

Make /root/hdp/streaming/mapper2.py emit 날짜(yyyy-MM-dd)<TAB>경로<TAB>1 (date, TAB, path, TAB, 1) for each valid line, use reducer.py as the reducer, and write to /user/root/streaming/daily with the job name hdp-daily, -D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1 (generic options, placed at the front) and -partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2 (command options, placed after).

The key is the first two columns (date and path), so sorting is by two columns, but partitioning is only by the first column (the date). So all the paths of one date go to one reducer, and one output file holds exactly one date. If you wrote reducer.py to split at the last tab, it works as it is even when the key has two columns.

Count broken lines with a user counter

Make mapper.py write reporter:counter:Lab,BadLines,1 to standard error for each broken line, and run it with the job name hdp-bad to /user/root/streaming/bad. The counter Lab·BadLines of that job must equal the number of broken lines in the three logs.

Streaming understands a reporter:counter:<그룹>,<이름>,<증가량> line on standard error (group, name, increment) as a counter update. This counter remains in the job history, so even after the job has finished you can answer "how many lines did we discard".

A mapper that dies: failure and rerunning

Write /root/hdp/streaming/crash_mapper.py, which ends with an exception when it meets the path /event/spring, run it with the job name hdp-crash and -D mapreduce.map.maxattempts=1 to /user/root/streaming/crash to make it fail, and write that job ID to /root/hdp/streaming/crash.txt. Then run the step 1 mapper.py again under the same job name to /user/root/streaming/crash-ok to make it succeed.

When a map task fails, the framework reruns it as another attempt (4 times by default). If you reduce the attempts to 1, the first failure is the job's failure. A failed job leaves a history too, so to find what died and why, read the task's standard error with yarn logs -applicationId application_<같은 숫자> (where the placeholder is the same number as the job ID).

Record the contract, the combiner and the failure

In /root/hdp/streaming/report.md, write three sections: ## 표준 입출력 계약, ## 컴바이너와 분할 and ## 실패와 카운터 (use exactly these Korean headings in this order; they mean "The standard I/O contract", "Combiner and partitioning" and "Failures and counters"). Put map_output and reduce_input from step 4 in the second section as numbers.

Write what the reducer receives, by what fraction the combiner reduced the shuffle, and how you noticed the broken lines and the failed task.