Big Data Analysis with PySpark

Tutorial 05 — Lecture 1: Concepts

MapReduce & Distributed Computing: A Gentle, Complete Introduction

Every tutorial so far scaled a database. This lecture scales computation. No installs, no code to run — just the ideas, explained in plain language: what MapReduce is, how Hadoop implements it, why Spark exists, and the vocabulary you'll need before writing a single line of PySpark. Lecture 2 is where you actually build things.

Before You Start

Nothing to install. This lecture is pure concepts — readable on a phone, no terminal needed. Basic programming familiarity (loops, functions) is enough; no prior big data or distributed systems knowledge assumed.

This moves from document databases (Tutorial 4) to distributed computation — a different problem, same "spread it across many machines" philosophy.

Suggested time: 45-60 minutes

What Is MapReduce, and Why Do We Need It?

MapReduce solves the sister problem to the one you just fixed in Tutorial 4 — and fixes it the same way: stop trying to fit everything on one computer, and spread the work across many.

Tutorial 4: one machine can't store all the data → shard it across machines.
This lecture: one machine can't process all the data in time → MapReduce it across machines.

Imagine counting word frequencies in one book using a loop and a dictionary; quite a trivial problem. But now do it for ten million books:

  • The text alone won't fit in one computer's storage (HDD/SSD)
  • A single CPU core would take days
  • You need many computers counting different books at the same time, then a way to merge their partial counts into one answer
That splitting data, then calculating partial results, then merging them back together — that is exactly what MapReduce standardizes.
When do you actually need it?
  • Your dataset no longer fits in one machine's memory or disk
  • A sequential job would take too long, even with a cluster available

Otherwise, a normal Python script (or a MongoDB aggregation pipeline) is simpler and faster — MapReduce is a tool for scale, not a default.

Why Not Just Write Custom Distributed Code?

Before MapReduce existed, programmers solved this by hand — writing their own code to make computers talk to each other over the network, often using an older toolkit called MPI (still used on supercomputers today). Written that way, you become responsible for three hard problems:

  • Splitting the work — deciding which machine gets which chunk of data
  • Moving results — sending partial results across the network
  • Handling failure — noticing when a machine crashes mid-job and re-running just that piece

MapReduce's entire pitch: it handles all three transparently. You write two small functions, map and reduce, and the framework takes care of distribution, data movement, and fault tolerance.

Abstraction

You only write map and reduce. The framework decides which node runs what.

Data Locality

Computation moves to where the data already lives, instead of dragging huge files across the network.

Fault Tolerance

If a node dies mid-job, the framework reruns just its piece elsewhere — you never notice.

The MapReduce Model: Map, Shuffle & Reduce

MapReduce is functional programming applied to distributed data. Python's own built-in map() and reduce() are the same idea on a single machine — understanding them first makes the distributed version click.

Everything Is Key–Value Pairs

MapReduce operates on (key, value) tuples. For a text file, the key might be the line number and the value the line's text. A full MapReduce job runs in three phases:

1
Map

Runs on every node in parallel. Turns each input (k, v) into zero or more intermediate (k', v') pairs.

2
Shuffle

Happens automatically. All pairs with the same key' are grouped together, across the network, into (k', list(v')).

3
Reduce

Runs on every node in parallel. Combines each key's list of values into one final (k', v'').

Wait — there's no literal "node" in PySpark

The words above ("every node", "across the network") describe Hadoop's picture: real, separate machines physically shipping data to each other. Spark keeps the same three phases but runs them differently:

  • Your dataset is split into partitions, not files on separate disks
  • A "node" becomes an executor — which can be a separate machine in a real cluster, or, in local[*] mode (what we use in Lecture 2), just a CPU core on your own laptop
  • You never write the map/shuffle/reduce plumbing yourself — Spark decides how to spread the work across whatever executors it has

Same mental model, different machinery underneath. We'll meet these exact terms in Step 4.

A Complete MapReduce Pass, Worked Through With Real Numbers

