Apache Hadoop — Stand up and run HDFS and YARN in one pod
Aggregate access logs with MapReduce written in Python
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
- 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. - 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. - Upload
/data/logs/access-2026-03-01.logto03.logto HDFS /user/root/streaming/in/, run the streaming job with the job namehdp-status, and write to /user/root/streaming/status. - Run the same job with the reducer attached as
-combiner, with the job namehdp-status-comband the output /user/root/streaming/status_comb, and writeMAP_OUTPUT_RECORDS,COMBINE_OUTPUT_RECORDSandREDUCE_INPUT_RECORDSof that job to /root/hdp/streaming/shuffle.json as{"map_output": …, "combine_output": …, "reduce_input": …}. - 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 namehdp-daily, 2 reducers, a two-column key and partitioning by the first column only, writing to /user/root/streaming/daily. - In the step 1 mapper, add a line that raises the user counter
Lab·BadLinesby 1 for each broken line (writereporter:counter:Lab,BadLines,1to standard error), and run it with the job namehdp-badto /user/root/streaming/bad. - With /root/hdp/streaming/crash_mapper.py, which ends with an exception when it meets the path
/event/spring, run it with the job namehdp-crashand 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. - 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"). Putmap_outputandreduce_inputfrom step 4 in the second section.
Notes
- Streaming job:
mapred streaming -D mapreduce.job.name=<이름> -files mapper.py,reducer.py -mapper 'python3 mapper.py' -reducer 'python3 reducer.py' -input <입력> -output <출력>(the placeholders are the job name, the input and the output). Put the-Doption before the other options. - Local test:
head -1000 /data/logs/access-2026-03-01.log | python3 mapper.py | sort | python3 reducer.py. - Two key columns, partitioned by the first column: put the generic options
-D stream.num.map.output.key.fields=2 -D mapreduce.partition.keypartitioner.options=-k1,1before-files, and the command options-partitioner org.apache.hadoop.mapred.lib.KeyFieldBasedPartitioner -numReduceTasks 2after-mapper. If the order is wrong, it stops with 'Unrecognized option'. - Job history is left in
/tmp/hadoop-yarn/staging/history/done_intermediate/root/(a-in the job name is written as%2Din the file name). Failed jobs are left too. - Common mistakes: thinking the reducer is called only once per key (it is a flow called line by line); making the combiner output a different shape from the input; and rerunning without deleting the output directory.
- Official documentation: Hadoop Streaming · MapReduce Tutorial · mapred-default.xml
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.