3 · Under the hood · lesson 11 of 20
Partitions
The unit of parallelism — one partition = one task per stage.
A DataFrame is split into partitions. Spark launches one task per partition per stage. Too few partitions and cores sit idle; too many and scheduling overhead dominates.
Python
df.rdd.getNumPartitions() # inspect
df.repartition(200) # full shuffle, exact count
df.repartition("country") # shuffle so same country co-locates
df.coalesce(10) # merge partitions, no shuffle
spark.conf.set("spark.sql.shuffle.partitions", 200) # default after a shuffleTIP
A good starting rule: aim for partitions of ~100-200 MB each. On a 100 GB dataset that's 500-1000 partitions.
Loading 3D scene…
Key takeaways
- ✓Partitions are the unit of parallelism — 1 partition = 1 task.
- ✓repartition() causes a shuffle; coalesce() only merges (cheap).
- ✓spark.sql.shuffle.partitions controls partition count after a shuffle (default 200).