What you will be able to do
- Explain why broadcasting a small table to every executor makes a join faster than a shuffle join
- Mark a DataFrame for a broadcast join with pyspark.sql.functions.broadcast() or DataFrame.hint("broadcast")
- Explain how spark.sql.autoBroadcastJoinThreshold controls automatic broadcasting, and how to turn it off
- Write SQL broadcast hints and predict which hint wins when hints conflict
- Recognise when Spark will not broadcast, and how Adaptive Query Execution changes join strategy at runtime
Key concept
Broadcast join — A join where Spark copies the smaller relation in full to every executor and holds it in memory as a hash table. Each partition of the larger relation can then be matched locally, so the larger side does not need a shuffle join.
1.Why broadcast a table at all?
A typical join pairs a large fact table with a small lookup table, such as country codes or product categories. In a shuffle join, rows from both sides are redistributed across the cluster so that matching keys end up on the same partition. That redistribution moves a lot of data over the network, even though one side is small.
A broadcast join avoids that. Spark sends a complete copy of the small table to every executor and keeps it in memory as a hash map. Each executor then reads its own partitions of the large table and looks up matches in its local copy. The Databricks documentation describes this pattern for a small lookup table: the optimizer turns the lookup into a broadcast hash join, and the table is copied to each executor and stored in memory as a hashmap.
So the purpose of a broadcast join is specific: it pays off when one side is small. The Databricks AQE blog puts it this way: among Spark's join strategies, broadcast hash join is usually the most performant if one side of the join fits well in memory. If neither side is small, copying one of them to every executor costs more than it saves, and a shuffle join is the right choice.
Broadcasting does not make every row free to process. Each row of the large table still needs a hash-table lookup. The join condition also matters: the same documentation warns that a complex lookup predicate, meaning anything other than a simple equality check, can make broadcast join ineligible.
Checkpoint 1 of 8· Check yourself
Which situation best fits a broadcast join?
Broadcast hash join pays off when one side fits well in memory and the join uses an equality key. Two huge tables, or a complex predicate, are the cases where it stops being the best choice or stops being eligible.
“broadcast hash join is usually the most performant if one side of the join can fit well in memory”Source: www.databricks.com
2.Requesting a broadcast in the DataFrame API
In PySpark you request a broadcast explicitly with pyspark.sql.functions.broadcast(df). Its documented purpose is to mark a DataFrame as small enough for use in broadcast joins. It takes one DataFrame and returns a DataFrame marked for broadcast. It does not run a join itself; you pass the marked DataFrame to an ordinary join. The function is new in version 1.6.0 and has supported Spark Connect since 3.4.0.
from pyspark.sql import functions as dbf
df = spark.createDataFrame([1, 2, 3, 3, 4], "int")
df_small = spark.range(3)
df_b = dbf.broadcast(df_small)
df.join(df_b, df.value == df_small.id).show()There is a second way to ask for the same thing: DataFrame.hint(name, *parameters), which returns a hinted DataFrame. Passing "broadcast" as the hint name marks that DataFrame for broadcasting. The hint reference shows how to confirm it worked. Compare the physical plans from explain(): without the hint the two small DataFrames are joined with SortMergeJoin, and with the hint the plan shows BroadcastHashJoin.
df.join(df2, "name").explain()
# == Physical Plan ==
# ...
# ... +- SortMergeJoin ...
# ...
df.join(df2.hint("broadcast"), "name").explain()
# == Physical Plan ==
# ...
# ... +- BroadcastHashJoin ...
# ...Checkpoint 2 of 8· Fill the gap
Which function completes this sample so that df_small is marked for a broadcast join?
from pyspark.sql import functions as dbf
df = spark.createDataFrame([1, 2, 3, 3, 4], "int")
df_small = spark.range(3)
df_b = dbf. ? (df_small)
df.join(df_b, df.value == df_small.id).show()broadcast(df) is the function in pyspark.sql.functions that marks a DataFrame as ready for a broadcast join. hint is a method on DataFrame, not a function in the functions module.
Source: docs.databricks.comCheckpoint 3 of 8· Exam question
Why does a broadcast join avoid the shuffle stage that a standard join between two large tables requires?
Correct answer: A — The smaller dataset is copied in full to every executor, letting each executor complete the join locally without exchanging the larger dataset's partitions.
- A. This is correct: a broadcast join sends a full copy of the small side to every executor, so the join against the large side's local partitions never needs to move rows across the network. That removal of the shuffle exchange is the entire performance benefit of the strategy.
- B. Spark does not build a compressed columnar cache of the larger table on each executor; the join strategy is about sending the small side, not compressing the large side. This describes a caching mechanism Spark does not use for broadcast joins.
- C. Pre-sorting both sides and merging matching rows describes a sort-merge join, a separate strategy Spark uses when neither side is small enough to broadcast. It is not how a broadcast join avoids network movement.
- D. Hash-partitioning both sides by join key and still exchanging every partition describes a shuffle hash join, which keeps the shuffle exchange this question asks about avoiding. A broadcast join is chosen precisely to skip that exchange.
- E. Spark has no shared distributed cache layer that both executors read from during a join; the broadcast mechanism physically pushes a copy of the data to each executor's local memory instead. This option invents a caching tier that does not exist in Spark's join execution.
3.Automatic broadcasting: spark.sql.autoBroadcastJoinThreshold
You do not always have to ask. During planning, Spark estimates the size of each join relation. If an estimate falls below a size threshold, Spark plans a broadcast hash join on its own. That threshold is the configuration property spark.sql.autoBroadcastJoinThreshold.
| Property | Default | Meaning | Since |
|---|---|---|---|
| spark.sql.autoBroadcastJoinThreshold | 10485760 (10 MB) | Maximum size in bytes for a table that will be broadcast to all worker nodes when performing a join; -1 disables broadcasting | 1.1.0 |
| spark.sql.broadcastTimeout | 300 | Timeout in seconds for the broadcast wait time in broadcast joins | 1.3.0 |
No. The threshold only controls *automatic* broadcasting. An explicit BROADCAST hint takes priority over it. When the hint is on table t1, Spark prioritizes a broadcast join with t1 as the build side, even if the statistics put t1 above spark.sql.autoBroadcastJoinThreshold. The resulting plan is a broadcast hash join when there is an equi-join key, and a broadcast nested loop join when there is not.
The automatic decision is only as good as the size estimate. The AQE blog notes that estimates can go wrong, for example after a very selective filter or when the relation comes from a series of complex operators rather than a plain scan. The final section shows how Spark corrects for that at runtime.
Checkpoint 4 of 8· Check yourself
You want Spark never to choose a broadcast join automatically. What should you set?
The documented way to disable broadcasting is to set the threshold to -1. broadcastTimeout only limits how long Spark waits for a broadcast, and disabling AQE does not turn off static broadcast planning.
“By setting this value to -1, broadcasting can be disabled.”Source: spark.apache.org
Checkpoint 5 of 8· Exam question
A Spark job automatically uses a broadcast hash join for a small dimension table with no explicit hint applied anywhere in the code. Which configuration property controls this automatic behavior, and what does its default value represent?
Correct answer: A — `spark.sql.autoBroadcastJoinThreshold` sets the largest estimated table size Spark broadcasts automatically, defaulting to 10 MB unless the session overrides it.
- A. This is correct: `spark.sql.autoBroadcastJoinThreshold` is the size cutoff the planner compares a table's estimated size against, and its default of 10 MB is why small dimension tables get broadcast without any explicit hint. Setting it to -1 disables this automatic behavior entirely.
- B. `spark.sql.shuffle.partitions` only controls the number of partitions used after a shuffle occurs and has nothing to do with the size threshold that triggers automatic broadcasting. Its default of 200 is unrelated to table size estimation.
- C. `spark.sql.broadcastTimeout` governs how long a broadcast exchange is allowed to run before timing out, not what size table qualifies for broadcasting in the first place. Its 300-second default is a wait limit, not a size limit.
- D. `spark.driver.maxResultSize` limits how much data a driver can collect from actions like `collect()`, and it is not the property the planner consults when deciding whether a table is small enough to broadcast. Raising it does not change broadcast eligibility.
- E. Automatic size-based broadcasting is independent of `spark.sql.adaptive.enabled` and works whether or not adaptive query execution is turned on. AQE can additionally convert a join to broadcast at runtime, but it is not a prerequisite for the size-based default behavior.
4.Broadcast hints in SQL and how conflicting hints are resolved
In SQL the same request is written as a comment-style hint inside the SELECT. Three names are accepted for the broadcast hint: BROADCAST, BROADCASTJOIN and MAPJOIN. They all ask for the same strategy. The hint takes the name of the table to broadcast. If that table has an alias in the query, the hint must use the alias.
> SELECT /*+ BROADCAST(t1) */ * FROM t1 INNER JOIN t2 ON t1.key = t2.key;
> SELECT /*+ BROADCASTJOIN (t1) */ * FROM t1 left JOIN t2 ON t1.key = t2.key;
> SELECT /*+ MAPJOIN(t2) */ * FROM t1 right JOIN t2 ON t1.key = t2.key;These examples show syntax only. Whether Spark can honour a hint depends on the join type, which the next section covers. Broadcast is one of four join strategy hints. When the two sides of a join carry different strategy hints, a fixed priority order decides which one wins. Spark logs a warning that the losing hint is overridden and will not take effect.
| Hint (and aliases) | Strategy requested |
|---|---|
| BROADCAST / BROADCASTJOIN / MAPJOIN | Broadcast join |
| MERGE / SHUFFLE_MERGE / MERGEJOIN | Shuffle sort merge join |
| SHUFFLE_HASH | Shuffle hash join |
| SHUFFLE_REPLICATE_NL | Shuffle-and-replicate nested loop join |
If both sides carry the BROADCAST hint, Spark picks the build side from the join type and the sizes of the relations.
Checkpoint 6 of 8· Match them up
Match each hint to the join strategy it requests
Tap a term, then the definition that fits it.
MAPJOIN and BROADCASTJOIN are aliases for BROADCAST. MERGE asks for a sort merge join, SHUFFLE_HASH for a shuffle hash join, and SHUFFLE_REPLICATE_NL for a nested loop join.
“We accept BROADCAST, BROADCASTJOIN and MAPJOIN for broadcast hint”Source: spark.apache.org
BROADCAST outranks MERGE, so Spark tries a broadcast join with t1. It logs a HintErrorLogger warning that the merge hint is overridden and will not take effect.
5.When Spark won't broadcast, and what AQE changes
The Spark tuning guide states the general rule: there is no guarantee that Spark will choose the join strategy specified in the hint, because a given strategy may not support every join type. The LEFT OUTER JOIN case above is the example to remember.
Adaptive Query Execution (AQE) adds a runtime layer on top of this. AQE has been enabled by default since Spark 3.2.0 and is controlled by spark.sql.adaptive.enabled. It re-optimizes the plan using statistics gathered at the end of shuffle and broadcast exchanges. One of its four main features is that it dynamically changes sort merge join into broadcast hash join. This fixes the estimation problem from earlier: if a relation turns out at runtime to be much smaller than estimated, a statically planned sort merge join can become a broadcast hash join.
Yes. Databricks recommends it when you know your query well. A statically planned broadcast join is usually faster than one AQE plans dynamically, because AQE might not switch to broadcast until it has already shuffled both sides to learn their real sizes. AQE respects hints the same way static optimization does.
AQE can also leave a relation unbroadcast even when it is under the threshold. One reason is the join type, as above. Another is that the relation has many empty partitions: AQE avoids the conversion if the share of non-empty partitions is below spark.sql.adaptive.nonEmptyPartitionRatioForBroadcastJoin. Finally, AQE's skew handling applies only to shuffle-based joins. Broadcast joins are never skew-optimized.
Checkpoint 7 of 8· Check yourself
With AQE enabled, which statement about broadcast joins is correct?
AQE may only switch to broadcast after shuffling both sides, so planning the broadcast up front with a hint is usually faster. AQE respects hints, and it never skew-optimizes broadcast joins.
“A statically planned broadcast join is usually more performant than a dynamically planned one by AQE”Source: docs.databricks.com
Checkpoint 8 of 8· Exam question
A data engineer joins a 500 GB fact table `orders` with a 2 MB dimension table `regions` on `region_id` using `orders.join(regions, "region_id")`. The job runs slowly because Spark chooses a sort-merge join instead of broadcasting `regions`. Which code change forces Spark to broadcast `regions` for this join?
Correct answer: A — ```python from pyspark.sql.functions import broadcast result = orders.join(broadcast(regions), "region_id") ```
- A. This is correct: wrapping `regions` with the `broadcast()` function explicitly marks it for a broadcast hash join, overriding whatever plan the optimizer would otherwise pick regardless of the table's estimated size. This is the standard DataFrame-API way to force the strategy.
- B. Calling `.cache()` on `regions` only persists it in memory for reuse across actions; it does not tell the join planner to use a broadcast strategy. The join could still execute as a sort-merge join.
- C. Repartitioning `orders` into 200 partitions changes how the large side is distributed for a shuffle-based join but does nothing to mark `regions` for broadcasting. This still leaves the choice of join strategy to the optimizer's default logic.
- D. Coalescing `regions` down to one partition changes its physical layout but does not instruct Spark to broadcast it; the optimizer still decides the join strategy based on its own cost estimates. Coalesce is a partition-count operation, not a join hint.
- E. Sorting `regions` within its partitions prepares data for a sort-based join rather than requesting a broadcast; it does not change which join strategy the planner selects. This action has no effect on broadcast eligibility.
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Spark ignores a broadcast() or BROADCAST hint when the table is larger than spark.sql.autoBroadcastJoinThreshold.Why is that wrong?
The threshold only governs automatic broadcasting. A BROADCAST hint gives priority to a broadcast join with that table as the build side, even when its estimated size is above the threshold.
Covered in Automatic broadcasting: spark.sql.autoBroadcastJoinThreshold
2.A broadcast hint always forces a broadcast join, whatever the join type.Why is that wrong?
Hints are not guaranteed. A strategy may not support every join type; for example, the left relation of a LEFT OUTER JOIN cannot be broadcast.
3.With AQE enabled, broadcast hints are pointless because AQE picks the best join at runtime.Why is that wrong?
AQE may only switch to broadcast after shuffling both sides, so a statically planned broadcast join from a hint is usually more performant. AQE respects hints.
4.AQE's skew-join handling also rebalances skewed partitions in broadcast joins.Why is that wrong?
AQE skew handling applies only to shuffle-based joins (sort merge and shuffle hash). Broadcast joins are never skew-optimized.
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“The lookup table is copied to each executor and stored in memory as a hashmap, which enables fast filtering during the table scan.”
↩︎ Why broadcast a table at all?“If the lookup predicate is complex (not a simple equality check), broadcast join may also become ineligible.”
↩︎ Why broadcast a table at all?“Even with broadcast hash join, each row still incurs the cost of a hash table lookup during execution.”
↩︎ Why broadcast a table at all?“The lookup table is copied to each executor and stored in memory as a hashmap, which enables fast filtering during the table scan.”
↩︎ Key concept“If the lookup table is large, the optimizer falls back to a shuffle join, which is slower.”
↩︎ Prediction - 2.https://www.databricks.com/blog/2020/05/29/adaptive-query-execution-speeding-up-spark-sql-at-runtime.htmlSecondary source
“broadcast hash join is usually the most performant if one side of the join can fit well in memory”
↩︎ Why broadcast a table at all?“Spark plans a broadcast hash join if the estimated size of a join relation is lower than the broadcast-size threshold.”
↩︎ Automatic broadcasting: spark.sql.autoBroadcastJoinThreshold - 3.
“Marks a DataFrame as small enough for use in broadcast joins.”
↩︎ Requesting a broadcast in the DataFrame API - 4.https://spark.apache.org/docs/latest/api/python/reference/pyspark.sql/api/pyspark.sql.functions.broadcast.htmlSecondary source
“Changed in version 3.4.0: Supports Spark Connect.”
↩︎ Requesting a broadcast in the DataFrame API - 5.
“Specifies some hint on the current DataFrame.”
↩︎ Requesting a broadcast in the DataFrame API - 6.https://spark.apache.org/docs/latest/sql-performance-tuning.htmlSecondary source
“Configures the maximum size in bytes for a table that will be broadcast to all worker nodes when performing a join.”
↩︎ Automatic broadcasting: spark.sql.autoBroadcastJoinThreshold“broadcast join (either broadcast hash join or broadcast nested loop join depending on whether there is any equi-join key)”
↩︎ Automatic broadcasting: spark.sql.autoBroadcastJoinThreshold“Spark prioritizes the BROADCAST hint over the MERGE hint over the SHUFFLE_HASH hint over the SHUFFLE_REPLICATE_NL hint.”
↩︎ Broadcast hints in SQL and how conflicting hints are resolved“When both sides are specified with the BROADCAST hint or the SHUFFLE_HASH hint, Spark will pick the build side based on the join type”
↩︎ Broadcast hints in SQL and how conflicting hints are resolved“which is enabled by default since Apache Spark 3.2.0”
↩︎ When Spark won't broadcast, and what AQE changes“will be prioritized by Spark even if the size of table ‘t1’ suggested by the statistics is above the configuration spark.sql.autoBroadcastJoinThreshold”
↩︎ Exam trap 1“there is no guarantee that Spark will choose the join strategy specified in the hint”
↩︎ Exam trap 2“By setting this value to -1, broadcasting can be disabled.”
↩︎ Checkpoint“We accept BROADCAST, BROADCASTJOIN and MAPJOIN for broadcast hint”
↩︎ Checkpoint - 7.https://docs.databricks.com/aws/en/optimizations/aqeOfficial docs
“If the size of the relation expected to be broadcast does fall under this threshold but is still not broadcast:”
↩︎ Automatic broadcasting: spark.sql.autoBroadcastJoinThreshold“Dynamically changes sort merge join into broadcast hash join.”
↩︎ When Spark won't broadcast, and what AQE changes“Broadcast is not supported for certain join types, for example, the left relation of a LEFT OUTER JOIN cannot be broadcast.”
↩︎ When Spark won't broadcast, and what AQE changes“AQE might not switch to broadcast join until after performing shuffle for both sides of the join”
↩︎ When Spark won't broadcast, and what AQE changes“AQE will respect query hints the same way as static optimization does”
↩︎ Exam trap 3“Broadcast joins are never skew-optimized.”
↩︎ Exam trap 4“A statically planned broadcast join is usually more performant than a dynamically planned one by AQE”
↩︎ Checkpoint - 8.
“When a table name is occluded by an alias you must use the alias name in the hint”
↩︎ Broadcast hints in SQL and how conflicting hints are resolved