§ 1.3Module 1

Case Study of Big Data Solutions

On this page

1.3 Case Study of Big Data Solutions: Web-Search Indexing

Recall first

A search engine receives documents and user queries. Which task is naturally parallel: reading documents to produce terms, grouping all documents for one term, or both? What must survive a worker failure?

First principles

The case

Use a web-search indexing pipeline as the case study. The goal is to turn a very large document collection into an inverted index: for each term, store the documents in which it occurs, optionally with counts or positions. The Google MapReduce paper presents this style of large-cluster processing and motivates hiding partitioning, scheduling, communication, and failure recovery inside a framework (Dean and Ghemawat, MapReduce).

Solution architecture

  1. Collect: crawlers write fetched documents and metadata to distributed storage.
  2. Parse: map workers read document chunks, clean text, tokenize, and emit (term, document_id) pairs.
  3. Group: the framework sends equal terms to the same logical reduce group.
  4. Reduce: each term’s values are sorted/deduplicated and written as its posting list.
  5. Serve: an index-serving layer looks up query terms and ranks matching documents. Serving is a different low-latency workload from offline indexing.

A Hadoop-style platform fits the offline stage because HDFS is designed for large datasets, streaming access, replication, and computation close to data (HDFS Design).

Worked trace

Documents:

D1: big data data
D2: big systems

A mapper emits (after tokenization):

(big,D1) (data,D1) (data,D1)
(big,D2) (systems,D2)

The grouping stage creates:

big    -> [D1,D2]
data   -> [D1,D1]
systems-> [D2]

The reducer can deduplicate and count:

big    -> [(D1,1),(D2,1)]
data   -> [(D1,2)]
systems-> [(D2,1)]

If the desired index is sorted by term, the key ordering supplied by the framework helps. If the desired output is sorted by document or rank, a later job or different key design is needed. This illustrates a general lesson: the emitted key determines the communication and grouping pattern.

What makes it a big-data solution?

The trade-offs are also visible: indexing has batch latency, the shuffle can be network-heavy, hot terms can create reducer skew, and a changed tokenizer may require rebuilding or incrementally maintaining derived data. The solution is not “MapReduce everywhere”; a search query needs a serving system optimized for low latency.

Exercise — revealed answer

Exercise: Why emit (term, document_id) rather than (document_id, term) for building an inverted index?

Answer: Grouping by the key sends all occurrences of one term to one reduce group, exactly where its posting list can be built. Reversing the pair groups a document’s terms together, which is useful for a document-centric index but not for term lookup.

Exam lens

Draw the pipeline documents → map(term, doc) → group by term → reduce(posting list) → serving index. State one role each for partitioning, data locality, replication, and rerunning failed tasks. Separate offline indexing latency from online query latency.

Rapid revision checklist

Key takeaways

  1. A big-data solution matches a workload: offline indexing and online serving need different designs.
  2. MapReduce’s key controls where related records meet.
  3. Partitioning, locality, replication, and retry turn a cluster of unreliable machines into a usable batch platform.

Sources