skip to content

In a Spark join, why does broadcasting the small side remove the shuffle and the skew?

level: middleimportance: should knowfreq 62%

answer

  1. the shuffle is what concentrates a hot key
  2. one side goes everywhere, one side stays put
  3. the driver holds it before the executors do
  4. ten megabytes, by default
  5. which side can be the build side depends on join type

basics

~20 s

A broadcast hash join ships one whole small table to every executor and probes it from inside each existing partition of the large table. No rows are redistributed by key, so a dominant key value never concentrates into a single task.

solid answer

~50 s

In a shuffle join both sides are re-partitioned by the join key, which is exactly what concentrates a hot key: every row carrying it lands in one partition and one task. A broadcast hash join skips that. Spark collects the small side to the driver, ships it to every executor once, builds a hash table there, and each task of the large side probes it locally — so the large side keeps whatever partitioning it already had and the work stays even. Spark does this automatically when a side's statistics fall under `spark.sql.autoBroadcastJoinThreshold` (10 MB by default; `-1` disables), or when you force it with `broadcast(df)` / the `BROADCAST` hint. The costs are real: the build side must fit in the driver and in every executor, and a too-large broadcast shows up as driver memory pressure or a `spark.sql.broadcastTimeout` failure after its 300-second default.

code

python · 5 lines
python
from pyspark.sql.functions import broadcast

# force the small dimension to be the build side
joined = orders.join(broadcast(customers), "customer_id", "inner")
joined.explain()

go deeper

for a junior

Know that a join normally shuffles both sides by the key, and that when one side is small enough Spark can instead send it to every executor and avoid moving the big table at all.

for a middle

Explain the mechanics — collect to driver, broadcast, build a hash table, probe locally — and name spark.sql.autoBroadcastJoinThreshold with its 10 MB default and the broadcast() hint that overrides it.

for a senior

Demonstrate the judgment: recognise broadcast as the first-line fix for a skewed join, size the build side against driver and executor memory, and diagnose broadcast timeouts and driver OOMs when it goes wrong.

for a principal

Set the policy: which pipelines may hint broadcasts, what the cluster-wide threshold should be, and how dimension tables are kept small enough to broadcast rather than growing until every join becomes a shuffle.

## Two ways to execute an equi-join Spark has two main shapes for a join on equal keys. A **shuffle join** — sort-merge join being the usual physical operator — inserts an `Exchange` under each side, hash-partitioning both by the join key so that matching keys land in the same reducer partition, sorts each side, and merges. It works at any scale and it is the only option when both sides are large. A **broadcast hash join** takes the smaller side, collects it to the driver, broadcasts it to every executor, builds an in-memory hash table from it, and then streams the large side through, probing per row. The large side is never moved. ## Why this kills skew Skew in a join is created by the redistribution step, not by the data itself. If 40% of your fact rows carry `customer_id = -1`, hashing that key sends 40% of the rows to a single partition, and one task must read and join them all. Broadcasting removes the redistribution entirely: the large side's rows stay in whatever partitions the scan produced, which are sized by bytes read rather than by key, so 40% of the rows are spread across 40% of the tasks. Each of them probes its local copy of the hash table. The skewed key is still skewed in the data, but no operator concentrates it. This is why "can I broadcast one side?" is the first question to ask about a skewed join, before salting or any AQE tuning. It is the cheapest fix available and it removes the shuffle cost too. ## How Spark decides `spark.sql.autoBroadcastJoinThreshold` sets the size in bytes below which a relation is broadcast automatically; the default is 10485760 (10 MB), and `-1` turns automatic broadcasting off. The comparison uses the planner's size estimate, which comes from data-source metadata or catalog statistics collected by `ANALYZE TABLE`. Estimates on a filtered or joined intermediate are often wrong, which is why the decision can be poor and why AQE re-checks it at runtime. You can force the choice with the `broadcast()` function in the DataFrame API, the `.hint("broadcast")` method, or the `BROADCAST` hint in SQL. A hint takes priority over the size estimate, and Spark prioritises `BROADCAST` above `MERGE`, `SHUFFLE_HASH` and `SHUFFLE_REPLICATE_NL` when several hints are present. A hint is not a guarantee: if the join type does not support broadcasting the requested side, Spark ignores it. ## The costs and failure modes - **The driver is a bottleneck.** The build side is materialised in driver memory before it is sent. Broadcasting a multi-gigabyte relation typically shows up as driver GC thrashing, a driver OOM, or a `spark.sql.broadcastTimeout` (300 seconds by default) expiring while the broadcast is being prepared. - **Every executor pays memory.** One copy of the hash table lives in each executor, competing with execution memory. Broadcasting a 500 MB table across 100 executors means 50 GB of cluster memory holding the same rows. - **Raising the threshold blindly is dangerous.** Bumping `spark.sql.autoBroadcastJoinThreshold` to hundreds of megabytes makes Spark broadcast far more often, including relations whose size was under-estimated. ## When you cannot broadcast The build side is constrained by join type. In a `LEFT OUTER` join the left side must be streamed, so only the right side can be the broadcast build side (and mirror-image for `RIGHT OUTER`); a `FULL OUTER` equi-join cannot use a broadcast hash join at all. And if both sides are genuinely large, broadcasting is off the table and you fall back to AQE's skew handling, salting, or isolating the hot key. ## A useful intermediate If the small side is small only *after* a filter or aggregation that Spark cannot estimate, materialise it — cache it, or write and re-read it — so the size is known and the broadcast decision is made on a real number. AQE gives you a second chance at this automatically by re-planning once the real statistics exist. ## Reading the plan A broadcast join appears as `BroadcastHashJoin` with a `BuildLeft`/`BuildRight` marker and a `BroadcastExchange` under the build side only. A shuffle join appears as `SortMergeJoin` with an `Exchange hashpartitioning` under **both** sides. Confirming which one you got from `df.explain()` takes seconds and settles most arguments about why a join is slow.

  • Is raising spark.sql.autoBroadcastJoinThreshold to 1 GB a reasonable way to get more broadcast joins?
    Rarely. The threshold is compared against an estimate, so a high value makes Spark broadcast relations whose real size is much larger, and every executor then holds a copy while the driver materialises it first. Prefer an explicit broadcast hint on the joins you have measured, and leave the global threshold near its default.
  • Which side gets broadcast in a LEFT OUTER join?
    Only the right side. The left side must be streamed so that its unmatched rows can still be emitted with nulls, which means it cannot be the build side. RIGHT OUTER is the mirror image, and a FULL OUTER equi-join cannot use a broadcast hash join at all, since both sides need their unmatched rows preserved.
  • Your small table is small only after a filter, but Spark still plans a sort-merge join. What do you do?
    The planner is working from a stale or missing estimate of the filtered relation. Either add an explicit broadcast hint, materialise the filtered side so its real size is known, or rely on AQE, which re-checks the actual size once the side's shuffle has been written and can convert the join to a broadcast hash join at runtime.

Instead of sending every customer to one clerk who happens to hold their file, you photocopy the whole (thin) file cabinet for every clerk and let each of them serve whoever is already in front of them.

saying these in an interview costs you the question

  • Claims broadcasting shuffles the small table to matching partitions
  • Thinks the threshold compares against exact sizes, not estimates
  • Ignores that the driver materialises the build side first
  • Says broadcast joins work for any join type and either side
  • Raises the broadcast threshold to gigabytes as a routine fix

context