Now walk through what happens to five numbers. Say map turns each (k, v) pair into (abs(v), 1), and reduce sums the 1s per key. Watch the two highlighted rows — both map to key 1, so the shuffle phase merges them into a single group before reduce ever runs:

Input (k, v)
Map → (abs(v), 1)
Shuffle — group by key
Reduce → (k', Σv')
(1, 2)
(2, −1)
(3, 1)
(4, 3)
(5, 6)
(2, 1)
(1, 1)
(1, 1)
(3, 1)
(6, 1)
(2, [1])
(1, [1, 1])
(3, [1])
(6, [1])
(2, 1)
(1, 2)
(3, 1)
(6, 1)

Four input pairs, four distinct keys after mapping — but (2,−1) and (3,1) (outlined in pink) both map to key 1. Shuffle merges them into one group before reduce ever sees them, so reduce only runs once per key, not once per input row.

Worked numeric example adapted from Data Analytics with Python, Ch. 2 (MapReduce).

The Same Idea, in Python

Now that you've traced it by hand, here's the identical pattern in plain Python — no Spark, no cluster, just a single machine. This is the exact building block Spark scales up.

Python — map() and reduce() on a single machine python
from functools import reduce

lst = [1, 2, 3, 4]

# map: apply a function to every element -> a new sequence
squares = list(map(lambda x: x * x, lst))
print(squares)                       # [1, 4, 9, 16]

# reduce: combine all elements pairwise into ONE value
def add_reduce(x, y):
    return x + y

total = reduce(add_reduce, lst)
print(total)                         # 10  (((1+2)+3)+4)

# Chained: square every element, then sum the squares
print(reduce(add_reduce, map(lambda x: x * x, lst)))   # 30
Critical rule: reduce must be commutative & associative

add_reduce(x, y) must equal add_reduce(y, x), and grouping must not matter. A distributed job never guarantees value order — an order-dependent function (like averaging pairs) silently gives a wrong, non-reproducible answer.

Combiners: A Local Reduce Before the Shuffle

If a node produces the pair ("Big", 1) three times, sending all three across the network wastes bandwidth. A combiner pre-aggregates values that share a key locally, before the shuffle — a pure optimization, only possible because reduce is already commutative and associative.

groupByKey()

Ships every single value across the network first, combines after. No local combiner.

reduceByKey()

Combines locally first, ships only the partial totals. Same result, less network traffic.

Both produce the identical final result — we'll prove it side by side with real, running PySpark code in Lecture 2.

Fault tolerance, briefly
  1. A mapper/reducer node fails or runs slow
  2. The framework relaunches the same task on a healthy node, using a replica of the input data
  3. Whichever copy finishes first wins; the other is discarded

This is why MapReduce demands idempotent, side-effect-free functions — a task might run more than once.

Hadoop: The First Implementation — and Why Spark Replaced It

Hadoop: The First Implementation

MapReduce is a model; something has to actually implement it. Hadoop (2006) was the first popular open-source implementation, and it ships two pieces that matter conceptually:

YARN (Resource Manager)

A master/worker scheduler. A Resource Manager on the master node tracks every worker's free CPU/RAM; a Node Manager on each worker reports status and runs containers — sandboxed slices of a machine where mapper/reducer tasks actually execute.

HDFS (Storage)

The Hadoop Distributed File System. Files are split into large (128 MB) blocks, spread across worker Data Nodes, and replicated 3× for fault tolerance. A Name Node tracks which blocks live where — this is what makes data locality possible.

Hadoop's Master/Worker Architecture, Visually
How to read this diagram
Master Node — coordinates
Worker Node — does the work
Resource Manager — assigns tasks
Name Node — tracks file locations
Node Manager — runs tasks here
Data Node — stores file blocks
Master Node
Resource Manager Name Node
Worker Node 1
Node Manager Data Node
Worker Node 2
Node Manager Data Node
Worker Node 3
Node Manager Data Node

