Spark splits the data into partitions and plans the whole job before running anything, and a group-by triggers a shuffle so each store's rows meet. 100 million rows in 4 partitions: store 3 leads at $479,607,696.
python3 --version.pip install "pyspark>=3.5"brew install openjdk@21git clone https://github.com/DayanEbrar0X/data-anatomy.ai.git cd data-anatomy.ai
python3 -m venv .venv source .venv/bin/activate
pip install -r requirements.txt # or just this lesson: pip install "pyspark>=3.5"
brew install openjdk@21 # macOS sudo apt install openjdk-21-jdk # Debian / Ubuntu java -version
cd data-engineering/13-pyspark python3 src/sales.py
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = (SparkSession.builder.master("local[4]")
.config("spark.ui.enabled", "false")
.getOrCreate())
spark.sparkContext.setLogLevel("ERROR")
sales = (spark.range(100_000_000, numPartitions=4)
.withColumn("store", F.col("id") % 50)
.withColumn("amt", F.abs(F.hash("id")) % 500))
print("partitions:", sales.rdd.getNumPartitions())
big = sales.filter(F.col("amt") >= 100)
top = (big.groupBy("store")
.agg(F.sum("amt").alias("revenue"))
.orderBy(F.desc("revenue")))
print("plan built, nothing ran yet")
rows = top.take(3) # action
for r in rows:
print(f"store {r.store:>2}: ${r.revenue:,}")
print("top store:", rows[0].store)Your laptop chokes on a billion rows. Spark splits the job across a hundred machines. The trick: cut the data into partitions. Every core or machine takes its own slice, all at once.
Our job: a hundred million sales. Which store sells the most? First, a Spark session. local[4] means four local cores.
Now the data: a hundred million rows, in four partitions. Each row gets a store, one of fifty, and an amount from a hash. Spark reports the partition count without running anything. Keep sales of a hundred or more.
That's a filter. Group by store, add up the amounts, and sort by revenue. Here's the catch: none of that has run yet. Spark is lazy.
Each transformation only adds a step to a plan. The work starts at an action: here, take(3). Now the plan lights up, and every partition runs at once. Each one filters its own rows and sums its own stores.
Then comes the shuffle. Totals for one store cross partitions and meet in one place. Then one sum per store, sorted, and the top three come back. Let's run it.
Four partitions. And the plan was built before anything ran. Then take(3) runs the job. Store 3 leads, with 479.
6 million. One winner, from a hundred million rows on one laptop. On a real cluster, you drop local[4]. The same code then spreads across hundreds of machines.
That's why Spark runs big batch jobs and ML feature pipelines, on Databricks, EMR and Dataproc. The gotcha: a shuffle moves data over the network. It's costly. Split.
Plan. Act once. That's how Spark outgrows one machine.
Read the lesson on GitHub →