A broadcast hint is one comment in a query, /*+ BROADCAST(d) */, and it is the cheapest way to delete a shuffle from a Spark job. It is also a common way to crash a driver. This article covers what happens after the planner accepts the hint: where the build-side rows go, how much memory they take at each stage, the hard limits and the timeout, and how to decide between a hint, a threshold change and letting Adaptive Query Execution choose. How Spark picks among its four join strategies, and the precedence rules between different hints, are covered in Spark join strategy hints; this piece assumes them and goes deep on broadcast only.

Configuration names and defaults below are from the Apache Spark 4.x SQL performance tuning documentation.

What the hint does, and what it cannot do

Spark accepts BROADCAST, BROADCASTJOIN and MAPJOIN as the same hint. In the DataFrame API the equivalents are broadcast(df) from pyspark.sql.functions or df.hint("broadcast"). The hint does exactly one thing: it tells the planner to prefer a broadcast join with the hinted relation as the build side even if statistics put that relation above spark.sql.autoBroadcastJoinThreshold (10 MB by default). With an equality key Spark plans a broadcast hash join; without one, a broadcast nested loop join, which compares every probe row with every build row.

It does not change legality. The build side must be the side whose unmatched rows are not kept, so in a LEFT JOIN b only b can be broadcast. Hinting a is ignored and the join falls back to whatever the planner would otherwise choose, usually a sort-merge join. Spark logs a warning for hints it cannot apply, which is the first place to look when a hint seems to do nothing.

from pyspark.sql.functions import broadcast

sales = spark.table("fact_sales")                       # billions of rows
stores = spark.table("dim_store").where("is_active")    # ~40k rows

joined = sales.join(broadcast(stores), "store_id")      # build side: stores

# SQL form: reference the alias used in FROM, not the table name
spark.sql("""
  SELECT /*+ BROADCAST(s) */ f.sale_id, s.region, f.amount
  FROM fact_sales f JOIN dim_store s ON f.store_id = s.store_id
  WHERE s.is_active
""")

The life of a broadcast

What one BROADCAST hint sets in motionBuild-side queryscan + filter dim tableCollect to driverrows as UnsafeRow bytesBuild HashedRelationon the driver heapjob 1TorrentBroadcastchunked, compressed blocksExecutor 1one copy, sharedExecutor 2one copy, sharedExecutor 3one copy, sharedExecutor 4one copy, sharedProbe side: big table, no shuffleeach task streams its partition through the hash tableThe timeout clock covers job 1, the collect and the build; the driver must hold the rows and the relation together.
Figure 1. The build side is computed by its own job, collected to the driver, turned into a hash relation and shipped to every executor once.

The physical operator is BroadcastExchangeExec. When the join first needs its input, it starts a separate Spark job that computes the build side, including any scans, filters and upstream joins. The resulting rows are collected to the driver as serialised binary rows. The driver then builds a HashedRelation: a specialised long-keyed map when the key is a single integral column, and a general binary-keyed map otherwise. That relation is handed to Spark's torrent broadcast, which splits it into compressed blocks; executors fetch blocks from the driver and from each other, so the driver is not the only source after the first copies spread.

Each executor keeps one deserialised copy in its block manager, shared by all tasks on that executor, and the join runs as a narrow stage: every task streams its partition of the large table through the local hash table. No shuffle of the large side happens at all, which is the whole point.

Two consequences matter. The driver holds the collected rows and the built relation at the same time during construction, so its peak is a multiple of the data size. And every executor pays memory for the relation whether its tasks need most of it or not.

Sizing the build side

The hint bypasses the size check, so you are the size check. Spark's estimate comes from statistics: for files, roughly the on-disk size (scaled by spark.sql.sources.fileCompressionFactor, which defaults to 1.0); for a filtered relation without cost-based optimisation, often the unfiltered size or a crude guess. Compressed columnar files are typically several times smaller than the same rows held as binary rows in memory, and a hash relation adds its own index on top. Measure rather than trust the estimate:

dim = spark.table("dim_customer").where("country = 'DE'")

# 1. What the planner believes (bytes); classic sessions only, not Spark Connect
est = dim._jdf.queryExecution().optimizedPlan().stats().sizeInBytes()

# 2. A lower bound: the in-memory cache is compressed columnar, smaller than the broadcast
dim.cache().count()
# Spark UI > Storage shows "Size in Memory" for the cached relation

# 3. Confirm the plan actually broadcasts
facts.join(broadcast(dim), "customer_id").explain()
#  ... BroadcastHashJoin [customer_id], [customer_id], Inner, BuildRight
#  ... +- BroadcastExchange HashedRelationBroadcastMode(...)
dim.unpersist()

# 4. The real figure: run once in staging, then open Spark UI > SQL, click the query and read
#    the BroadcastExchange node: data size, time to collect, time to build, time to broadcast

The cache figure is useful as a quick floor, but it understates the broadcast, because the cache stores compressed columns while the broadcast holds uncompressed binary rows plus a hash index. The number to plan with is the data size metric on the exchange node in the SQL tab, captured from one staging run on production-sized data.

Worked sizing. A customer dimension has 30 million rows and six columns and occupies 1.1 GB as Parquet. A staging run shows the exchange's data size as 4.2 GB. During the build the driver holds the collected rows plus the relation, so plan for roughly two to three times the in-memory size, call it 10 to 12 GB of driver heap headroom for this one join. Executors each hold one copy: with 4.2 GB per copy and 16 GB executors, a quarter of every executor's memory disappears before any task runs. The hint is legal; the better plan is to project only the two columns the query uses and filter first, which in this case cut the broadcast data size to 700 MB and made the broadcast comfortable.

