Input
Map
Partitions
Reduce
| Word | Count |
|---|
| Word | Count |
|---|
This C MapReduce library runs a
word-count pipeline. Word-count is a classic
example: you want to count how often each word appears, so the
pipeline first turns each word into a (word, 1) pair,
then groups pairs by word, then sums the ones to get the final
count. Input is processed in parallel by mappers,
their output is grouped by key into partitions, and
reducers aggregate counts per key. Above, you can
paste text, choose how many partitions to use, and run to see each
stage of the pipeline plus the word cloud, bar chart, and table of
results. The same design scales to many input files on disk.
The library uses a thread pool and runs one
mapper per input (e.g. one thread per file). Each mapper reads
its input, tokenizes on whitespace, and for every word calls
MR_Emit(word, "1"). Emitting
"1" means “this word appeared once”; the reduce
phase will later sum those ones to get the total count per
word. The library does not send key-value pairs straight to
the reducers. Instead, a partitioner hashes
each key to decide which partition (0 to
num_parts - 1) will receive that pair. The same
key always hashes to the same partition, so all pairs with the
same word end up in one partition. That way, a single reducer
can process all values for that word. Map output is therefore
sharded by key across partitions, and multiple mappers can run
in parallel without blocking each other.
Partitions sit between map and reduce and act as the
shuffle step. After map, you have a stream of
(key, value) pairs from many mappers; before
reduce, all values for the same key must be in one place. Each
partition holds every pair whose key hashed to that partition
index, so within a partition all pairs for a given key are
grouped together. Reducers do not see raw map output; they see
keys and the list of values per key within their partition.
The library moves data into partitions automatically
(buffering and grouping by key). Your code only emits in map
and consumes in reduce. The number of partitions controls how
much parallelism you get in the reduce phase: one reducer runs
per partition, so more partitions mean more concurrent
reducers, at the cost of more overhead.
One reducer runs per partition, again using the thread pool,
so reduce is parallel as well. For each distinct key in that
partition, the library calls your reducer with that key and
the partition index. The reducer’s job is to combine all the
values for that key. It calls
MR_GetNext(key, partition_idx) repeatedly; each
call returns the next value for that key, or NULL
when there are no more. For word-count, the values are the
"1" strings emitted by the mappers, so the
reducer typically counts or sums them and writes the final
count (e.g. to a file or a shared structure). When the reducer
returns, the library moves on to the next key in the
partition. All partitions are processed in parallel, so the
total work is spread across multiple threads.
Under the hood, the implementation uses pthread
synchronization primitives to keep the pipeline thread-safe
and efficient. The thread pool uses
pthread_mutex_t and
pthread_cond_t to protect the job queue and to
let workers sleep while waiting for work. It also uses a
separate mutex and condition variable to coordinate clean
termination so that threads can exit without busy-waiting.
The thread pool also supports a Shortest Job First (SJF)
scheduling policy by tracking a size field on
each job and prioritizing smaller jobs first. This can improve
overall throughput and reduce average waiting time when the
workload is a mix of short and long tasks.
The MapReduce implementation also uses
pthread_mutex_t
to guard intermediate results. Each partition has its own
mutex, which prevents races when multiple mapper threads emit
keys that land in the same partition. Different partitions can
still be written in parallel because they do not share a lock.
Partitions are implemented as an array of partition
structures. Each partition stores intermediate
(key, value) pairs in a linked list with
head and tail pointers for efficient
insertion and traversal. Keys are maintained in
ascending order within each partition, which
makes it straightforward to process distinct keys in reduce. A
per-partition mutex keeps insertions and reads thread-safe
during concurrent emits.
The library was tested with automated scripts that vary input
sizes and concurrency settings. The included
test.sh script validates correctness by checking
the generated result files, and the Makefile includes a
Valgrind target to help catch memory errors and leaks. The
code was also exercised with different combinations of worker
threads and partition counts to validate behavior under
different workloads.
You plug in your own Map and
Reduce functions and call
MR_Run with the input files and parameters. The
following are the main entry points and callbacks.
MR_Run(file_count, file_names, Map, Reduce, num_workers,
num_parts)
is the entry point. You pass the number of input files, the
array of file paths, your Map and Reduce function pointers,
the size of the thread pool (num_workers), and
the number of partitions (num_parts). It runs
the full map then reduce pipeline.
MR_Emit(key, value) is called from your mapper
to output a key-value pair. The library takes ownership of
storing and routing the pair to the correct partition.
MR_GetNext(key, partition_idx) is called from
your reducer to read the next value for the current key in
the given partition. It returns NULL when there
are no more values for that key. You typically call it in a
loop to consume all values before returning from the
reducer.
flowchart LR
subgraph input [Input]
T[Files]
end
subgraph map [Map]
M[Mapper]
E[MR_Emit]
end
subgraph part [Partitions]
P0[Partition 0]
P1[Partition 1]
Pn[Partition n]
end
subgraph reduce [Reduce]
R[Reducer]
Out[Counts]
end
T --> M
M --> E
E --> P0
E --> P1
E --> Pn
P0 --> R
P1 --> R
Pn --> R
R --> Out