A join method that copies the smaller table to every executor so that the larger table needs no shuffle. You can force it with the BROADCAST hint (or BROADCASTJOIN / MAPJOIN), or AQE may pick it while the query runs.
Also called broadcast hash join.
Read more: Microsoft Learn
In the Ultra Transcenders books
Each book explains Broadcast join in context, with comparison tables and the common traps.
Terms in this definition
- JOIN
Combines rows from two or more tables in one SELECT, most often by pairing a primary key with the foreign key that refers to it.
- Event
Table in Log Analytics where entries from Windows event logs are kept.
- Executor
Runs tasks on a worker node; Azure Databricks places exactly one on each worker. Autoscaling, losing a spot instance or out-of-memory errors can all take one away.
- Shuffle
When Spark has to move rows between executors, as joins, sorts and aggregations require. Per stage, the Spark UI lists Shuffle Read and Shuffle Write; heavy shuffling frequently causes spill.
- AQE
Adaptive query execution: Spark re-plans a query while it runs, for instance turning a sort-merge join into a broadcast join, merging shuffle partitions, dealing with skew and propagating empty relations. Databricks advises leaving it enabled.
- Agents (classic) API
First-generation Foundry Agent Service API, based on threads, messages and runs. It is deprecated, replaced by conversations and responses, and retires on 31 March 2027.
Related terms
- Autotune
A preview option in Fabric Spark that studies previous runs to pick better values for shuffle partitions, broadcast join threshold and max partition bytes query by query. It's disabled by default (
spark.ms.autotune.enabled), limited to Runtime 1.2, and incompatible with private endpoints and high concurrency mode.