Hard limits and the timeout

Beyond memory, three hard edges apply. Spark refuses to broadcast a relation larger than 8 GB or with 512 million rows or more, failing the query rather than silently switching strategy. And spark.sql.broadcastTimeout (300 seconds by default) bounds how long the join waits for the broadcast to be ready.

The timeout is the misunderstood one. Its clock includes computing the build side, not just shipping it. A small dimension produced by an expensive subquery, a slow JDBC source, or a stage stuck behind other jobs on a busy cluster times out even though the final relation is tiny. Raising the timeout hides that; materialising the build side first (write it out, or cache and count it in a prior step) fixes it, because the broadcast job then reads ready data.

Large broadcasts can also surface as driver-side result-size errors or long garbage-collection pauses before any explicit failure. If driver GC time climbs during a stage that should be fast, check whether a broadcast is being built.

Hints, thresholds and AQE

With Adaptive Query Execution enabled (the default since Spark 3.2), the planner gets a second chance after shuffle stages finish and real sizes are known. If a sort-merge join's input turns out smaller than spark.sql.adaptive.autoBroadcastJoinThreshold (which defaults to the static threshold), AQE rewrites it to a broadcast hash join, and with spark.sql.adaptive.localShuffleReader.enabled (default true) reads the already-shuffled probe side locally instead of reshuffling it. See Spark AQE architecture for the mechanics.

That changes when a hint is worth writing. AQE can only promote a join after a shuffle has produced runtime statistics, so it saves the join's work but not the shuffle that fed it. A hint acts at planning time and removes the shuffle of the large side entirely. Hints are respected under AQE; the hinted join is planned as a broadcast from the start.

SituationBest toolWhy
Dimension reliably small, stats missing or staleHintRemoves the shuffle up front; stats cannot mislead
Many dimensions, all well under a known sizeRaise the static threshold moderatelyOne setting covers every query
Build side small only after selective filtersLet AQE decideRuntime sizes are accurate; no guessing
Build side grows over timeAQE, or a hint with a size assertionA static hint becomes an OOM later
Non-equi join with a small sideHint with careNested loop cost grows with probe rows x build rows

Where hints get lost

Hints attach to relation names as the analyser resolves them, and a few patterns lose them. A hint must name the alias visible in that FROM clause; naming the base table when it is aliased does nothing. A hint inside a view definition applies inside the view, while a hint naming the view from outside applies to the whole view's result. Hints are resolved during analysis, but the planner applies them to whichever join survives optimisation, after filters move, subqueries are rewritten and joins are reordered, so the join you hinted may not be the physical join you get. Always confirm with explain() that BroadcastHashJoin appears with the build side you expect.

Guarding a hint against data growth

A hint encodes an assumption about data size that was true on the day it was written. Make the assumption executable, so the job fails with a clear message instead of an out-of-memory error on the driver months later. Counting rows is cheap for a small dimension and catches most growth; for wide rows, assert on the cached size measured during development and re-measure when the schema changes.

from pyspark.sql.functions import broadcast

MAX_BROADCAST_ROWS = 5_000_000   # chosen from measured in-memory size, not from disk size

def broadcast_checked(df, name):
    """Broadcast df only while it is still small; otherwise let the planner decide."""
    n = df.count()
    if n > MAX_BROADCAST_ROWS:
        print(f"WARN {name}: {n} rows, skipping broadcast hint")
        return df
    return broadcast(df)

stores = spark.table("dim_store").where("is_active").select("store_id", "region")
joined = spark.table("fact_sales").join(broadcast_checked(stores, "dim_store"), "store_id")

The extra count() runs a small job over the dimension; cache the dimension first if you also use it elsewhere so the scan is not repeated. Falling back to an unhinted join keeps the pipeline running and still leaves AQE free to broadcast if the runtime size permits. Whether to fail or fall back is a policy choice: fail when the dimension growing tenfold signals a data bug, fall back when growth is expected.

Failure modes

  • Driver out of memory during the build: the relation was far larger in memory than on disk. Project and filter the build side, or drop the hint.
  • Broadcast timeout on a small table: the build side's own computation was slow. Materialise it first instead of raising spark.sql.broadcastTimeout.
  • Executor memory pressure and spilling across unrelated stages: several large broadcasts are resident at once. Check the Storage tab and executor memory metrics.
  • Nested-loop explosion: a broadcast on a join with no equality key compares every pair of rows. Add an equi-key, even a coarse bucketing key, before the range condition.
  • Hint ignored: illegal build side for the join type or an alias mismatch. Read the warning in the driver log.
  • Stale hint: a dimension that was 50 MB last year is 9 GB now and hits the hard limit. Guard hinted tables with a size check in the job.

What to do next

  1. List every broadcast hint in your codebase and record the measured in-memory size of each build side.
  2. Remove hints where the build side is filtered heavily at runtime and let AQE choose instead.
  3. For hints you keep, project only needed columns and add a row or size assertion before the join.
  4. Size driver memory for the largest broadcast at two to three times its in-memory size.
  5. Replace any timeout increase with materialisation of the build side.
  6. Review broadcast variables for the torrent mechanism, and sort-merge join and skew handling for the joins a broadcast replaces.
Key takeaway: A broadcast hint overrides the size threshold but not join legality, and once accepted it runs a separate job, collects the build side to the driver, builds a hash relation there and ships one copy to every executor. In-memory size is often several times the on-disk size, the driver needs two to three times that during the build, and Spark refuses relations over 8 GB or 512 million rows. Hint dimensions that are reliably small, let AQE decide when size depends on runtime filters, and fix timeouts by materialising the build side.