MapReduce for Dummies
Say you need to grep for "Hi There" across a petabyte of text spread over thousands of files. A single machine reading that much data sequentially could take hours — I/O alone becomes the bottleneck, before you even count the CPU time to scan every line.
The fix: split the work. Hand different files to different machines (workers) so they scan in parallel, and use a central coordinator to hand out tasks, track progress, and reassign anything that fails.
That's MapReduce in one sentence: split, distribute, collect.
You can find the full implementation of this lab in this commit.
Sample Working example
We use mapReduce to perform grep over 100 files using n=10 workers.
Step 1- Define the coordinator and pass on the tasks
c := Coordinator{
mu: sync.Mutex{},
mapTasks: make([]Task, len(files)),
reduceTasks: make([]Task, nReduce),
nMap: len(files), //100
nReduce: nReduce, //10
}
for i, file := range files {
c.mapTasks[i] = Task{
Id: i,
File: file,}
}
Step 2- Initialize workers and run them until task is complete
// rpc call to co-ordinator
for {
task, err := fetchTask()
}
func fetchTask() (*Task, error) {
ok := call("Coordinator.FetchTask", &req, &reply)
if !ok {
return &Task{Type: IdleTask}, fmt.Errorf("failed to fetch task")
}
return &reply.Task, nil
}
Step 3- Perform Map Task (In this case calling grep)
This step contains the crux of MapReduce -
- Tasks should be atomic, tasks failed mid-way are discarded
- Each Map task is divided into N(=10) reduce tasks. So, in total there are map-tasks * NReduce(=10) distinct tasks across nodes.
- To prevent a stuck/failed worker from hoarding the task, ideally create a deadline which automatically fails the task on crossing it.
content, err := os.ReadFile(task.File)
if err != nil {
return err
}
if err := c.Err(); err != nil {
return err
}
// Call map function
kva := grepMyFile(task.File, string(content))
// Write intermediate key-value pairs to files
buckets := make([][]KeyValue, task.NReduce)
for _, kv := range kva {
// hash ... to be done
haskKey := ihash(kv.Key) % task.NReduce
buckets[haskKey] = append(buckets[haskKey], kv)
}
for y, bucket := range buckets {
// **IMP : create temp file for atomicity
fileName := fmt.Sprintf("mr-%d-%d", task.Id, y)
tmpFile, err := os.CreateTemp(".", "mr-tmp-*")
enc := json.NewEncoder(tmpFile)
for _, kv := range bucket {
enc.Encode(&kv)
}
tmpFile.Close()
os.Rename(tmpFile.Name(), fileName)
}
Step 4 - Perform Reduce Task
The reduce Step works in similar way as Map Step. You collect the results produced by various workers to create a unified result.
Top comments (0)