Apache Hadoop — Stand up and run HDFS and YARN in one pod
In Hadoop Streaming, one line of standard I/O is the whole contract
In one line
Hadoop Streaming is a tool for plugging in, as the mapper and reducer, any program that reads standard input and writes standard output. The contract is simple: one line is one record, what comes before the first tab is the key, and the reducer receives lines sorted by key. As simple as it is, the script must find the boundaries where the key changes by itself, and failures and counters must also be spoken through exit codes and standard error.
Why this was needed
MapReduce's original contract is a Java interface. You have to extend Mapper, match the Writable types, and bundle into a jar. That much formality is heavy for counting requests per status code from an access log. Moreover, most people who work with data already have a single-machine script in Python or shell.
The Hadoop Streaming documentation fills that gap. Anything that is an executable or a script can be used as a mapper and reducer, and the documentation's first example uses /bin/cat as the mapper and /usr/bin/wc as the reducer. A pipeline that ran on one machine as cat log | map.py | sort | reduce.py spreads out to hundreds of machines almost as it is. The sort in the middle is taken over by MapReduce's shuffle.
How it works
According to the documentation, a map task starts the specified executable as a separate process when it begins. It turns the input split into lines and feeds them into that process's standard input, and gathers the lines that come out of standard output and turns them into key-value pairs. The default rule is what comes before the first tab character is the key and what follows is the value, and if there is no tab, the whole line is the key and the value is empty. The reducer side is the same: the framework unpacks the sorted and merged pairs back into 키\t값 lines (key, tab, value) and feeds them into the reducer process's standard input.
The most important fact in this contract is that what the reducer receives is not a bundle per key but a stream of sorted lines. A Java reducer is called once per key with reduce(키, 값 목록) (key, list of values), but a streaming reducer has to read lines one at a time and judge by itself "did the key change". As the MapReduce tutorial says, the map output is sorted and then divided per reducer, so lines of the same key necessarily come adjacent. The reducer relies on that one guarantee: it remembers the key of the previous line, and at the moment the key changes, it emits the total and resets.
import sys
current, total = None, 0
for line in sys.stdin:
key, _, value = line.rstrip("\n").partition("\t")
if key != current:
if current is not None:
print(f"{current}\t{total}")
current, total = key, 0
total += int(value)
if current is not None:
print(f"{current}\t{total}")
Leaving out the line that emits the last key once more after the loop is the most common mistake. Then, in the result, the one key that comes last in sort order silently disappears.
The mapper and reducer scripts must be on the worker nodes. If you pass them with -files, a symbolic link with the same name is created in the task's working folder. The documentation warns that the command fails unless you put generic options such as -files or -D before the streaming options.
Knobs for handling the shuffle
The combiner. With -combiner, you can give one more executable. For a mapper that emits 1 per status code, you can apply the same summing script as the reducer as a combiner and greatly reduce the number of lines passed to the shuffle. The number of times a combiner runs is not guaranteed, so apply only the summing kind, whose result is the same even when applied several times.
Key fields and partitioning. When the key consists of several fields, like 2026-09-19.404, there are times you want to sort by the whole key but partition by only the first field. The documentation's KeyFieldBasedPartitioner example is that. You give, with stream.num.map.output.key.fields, how many fields the key has, with map.output.key.field.separator the field separator, and with mapreduce.partition.keypartitioner.options=-k1,1 the field to use for partitioning. Then all the lines of the same date go to the same reducer, and inside it they arrive sorted by the whole key. You use it when you want per-date result files, or when you want to compute per-date subtotals inside a reducer.
The number of reducers and map-only jobs. The default number of reducers is the default 1 of mapreduce.job.reduces in mapred-default.xml. According to the documentation, if you set it to 0, the map output becomes the result with no reducer. For work that only filters lines or changes the format, removing the shuffle itself is the fastest.
How to speak about failures and counters
A streaming script has no Java API, so it speaks to the framework only through two channels. One is the exit code. According to the documentation, by default a task that ends with a nonzero exit code is treated as a failure. A failed task is tried again, and because the default of mapreduce.map.maxattempts is 4, if the same map fails four times, the whole job fails. A mapper that throws an exception on one broken line dies four times on the split containing that line and drags the job down.
The other is standard error. If you write a line of the form reporter:counter:<그룹>,<카운터>,<양> (group, counter and amount) to stderr, the counter goes up, and reporter:status:<메시지> (message) changes the status text. So it is better for a broken line to be skipped instead of dying with an exception, while writing reporter:counter:logs,malformed,1. The job runs to the end, and how many lines were broken is left as a number in the counter. If you write it mixed into standard output, it becomes a result record, so always send it to stderr.
Configuration values can also be read as environment variables. According to the documentation, the dots in the configuration name are changed to underscores, so, for example, mapreduce.job.id appears as mapreduce_job_id.
What it looks like in the field
First, it is right locally but the last key is missing on the cluster. It is the case where the reducer did not emit the last group after the loop. When testing with sort on one machine, it sometimes does not show by chance with a small dataset.
Second, the job retries four times and dies. There is a line whose encoding is broken or which has too few columns. If the Python mapper throws an exception at that line, the retries of the same split die at the same line. Retries are a device for transient failures, not for bad data.
Third, the counter is printed in the result file. You wrote reporter:counter to standard output with print. A key starting with reporter:counter:... appears in the result.
Fourth, the error that the script is not on the node. You left out -files or put it after the streaming options. Also check the execute permission and the interpreter declaration on the first line.
What really matters in practice
- The contract is one line. What comes before the first tab is the key, and if there is no tab, the whole line is the key.
- The reducer finds the boundaries by itself. The only guarantee is that the same key comes adjacent. Do not forget the last key.
- Count broken lines instead of dying on them. If you write reporter:counter to stderr, the job finishes and the number remains.
- A nonzero exit code is a failure. If the same map fails 4 times by default, the job ends.
- You can separate partitioning and sorting. With KeyFieldBasedPartitioner, you choose the reducer by the first field only.
- Put generic options first. If -files and -D come after the streaming options, the command fails.
What you will do in the next lab
You make a Python mapper that emits 1 for each status code from an access log and a reducer that reads sorted input continuously and emits the sum, and run them as a streaming job. You apply the same summing script as a combiner and compare the map output and reduce input record counts with counters, and with KeyFieldBasedPartitioner you divide a two-field key, the date and the path joined by a tab, over two reducers by looking only at the first field (the date). You count broken lines by raising a user counter through standard error, run a mapper that deliberately dies on a certain path with one map attempt to see the job fail immediately, and then rerun the job of the same name with a normal mapper to make it succeed.