Big Data Analysis with PySpark

Tutorial 06 — Lecture 2: Implementation

PySpark Implementation: Step by Step

Now that you know what MapReduce is and why Spark exists, it's time to write code. Set up a real PySpark environment on your laptop and build four working programs — from a plain Word Count to a trained machine learning model — on real e-commerce data.

Before You Start

Required: complete Lecture 1: MapReduce & Distributed Computing Concepts first. This lecture assumes you already know what map/reduce, Hadoop, and a SparkSession are — we go straight to installing and coding.

We reuse the Olist E-Commerce dataset (already in data/olist_dataset/ if you completed Tutorial 3).

Suggested time: 100-130 minutes

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).

Ubuntu / Debian
sudo apt update
sudo apt install -y openjdk-21-jdk-headless
java -version
Fedora / RHEL
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
Remember to 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()
"
Expected output ends with "PySpark is running. Spark version: 4.2.0". A wall of WARN log lines above it (about hostname resolution, native libraries, etc.) is normal — that's Spark's own logger, not an error.
Where 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.

generic_template.py python
# 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.

Mental model to keep: imports & session → load → transform (lazy) → action (runs it) → stop. Every example from here on follows exactly this shape — only steps 4 and 5 change.

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).

Trace It By Hand First

Before opening an editor, run the four questions above on two tiny lines: "the cat sat" and "the cat ran".

1–2. After splitting into words the, cat, sat, the, cat, ran
3. After pairing with a value (the,1) (cat,1) (sat,1) (the,1) (cat,1) (ran,1)
4. After combining by key (the,2) (cat,2) (sat,1) (ran,1)

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:

word_count.py python
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
Expected output (16 lines, alphabetical): starts with 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.

Apply the Framework from Step 3
  • Final answer shape: (month, order_count)
  • Key: the calendar month, sliced out of order_purchase_timestamp
  • Value: 1 per 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.

orders_per_month.py python
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
Expected output (25 lines): starts with 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()
Output ends with "Identical results: True". On this tiny dataset the speed difference is invisible, but on a real cluster with millions of repeated keys, 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.

revenue_by_category.py python
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
Expected output (top 10 of 72 categories):
+-----------------------------+-------------+----------+
|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      |
+-----------------------------+-------------+----------+
Under the hood, it's still MapReduce. .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.

predict_shipping_cost.py python
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
Expected output ends with "RMSE: 11.99". That means the model's shipping-cost predictions are off by about R$12 on average — a reasonable first model, not a great one. Coefficients will be small positive numbers; the largest is on product_height_cm, meaning height moves predicted freight more than the other dimensions per unit.
Same MapReduce rules apply. Fitting a linear regression model requires an iterative numerical optimization (gradient descent under the hood) that reduces partial gradients across all partitions on every pass. This is precisely the disk-I/O-per-iteration problem that made Hadoop MapReduce impractical for ML (see Lecture 1, Step 3) — and precisely what Spark's in-memory RDDs were built to fix.

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.

Congratulations! You've gone from a plain-Python 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.

1. PySpark Installed python3 -c "import pyspark; print(pyspark.__version__)"

Prints 4.2.0 (or your pinned version)

2. Word Count Runs python3 word_count.py

16 lines; "Big: 4" and "World: 3" appear

3. Real Dataset Readable python3 orders_per_month.py

25 months; totals 99441 orders

4. ML Pipeline Trains 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