Environment Setup
In Lecture 1 you learned that a SparkSession is your single
entry point into Spark, and that local[*] mode turns your own
laptop's CPU cores into "workers." Let's install everything needed to actually create one.
Install Java
Spark runs on the JVM, so PySpark needs a Java runtime even though you only write Python. Java 17 or newer works (this tutorial was verified end-to-end on Java 21).
sudo apt update
sudo apt install -y openjdk-21-jdk-headless
java -version
sudo dnf install -y java-21-openjdk
java -version
Install PySpark
Create an isolated virtual environment inside the project so PySpark doesn't clash with any other Python setup, then install it:
# From the project root (big-data-analysis-with-mogodb)
python3 -m venv .venv
source .venv/bin/activate
pip install --upgrade pip
pip install pyspark==4.2.0
source .venv/bin/activate again
in every new terminal tab before running the scripts below. The download is around
300–450 MB, so this step needs a decent connection.
Verify the Install
This is the exact SparkSession
pattern previewed in Lecture 1 — now it actually runs:
python3 -c "
from pyspark.sql import SparkSession
spark = SparkSession.builder.master('local[*]').appName('verify-pyspark').getOrCreate()
print('PySpark is running. Spark version:', spark.version)
spark.stop()
"
WARN log lines
above it (about hostname resolution, native libraries, etc.) is normal — that's
Spark's own logger, not an error.
local[*] fits in the picture:
This tells Spark's Cluster Manager to use your own CPU cores as "workers" instead of connecting to YARN or a real cluster. Every concept from Lecture 1 — the driver submitting a job, executors running tasks in parallel, partitions of data — still applies; it's just all happening inside one process on your laptop.
Anatomy of a PySpark Script
Every PySpark script you write — whether it's a five-line word count or a full ML pipeline — follows the same skeleton. Before jumping into the examples, let's build a mental template you can reuse for any PySpark program.
# 1. IMPORTS - bring in only what you need
from pyspark.sql import SparkSession
from pyspark.sql import functions as F # common alias for built-in column functions
# 2. ENTRY POINT - one SparkSession per script, built once
spark = (SparkSession.builder
.master("local[*]") # run locally, use all available cores
.appName("my-script-name") # shows up in the Spark UI / logs
.getOrCreate())
sc = spark.sparkContext # only needed if you use the older RDD API
# 3. HOUSEKEEPING - quiet the noisy default logger (optional, but recommended)
sc.setLogLevel("ERROR")
# 4. LOAD DATA - read a file, or build a small dataset from Python values
df = spark.read.csv("data/some_file.csv", header=True, inferSchema=True)
# 5. TRANSFORM - describe the pipeline; nothing runs yet (lazy evaluation)
result = (df
.filter(F.col("status") == "delivered") # keep only relevant rows
.groupBy("category") # MAP/SHUFFLE step
.agg(F.count("*").alias("total"))) # REDUCE step
# 6. ACTION - this is what actually triggers computation across the cluster
result.show() # or .collect(), .write.csv(...), .count(), etc.
# 7. CLEAN UP - release the cluster resources when you're done
spark.stop()
Nothing in this file does real work on its own — it's a shape you'll see in every example below. Let's walk through what each numbered section is actually doing.
Imports
Almost every script needs SparkSession from
pyspark.sql. If
you'll build column expressions (filters, aggregations, new columns), also import
functions as F
— you'll see F.col(...)
and F.count(...)
constantly from Example 3 onward.
Create the
SparkSession
This is the single entry point into Spark you met in Lecture 1 and built in Step 1.
.master("local[*]")
picks the cluster manager (your laptop's cores, here), and .appName(...) is just
a label for logs and the Spark UI. Build it once per script —
everything else hangs off this spark object.
Housekeeping
setLogLevel("ERROR")
suppresses the wall of WARN lines you
saw in Step 1. It's optional, but every example script below includes it so the
output you actually care about isn't buried.
Load Data
Data enters a Spark program in one of a few ways: spark.read.csv(...)
(or .json(...),
.parquet(...)) for
DataFrames, sc.textFile(...)
for raw RDDs, or sc.parallelize(...)
to turn an in-memory Python list into an RDD for quick experiments — exactly
what Example 1 does below.
Transform (lazy)
Calls like .filter(),
.map(), .groupBy(), and .agg() describe
what should happen to the data — they don't run it yet. Spark just
builds up a execution plan. Chaining several of these together, like the pipeline
above, is completely normal and reads top-to-bottom as one recipe.
Trigger with an Action
Nothing computes until you call an action — .show(), .collect(), .count(), or a .write...() call. That
single line is what sends the whole transform pipeline above it out to run.
Stop the Session
spark.stop() frees
the resources your session was holding. Easy to forget in a notebook, but good habit
in standalone scripts like the ones you'll write here.
Example 1: Word Count — the "Hello World" of Big Data
Every big data course starts here for a reason: it's simple enough to reason about by hand,
yet it exercises the full map
→ shuffle → reduce pipeline from Lecture 1.
The Mindset: Frame the Problem Before You Type Code
The hard part of MapReduce-style code was never the Python syntax — it's knowing what to write before you write it. Almost every "count things" or "total things per group" problem, word count included, answers the same four questions in the same order. Answer them on paper first; the code below is just those answers translated into PySpark calls.
What should the final answer look like?
Name the shape of the output before touching code. Here: a list of pairs, (word,
count). Once you know the shape you're aiming for, everything else
follows from it.
What's the "group by" — your key?
Ask: "if I had a million of these, what would I bucket them by?" For word count, each distinct word is a bucket. That word is your key.
Turn every
input record into (key,
value) — this is MAP
Pick a value that's safe to combine in any order, on any machine, in any
partition — usually 1
(to count) or the number itself (to sum). Addition is commutative and
associative, so Spark can add up partial results anywhere and combine them later
without changing the answer. Reshaping raw input into (key,
value) pairs is exactly what .map()
does — here it takes two calls, see below.
Combine values sharing a key — this is SHUFFLE + REDUCE
Spark gathers every pair with the same key together (the shuffle) and folds their
values into one (the reduce). In PySpark that's almost always a single line:
reduceByKey(lambda
a, b: a + b).
Before opening an editor, run the four questions above on two tiny lines: "the cat sat" and "the cat ran".
Each column above maps to exactly one line of PySpark below — column 1 is flatMap, column 2 is map, column 3 is reduceByKey. That's the
whole trick: every transformation in the code is one column in your
hand-trace.
With that mental model in hand, here's the same logic as PySpark. Create word_count.py in the project root:
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[*]").appName("word-count").getOrCreate()
sc = spark.sparkContext
sc.setLogLevel("ERROR") # silence the WARN noise from Step 1
text = """Welcome to the Big World of Big Big Data Welcome World bye
World Hello MapReduce GoodBye MapReduce
This Book on Big Data is fun"""
lines = sc.parallelize(text.strip().split("\n")) # one RDD partition per line
word_counts = (lines
.flatMap(lambda line: line.split(" ")) # MAP: line -> many words (flatMap, not map!)
.map(lambda word: (word, 1)) # MAP: word -> (word, 1)
.reduceByKey(lambda a, b: a + b)) # SHUFFLE + REDUCE: sum the 1s per word
for word, count in sorted(word_counts.collect()):
print(f"{word}: {count}")
spark.stop()
python3 word_count.py
Big: 4, includes World: 3, MapReduce: 2, Welcome: 2, Data: 2, and the rest appearing once each.
Why flatMap, Not map?
map() would turn each
line into a list of words — three lines in, three lists out, still
nested. flatMap() flattens
those lists into one continuous stream of words, which is what we need before pairing each
word with a 1. Every RDD
operation here is lazy — nothing actually runs until .collect() (an action)
triggers the whole pipeline.
Example 2: Key–Value RDDs — Orders Per Month
Now for real data. We'll process data/olist_dataset/olist_orders_dataset.csv
(99,441 orders, the same file from Tutorial 3) using plain RDDs to count how many orders were
placed each calendar month — the classic "max temperature per year" style of MapReduce
problem, applied to real e-commerce data.
- Final answer shape:
(month, order_count) - Key: the calendar month, sliced out of
order_purchase_timestamp - Value:
1per order — so counting orders is just summing 1s - Combine:
reduceByKey(lambda a, b: a + b), exactly like word count
Same four questions, same answer pattern — only where the key comes from changes.
import csv
import io
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[*]").appName("orders-per-month").getOrCreate()
sc = spark.sparkContext
sc.setLogLevel("ERROR")
DATA_PATH = "data/olist_dataset/olist_orders_dataset.csv"
raw = sc.textFile(DATA_PATH)
header = raw.first()
rows = raw.filter(lambda line: line != header) # drop the CSV header row
def parse_line(line):
return next(csv.reader(io.StringIO(line))) # handles quoted commas correctly
parsed = rows.map(parse_line)
# order_purchase_timestamp is column index 3, e.g. "2017-10-02 10:56:33"
orders_per_month = (parsed
.map(lambda cols: (cols[3][:7], 1)) # MAP: (k, v) -> ("2017-10", 1)
.reduceByKey(lambda a, b: a + b) # SHUFFLE + REDUCE: sum per month
.sortByKey()) # chronological order
for month, count in orders_per_month.collect():
print(f"{month}: {count} orders")
spark.stop()
python3 orders_per_month.py
2016-09: 4 orders, ends with 2018-10: 4 orders, peaking around 2018-01 at over 7,000 orders. Total across all
months: 99,441.
Proving the
Combiner Effect: reduceByKey
vs groupByKey
Both transformations below produce the exact same word-count result from Step 2, but
take different routes to get there. Add this to a new file, combiner_demo.py, and run it:
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[*]").appName("combiner-demo").getOrCreate()
sc = spark.sparkContext
sc.setLogLevel("ERROR")
text = """Welcome to the Big World of Big Big Data Welcome World bye
World Hello MapReduce GoodBye MapReduce
This Book on Big Data is fun"""
pairs = (sc.parallelize(text.strip().split("\n"))
.flatMap(lambda line: line.split(" "))
.map(lambda w: (w, 1)))
# groupByKey: ships every single (word, 1) pair across the network first,
# THEN sums each key's full list locally. No local pre-aggregation.
grouped = pairs.groupByKey().mapValues(sum)
# reduceByKey: combines same-key pairs on each partition FIRST (the combiner),
# and only sends already-summed partial totals across the network.
combined = pairs.reduceByKey(lambda a, b: a + b)
print("groupByKey :", sorted(grouped.collect()))
print("reduceByKey:", sorted(combined.collect()))
print("Identical results:", sorted(grouped.collect()) == sorted(combined.collect()))
spark.stop()
groupByKey() can shuffle orders of magnitude
more data across the network. Default to reduceByKey() whenever your operation is
commutative and associative — exactly the combiner rule from Lecture 1.
Example 3: Spark SQL & DataFrames — Revenue by Category
Hand-writing map/reduceByKey chains for every
join and aggregation gets tedious fast — which is exactly why Spark SQL and
DataFrames exist. A DataFrame is an RDD with a known
schema (named, typed columns), which lets Spark's query optimizer plan the map/shuffle/reduce
steps for you, the same way MongoDB's aggregation pipeline planned $lookup joins
behind the scenes in Tutorial 3.
We'll join three Olist CSV files — order items, products, and the category name translation table — to find total revenue and items sold per product category.
from pyspark.sql import SparkSession
import pyspark.sql.functions as F
spark = SparkSession.builder.master("local[*]").appName("revenue-by-category").getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
DATA = "data/olist_dataset"
order_items = spark.read.csv(f"{DATA}/olist_order_items_dataset.csv", header=True, inferSchema=True)
products = spark.read.csv(f"{DATA}/olist_products_dataset.csv", header=True, inferSchema=True)
categories = spark.read.csv(f"{DATA}/product_category_name_translation.csv", header=True, inferSchema=True)
revenue_by_category = (
order_items
.join(products, on="product_id", how="inner")
.join(categories, on="product_category_name", how="left")
.groupBy("product_category_name_english")
.agg(
F.round(F.sum("price"), 2).alias("total_revenue"),
F.count("order_id").alias("items_sold")
)
.orderBy(F.desc("total_revenue"))
)
revenue_by_category.show(10, truncate=False)
spark.stop()
python3 revenue_by_category.py
+-----------------------------+-------------+----------+
|product_category_name_english|total_revenue|items_sold|
+-----------------------------+-------------+----------+
|health_beauty |1258681.34 |9670 |
|watches_gifts |1205005.68 |5991 |
|bed_bath_table |1036988.68 |11115 |
|sports_leisure |988048.97 |8641 |
|computers_accessories |911954.32 |7827 |
|furniture_decor |729762.49 |8334 |
|cool_stuff |635290.85 |3796 |
|housewares |632248.66 |6964 |
|auto |592720.11 |4235 |
|garden_tools |485256.46 |4347 |
+-----------------------------+-------------+----------+
.join() and .groupBy().agg() compile down to the same
shuffle-and-combine machinery you used manually in Step 3 — Spark's Catalyst
optimizer just chooses a more efficient physical plan than a naive translation would
(for example, deciding which side of the join to broadcast). This is exactly why
how="left" matters on the second join: a
handful of products have no English category translation, and left keeps their revenue instead of silently
dropping it.
Example 4: Machine Learning with PySpark — Predicting Shipping Cost
Machine learning on big data is still map (transform data, compute local gradients) and
reduce (aggregate into one model) under the hood — Spark's ML library, MLlib, just wraps it in a much friendlier API. We'll
train a Linear Regression model to predict an order
item's freight_value
(shipping cost) from its price and the product's physical dimensions.
Two MLlib Concepts You Need First
Transformer
Takes a DataFrame, returns a new
DataFrame (adds columns). No training involved — e.g. VectorAssembler, or a fitted model making
predictions.
Estimator
Has a .fit() method: learns parameters from data
and produces a Transformer. LinearRegression is an Estimator; the model
it produces is a Transformer.
MLlib requires all input features combined into a single features vector column
— that's the assembler's job. Chaining the assembler and the model together into one
Pipeline means you fit and
apply both steps with a single call, on both training and test data.
from pyspark.sql import SparkSession
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.regression import LinearRegression
from pyspark.ml import Pipeline
from pyspark.ml.evaluation import RegressionEvaluator
spark = SparkSession.builder.master("local[*]").appName("predict-shipping-cost").getOrCreate()
spark.sparkContext.setLogLevel("ERROR")
DATA = "data/olist_dataset"
order_items = spark.read.csv(f"{DATA}/olist_order_items_dataset.csv", header=True, inferSchema=True)
products = spark.read.csv(f"{DATA}/olist_products_dataset.csv", header=True, inferSchema=True)
df = (order_items.join(products, on="product_id", how="inner")
.select("price", "freight_value", "product_weight_g",
"product_length_cm", "product_height_cm", "product_width_cm")
.dropna())
# 80/20 train/test split. seed=42 makes the split reproducible.
train_df, test_df = df.randomSplit([0.8, 0.2], seed=42)
feature_cols = ["price", "product_weight_g", "product_length_cm",
"product_height_cm", "product_width_cm"]
assembler = VectorAssembler(inputCols=feature_cols, outputCol="features")
lr = LinearRegression(featuresCol="features", labelCol="freight_value", maxIter=10, regParam=0.3)
pipeline = Pipeline(stages=[assembler, lr])
# fit() is an Estimator call -> produces a PipelineModel (a Transformer)
model = pipeline.fit(train_df)
# transform() applies the fitted pipeline to unseen test data
predictions = model.transform(test_df)
predictions.select("freight_value", "prediction").show(5)
evaluator = RegressionEvaluator(labelCol="freight_value", predictionCol="prediction", metricName="rmse")
rmse = evaluator.evaluate(predictions)
print(f"RMSE: {rmse:.2f}")
lr_model = model.stages[1] # the fitted LinearRegressionModel inside the pipeline
print("Coefficients:", lr_model.coefficients)
print("Intercept:", lr_model.intercept)
spark.stop()
python3 predict_shipping_cost.py
product_height_cm, meaning
height moves predicted freight more than the other dimensions per unit.
Practice Exercises: MapReduce & PySpark
Every solution below has been run against the real Olist dataset in this repository — the numbers you see are what you should get too.
Exercise 1: Distinct 2018 Customers
Task:
Using the RDD API on olist_orders_dataset.csv, count how many
distinct customers placed an order in 2018 (order_purchase_timestamp starts with
"2018"). customer_id is column index 1.
Exercise 2: Top 5 Busiest Sellers
Task:
Using the RDD API on olist_order_items_dataset.csv, find the 5
sellers (seller_id, column index 3) with
the most items sold, ranked highest first.
Exercise 3: Average Review Score by Category
Task:
Using DataFrames, join olist_order_reviews_dataset.csv (on order_id) through order items and products to
categories, and compute the average review_score per English category name.
Exercise 4: Does a New Feature Help the ML Model?
Task:
Extend Step 5's pipeline by adding product_photos_qty as a sixth input feature.
Retrain, re-evaluate the RMSE, and decide: did it help?
Exercise 5: Word Ranking Capstone
Task:
Recall the classic MapReduce ranking problem from Lecture 1: how would you rank
words by frequency using MapReduce? Read product_category_name_translation.csv, split
each English category name on underscores (e.g. "health_beauty" → ["health", "beauty"]), and find the 10 most
common words across all 71 category names.
map()/reduce() to a full PySpark ML pipeline on real
e-commerce data. When you're done experimenting, deactivate the virtual environment with
deactivate.
Tutorial Verification Checklist
Run these checks to confirm PySpark is installed and every example produces the expected numbers.
python3 -c "import pyspark; print(pyspark.__version__)"
Prints 4.2.0 (or your pinned version)
python3 word_count.py
16 lines; "Big: 4" and "World: 3" appear
python3 orders_per_month.py
25 months; totals 99441 orders
python3 predict_shipping_cost.py
Prints "RMSE: 11.99"
Troubleshooting Common Issues
Quick Reference: RDD & DataFrame Operations
Bookmark this cheat sheet for future PySpark work.
# --- Session setup ---
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[*]").appName("app-name").getOrCreate()
sc = spark.sparkContext
sc.setLogLevel("ERROR")
# --- Creating RDDs ---
sc.parallelize([1, 2, 3]) # from a Python collection
sc.textFile("path/to/file.csv") # from a text file, one element per line
# --- RDD transformations (lazy) ---
rdd.map(func) # 1 element -> 1 element
rdd.flatMap(func) # 1 element -> 0..N elements, flattened
rdd.filter(func) # keep elements where func() is True
rdd.distinct() # remove duplicates (needs a shuffle)
rdd.reduceByKey(func) # combine values per key, WITH a combiner
rdd.groupByKey() # group values per key, NO combiner (slower)
rdd.sortByKey() # sort a key-value RDD by key
# --- RDD actions (trigger execution) ---
rdd.collect() # bring ALL elements to the driver
rdd.take(n) # bring the first n elements
rdd.takeOrdered(n, key=func) # top-n by a custom key, without a full sort
rdd.count() # number of elements
rdd.reduce(func) # combine ALL elements into one value
# --- DataFrames (Spark SQL) ---
df = spark.read.csv("file.csv", header=True, inferSchema=True)
df.join(other_df, on="key", how="inner") # inner / left / right / full
df.groupBy("col").agg(F.sum("x"), F.avg("y"))
df.orderBy(F.desc("col"))
df.show(10, truncate=False)
# --- MLlib pipeline shape ---
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.regression import LinearRegression
from pyspark.ml import Pipeline
from pyspark.ml.evaluation import RegressionEvaluator
assembler = VectorAssembler(inputCols=[...], outputCol="features")
model_estimator = LinearRegression(featuresCol="features", labelCol="target")
pipeline_model = Pipeline(stages=[assembler, model_estimator]).fit(train_df)
predictions = pipeline_model.transform(test_df)
RegressionEvaluator(labelCol="target", metricName="rmse").evaluate(predictions)
spark.stop() # always stop the session when a script finishes