One Master Node per cluster coordinates; every Worker Node reports its status up and receives task/storage instructions down. Lose a worker, and the master reassigns its work & rebuilds its data from replicas elsewhere — the fault tolerance from Step 1, made concrete.

Why Spark Instead of Raw Hadoop?

Hadoop MapReduce writes intermediate results to disk between every map and reduce stage — safe, but slow for machine learning jobs that pass over the same data dozens of times. Spark (2014) fixes this:

  • Keeps data in RAM across steps, in a structure called a Resilient Distributed Dataset (RDD)
  • Same MapReduce ideas, same fault tolerance and data locality — but up to 100× faster for iterative jobs
  • Runs on top of Hadoop (using YARN and HDFS), or entirely standalone in local mode on one laptop — what we'll do in Lecture 2, where every partition becomes a thread instead of a network hop
Concept Hadoop MapReduce Spark (local mode, Lecture 2)
Between stages Written to HDFS (disk) Kept in memory (RDD)
Worker unit Mapper / reducer processes in YARN containers Executor threads over CPU cores
Data unit 128 MB HDFS blocks In-memory partitions

PySpark Vocabulary: The Building Blocks

Before Lecture 2 has you typing code, let's learn the words that code uses. Good news: you already know the concepts from Hadoop's master/worker architecture — Spark just uses different names for (mostly) the same roles.

Same Roles, New Names

In Hadoop, you met… In Spark, it's called… In plain English
Client submitting a job Driver Program The process running your Python script top to bottom.
Resource Manager Cluster Manager Finds available machines/cores and hands them to the driver.
Node Manager + Container Executor A worker process that holds data and runs tasks on it.
HDFS Block Partition One slice of your dataset that one executor works on.

The PySpark-Specific Glossary

Four more terms are specific to how Spark's Python API is organized — you'll see all four in almost every script you write in Lecture 2.

SparkSession

Think of it as the front desk. It's the single entry point for everything — reading files, running SQL, building ML models. You create exactly one per program, and every other Spark object is reached through it.

SparkContext

The lower-level radio the SparkSession uses internally to actually talk to the cluster manager. You'll rarely touch it directly — mostly via spark.sparkContext for RDD-specific work.

RDD (Resilient Distributed Dataset)

A distributed collection of Python objects, split into partitions across executors. "Resilient" means Spark remembers the recipe used to build it (its lineage), so it can rebuild any lost piece after a failure — no backup files needed.

DataFrame

An RDD with column names and types attached — a spreadsheet instead of a bag of arbitrary objects. Built for structured, tabular data, which is exactly what our CSV files are.

How They Fit Together

Your Script creates a SparkSession (the Driver)
Cluster Manager local[*] = your own CPU cores
Executors run tasks on RDD/DataFrame partitions
Preview only — don't run this yet. Here's what creating a SparkSession actually looks like in code. We'll properly install Java and PySpark and run this for real in Lecture 2, Step 1.
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .master("local[*]") \
    .appName("my-first-app") \
    .getOrCreate()

sc = spark.sparkContext   # low-level handle, only needed for RDD operations

Concept Check

No code here — just check that the ideas stuck before moving on.

Q1. In Hadoop, which process actually decides which worker machine runs a given task: the Resource Manager, or the Node Manager?

Q2. Why is Spark typically much faster than Hadoop MapReduce for training a machine learning model?

Q3. What's the one-sentence difference between an RDD and a DataFrame?

Q4. You create a SparkSession with .master("local[*]"). What does that actually do, in plain terms?

Recap & What's Next

You now know, in plain language: what problem MapReduce solves, how map/shuffle/reduce and combiners work, how Hadoop implements it with a master/worker architecture, why Spark replaced raw Hadoop for iterative workloads, and the core PySpark vocabulary — SparkSession, RDD, DataFrame, driver, executor, and partition.

Lecture 2 is entirely hands-on: install Java and PySpark, then build four real programs — Word Count, a key–value RDD job, a DataFrame join, and a trained machine learning pipeline — all on real e-commerce data.

Continue to Lecture 2: Implementation