Back to blog

Spark for Streaming Media Analytics · Part 4 of 8

Spark Performance: Shuffles, Partitions, and Broadcast Joins

Read the plan, find the expensive step, and fix it deliberately instead of adding cluster size.

Chris Eberl

Chris Eberl

Founder • Engineering Leader, GenAI, Data, ML

Data EngineeringOct 8, 20269 min read
Spark Performance: Shuffles, Partitions, and Broadcast Joins

Most slow Spark jobs are slow for a small number of reasons: too much data moved between machines, partitions that are badly sized, or one key holding far more rows than the rest. Each is visible in the plan before you spend money.

Spot the shuffle

In explain() output, an Exchange node is a shuffle: rows are redistributed by key across the network. groupBy, join, distinct, and window partitions usually create one.

Find exchanges
1plan = genre_metrics._jdf.queryExecution().executedPlan().toString()2print(plan.count("Exchange"), "shuffles")3 4# Friendlier: open the SQL tab in the Spark UI after running an action

Broadcast small dimensions

The catalog has 300 rows. Shuffling 200,000 events to join it is wasteful; sending the small table to every executor avoids moving the large one.

Broadcast join
1joined = events.join(F.broadcast(catalog), "title_id")2joined.explain()  # look for BroadcastHashJoin, not SortMergeJoin

Simulate skew

One viral title
1skewed = events.withColumn("title_id",2    F.when(F.rand(seed=7) < 0.4, F.lit(42)).otherwise(F.col("title_id")))3 4# A salted aggregation spreads the hot key across buckets5SALT = 166salted = (7    skewed.withColumn("salt", (F.rand(seed=1) * SALT).cast("int"))8    .groupBy("title_id", "salt").agg(F.sum("watch_seconds").alias("s"))9    .groupBy("title_id").agg(F.sum("s").alias("watch_seconds"))10)
SettingWhat it controlsStarting point
spark.sql.shuffle.partitionsPartitions after a shuffleLet AQE coalesce; tune if tasks are tiny or huge
spark.sql.adaptive.enabledRuntime re-planningOn (default in recent Spark)
spark.sql.adaptive.skewJoin.enabledSplits skewed join partitionsOn; verify in the plan
spark.sql.autoBroadcastJoinThresholdMax size to auto-broadcastRaise carefully; executors hold the copy

Exercises

  • Run the skewed aggregation without salt and find the straggler task in the Spark UI.
  • Disable AQE, re-run the genre metrics, and compare partition counts.
  • Force a sort-merge join on the catalog and compare the shuffle size with the broadcast version.
  • Write a one-paragraph performance note explaining which fix you would ship and why.

Official documentation

Read next

Newsletter

New posts, straight from Chris

A short note from me whenever a new article goes live — product engineering, AI workflows, IoT, indie apps, and engineering leadership. No spam, unsubscribe anytime.

By subscribing, you agree to our Privacy Policy. We do not share your email.