## TL;DR
Raise `spark.sql.autoBroadcastJoinThreshold` when a small dimension table is not broadcasting, lower it when an oversized broadcast OOMs the driver, and use the `broadcast()` hint for per-query control. Broadcast is a driver-memory decision, not an executor one.

```text
spark broadcast join threshold tuning
```

## Use this when
- A join on a small table runs as a slow sort-merge instead of broadcasting
- A broadcast join OOMs the driver
- You want one query to broadcast without changing the global setting

## Not for this skill when
- Executors OOM during a shuffle, thats partition sizing
- Column resolution fails at analysis time
- Timestamps parse to null

## Steps

1. Check the current threshold and whether the small side actually fits:

```python
print(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
dim.count(), dim.rdd.map(lambda r: len(str(r))).sum()  # rough size check
```
Expected output: the threshold in bytes (default 10MB, shown as 10485760) and a rough size of the dimension table. Broadcast is only considered when the small side is under the threshold.

2. Raise the threshold when the dimension table is small but over 10MB:

```python
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "104857600")  # 100MB
```
Expected output: set before the query. Joins against tables under 100MB now broadcast automatically. The broadcast data is collected on the driver first, so keep it comfortably under driver memory.

3. Force or forbid broadcast per query with hints, without touching the global config:

```python
from pyspark.sql.functions import broadcast
fact.join(broadcast(dim), "id")          # force
fact.join(dim.hint("shuffle_hash"), "id")  # forbid broadcast, use shuffle hash
```
Expected output: the physical plan in the UI shows BroadcastHashJoin or SortMergeJoin as hinted. Hints are the right tool when one query is special.

4. If the driver OOMs during broadcast, the table was bigger than you thought. Check the real size:

```python
dim.cache().count()
# then check Storage tab in the UI for the cached size
```
Expected output: the true in-memory size, often several times the on-disk size after deserialization. Lower the threshold or drop the hint so the join goes sort-merge.

5. Confirm the plan in the UI or with explain:

```python
fact.join(dim, "id").explain()
```
Expected output: `BroadcastHashJoin` when broadcasting, `SortMergeJoin` otherwise. Trust the plan, not your assumption about which side Spark picked.

## Variant phrasings

### spark broadcast join not happening
Either the small side exceeds the threshold (step 2) or stats are missing so Spark overestimates its size. Run ANALYZE or check `spark.sql.statistics` handling.

### spark driver out of memory broadcast
The broadcast table was too big for the driver (step 4). This is the failure mode of setting the threshold too high.

### spark disable broadcast join
Set the threshold to -1 to disable auto-broadcast globally, or use the `shuffle_hash` / `shuffle_merge` hints per query.

## Why it happens
A broadcast join sends the small table to every executor once, avoiding the shuffle of the big table entirely. Spark decides based on size estimates versus the threshold, collecting the small side on the driver. Wrong threshold in either direction gives you the worst of both: a shuffle you did not need, or a driver OOM.

## Edge cases
- AQE can flip the join strategy at runtime based on actual sizes, keep it enabled and treat the threshold as a starting hint.
- Skewed big-side keys still hurt broadcast joins less than sort-merge, but the broadcast side must still fit in driver memory.
- Setting the threshold absurdly high (gigabytes) is a classic driver-OOM footgun, size it from measured table sizes.
- `spark.sql.broadcastTimeout` matters when the broadcast collection is slow, not just large.

## Provenance

Resolved from the public thread: https://vectle.com/posts/pst_Z_oMrP9jMXELsxvIXUHVuQ
