11.8 Worked examples
① Working set: hash vs sort. 10 billion log lines, 2 million distinct URLs, average URL 60 bytes.
- Hash: 2e6 × (60 B + 8 B counter + ~50 B overhead) ≈ 236 MB → fits comfortably. Record count is irrelevant.
- Now 10 billion lines with 500 million distinct URLs (e.g. keyed by session ID): 5e8 × 118 B ≈ 59 GB → exceeds a laptop, and sorting (spilling to disk, sequential I/O) wins. The lesson: choose by cardinality, not volume.
② Shuffle volume. 1 TB input, 200 map tasks, 200 reduce tasks. Each mapper writes 200 files, so the shuffle produces 40,000 files and moves ~1 TB across the network. At 10 Gb/s per node over 20 nodes, that's ~1 TB / 25 GB/s ≈ 40 s of pure network time — before any computation. This is why advantage ② (operator fusion, avoiding unnecessary shuffles) matters more than raw CPU.
③ Skew. 100 reduce tasks; keys distributed so one key has 30% of the rows. That reducer does 30% of the work alone, so wall-clock ≈ 0.30 × total, versus 0.01 × total if perfectly balanced — a 30× worse stage duration than the ideal, no matter how many machines you add. Adding nodes cannot fix skew; only changing the key can.
④ Recomputation cost on spot instances. A 5-stage job, each stage 20 minutes, intermediate data in memory. An executor is preempted during stage 5 holding shuffle output from stage 4. Spark must recompute stage 4 for the lost partitions. With a 10% preemption rate per 20-minute window across 100 executors, you will lose executors most runs — and without an external shuffle service, expected runtime inflates substantially. This is the concrete reason MapReduce wrote everything to the DFS, and the concrete cost of Spark's choice not to.
⑤ The direct-write anti-pattern, quantified. 200 parallel tasks × 5,000 records/s each = 1,000,000 writes/s aimed at a production database sized for 20,000 writes/s. 50× over capacity — the database doesn't just slow the job, it takes down the user-facing application. Via Kafka, the same 1M/s is absorbed by the log's sequential writes and the consumer drains at 20,000/s over ~50× longer — with production traffic unaffected.
⑥ Block size and metadata. 1 PB of data.
- ext4-style 4 KiB blocks → 2.7 × 10¹¹ blocks of metadata. Impossible.
- HDFS 128 MB blocks → 8.4 million blocks. At ~150 bytes of NameNode memory per block, ≈ 1.2 GB — manageable.
- Same 1 PB stored as 1 billion 1 MB files → 1 billion blocks → ~150 GB of NameNode heap. This is the small-file problem, and it's a metadata problem, not a data problem.