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