I had already built a small HDFS-style distributed file system in Go. (Part I) It has a controller for metadata, storage nodes that hold replicated chunks, and a client. All three talk over TCP with Protocol Buffers. This post is about the next step: running MapReduce jobs on the data stored in it.
Here's the design goal. The data is already split into chunks and spread across storage nodes, so the computation should go to the chunks, not the other way around.
This is a course project, so the repository is private. If you’re interested, please comment below, and I’ll grant you viewer access.
The pieces
MapReduce adds one new binary and extends two existing ones:
-
Computation manager (new): accepts jobs, schedules map and reduce tasks, tracks progress, and merges the final output. It's a separate process from the controller and asks the controller for chunk placements with the same
GetFileRequestthe client uses. - Storage nodes (extended): they now also run map tasks on their local chunks, receive shuffle data, and run reduce tasks.
-
Client (extended): new
submitjobandjobstatuscommands.
All of this uses the same framing as the file system: a length-prefixed protobuf Wrapper, with new message types like SubmitJobRequest, MapTaskRequest, ShuffleData, ReduceTaskRequest, and GetJobStatusResponse.
Jobs are just executables. A mapper reads lines from stdin and writes key<TAB>value lines to stdout. A reducer reads sorted key<TAB>value lines and writes its results. The repo has three example jobs: word count, a log analyzer that counts log lines per domain, and a second-stage report job that reads that output and produces the top 10 domains. The mapper binary has to exist at the same path on every storage node, because each node runs it with exec.Command(req.PluginPath).
Step 0: chunks that don't split lines
The original file system split files at fixed byte offsets. That works for storing files but breaks text processing: a chunk boundary can land in the middle of a line, and then no mapper sees that line whole.
So the client now checks the file type before chunking. It combines http.DetectContentType on the first 8 KB with the file extension. Text-like types (text/*, JSON, CSV, YAML, and a few others) go through planLineChunks, which only cuts between lines. A single line longer than the chunk size gets a chunk of its own. Everything else is still split at fixed byte offsets.
Each text chunk also records the line number it starts at. That number travels through StoreFileRequest.chunk_start_lines into ChunkPlacement.start_line. When a node runs a mapper, formatMapInput prefixes every line with its line number in the whole file (lineNumber<TAB>line). So a mapper sees the same line numbers it would see if it read the whole file in one pass.
I also sped up uploads while I was changing the client: storeChunksConcurrently now uploads with a pool of 8 workers (storeWorkers) instead of sending one chunk at a time.
Submitting a job
The client sends SubmitJobRequest with the mapper path, the input file, the number of reducers, and the reducer path. The manager fetches the placements for the input file, creates a JobState, assigns an ID, and starts dispatching. A job moves through queued, running, reducing (or merging for map-only jobs), and ends as complete or failed. The client polls with jobstatus.
Reduce tasks are dispatched before map tasks, so each reducer is already waiting when mapper output starts arriving. selectReducerNodes assigns reducers round-robin over the primary node of each chunk.
For each map task, pickNode picks the replica of that chunk with the fewest in-flight map tasks (tracked in nodeActiveTasks). Each map task runs on a node that already has the chunk on local disk, so no input data moves over the network.
The map side: spill, sort, merge
On the storage node, executeMapTask reads the chunk from local disk, verifies its SHA-256 checksum, formats the input, and starts the mapper process. Mapper output can be much larger than memory, so it isn't collected in a slice. spillMapperOutput reads stdout line by line, buffers up to mapSpillMaxLines (50,000) lines, sorts them, and writes each batch to a spill_NNNN.run file. Every line must contain a tab, or the task fails.
When the mapper exits, the sorted runs are merged with a k-way merge over a min-heap. The heap orders by key, then by the full line:
for h.Len() > 0 {
item := heap.Pop(h).(lineHeapItem)
if err := emit(item.reader.line); err != nil {
return err
}
if err := item.reader.advance(); err != nil {
if err == io.EOF {
continue
}
return err
}
heap.Push(h, item)
}
emit does two things at once. It writes the line to the node's local map output file, and it partitions the pair to a reducer by hashing the key with FNV-1a:
emit := func(line string) error {
if _, err := writer.WriteString(line); err != nil {
return err
}
if err := writer.WriteByte('\n'); err != nil {
return err
}
if len(senders) == 0 {
return nil
}
key, value, err := parseKVLine(line)
if err != nil {
return err
}
reducerID := int(fnv32a(key) % uint32(len(senders)))
return senders[reducerID].append(key, value)
}
Each shuffleSender holds one open TCP connection to its reducer's node. It batches pairs and sends a ShuffleData message every shuffleBatchMaxKVs (2,048) pairs. At the end it sends a final message with done: true. The merge output is sorted, so every batch a reducer receives is already sorted too.
If a job has zero reducers, there's no shuffle. Each map task uploads its output to the DFS as __mr_tmp_<job>_map_chunk_<i>.out, and the manager concatenates those files in chunk order.
The reduce side: waiting for every mapper
This part has a race to handle. A reducer needs to know how many mappers to wait for, and that number comes in ReduceTaskRequest.num_mappers. But shuffle data from a fast mapper can arrive at the node before the ReduceTaskRequest does.
reducerBuffer handles both orders. If shuffle data arrives first, the buffer is created with numMappers = -1 (unknown). Each non-empty batch is written straight to disk as its own sorted run file, and each done message increments a counter:
func (s *StorageNode) acceptShuffleData(sd *pb.ShuffleData) {
buf := s.getOrCreateReducerBuffer(sd.JobId, sd.ReducerId, -1)
buf.mu.Lock()
defer buf.mu.Unlock()
if len(sd.Pairs) > 0 {
runPath := s.reducerRunPath(sd.JobId, sd.ReducerId, buf.nextRunID)
buf.nextRunID++
if err := s.writeShuffleRun(runPath, sd.Pairs); err != nil {
if buf.err == nil {
buf.err = err
}
} else {
buf.runFiles = append(buf.runFiles, runPath)
}
}
if sd.Done {
buf.doneCount++
if buf.numMappers >= 0 && buf.doneCount >= buf.numMappers {
buf.once.Do(func() { close(buf.ready) })
}
}
}
When the ReduceTaskRequest arrives, getOrCreateReducerBuffer fills in numMappers and runs the same check, so the ready channel closes whichever message comes last. sync.Once makes sure it only closes once. executeReduceTask blocks on <-buf.ready.
After that, the reducer reuses mergeSortedRuns from the map side and streams the merged lines directly into the reducer process's stdin. The reducer sees every value for a key together, in sorted order, without the node ever loading the whole partition into memory. The reducer's stdout is uploaded to the DFS as __mr_tmp_<job>_reducer_<id>.out. Once every reducer reports success, the manager concatenates the partitions in reducer-ID order into <job>_output.out.
Handling failed and hung map tasks
A map task can fail in two ways, and the manager handles each one differently.
It reports failure. The node sends back MapTaskResponse{success: false}. The manager moves to the next replica in that chunk's placement and dispatches there right away. Once no replicas are left, the job fails.
It never answers. The node accepts the task and then dies or hangs. watchJobTimeout checks every 10 seconds for chunks that have been running longer than the timeout and re-dispatches them to the next replica. The timeout adapts to how fast tasks are actually finishing:
timeout := fallbackTimeout
if job.MapsTotal > 0 && float64(job.MapsDone)/float64(job.MapsTotal) >= threshold && len(job.completionTimes) > 0 {
var total time.Duration
for _, d := range job.completionTimes {
total += d
}
avg := total / time.Duration(len(job.completionTimes))
adaptiveTimeout := time.Duration(adaptiveMultiplier) * avg
if adaptiveTimeout < minAdaptiveTimeout {
adaptiveTimeout = minAdaptiveTimeout
}
if adaptiveTimeout < timeout {
timeout = adaptiveTimeout
}
}
Until 70% of map tasks have finished, the timeout is the fallback of 15 minutes. After that, it's twice the average completion time, but never less than 3 minutes and never more than 15. Early in a job there's no data about how long a task normally takes, so a fixed 3-minute timeout would kill slow but healthy tasks. With enough completions to estimate from, a stuck task gets retried much sooner than the 15-minute fallback. The most recent change to this code made the adaptation more conservative.
If the original task and its retry both finish, completedChunks counts the chunk only once.
What I'd change next
Some limitations I know about:
-
Shuffle data isn't tied to a map attempt.
ShuffleDatacarriesjob_idandreducer_id, but no chunk index or attempt number. If a map task sends part of its shuffle and then fails, or if a timed-out task and its retry both finish, the reducer can receive duplicate pairs, and the extradonemessages can make it start before the real last mapper finishes. Tagging each stream with its chunk index, and countingdonemessages per chunk, would fix this. -
Retries can go back to the same node. Retries walk the placement list starting at index 1, whichever replica
pickNodeactually chose. IfpickNodepickedNodes[1]and that task fails, the first retry goes to the same node again. - Reducers aren't retried. One failed reduce task fails the whole job.
- The final merge happens in memory. The manager downloads every partition into one byte slice and uploads the result as a single chunk with replication factor 1. That limits output size to the manager's memory, and the output has no replicas.
-
Nothing gets cleaned up. Spill files, shuffle runs,
__mr_tmp_files in the DFS, and the per-job reducer buffers all stay around after a job finishes. - Job state lives in the manager's memory only. If the manager restarts, every job in flight is lost.
The ideas I learned the most from were external sorting with spill files and a heap merge, and handling shuffle data that shows up before the reduce task does. Scheduling map tasks on the nodes that already hold the chunk was the easy part, because the file system already tracked where every chunk lives.
Top comments (0)