Databricks Certified Data Engineer Associate: PySpark Joins at Scale
Large PySpark joins are rarely slow because joining is inherently expensive. They are slow because the two sides have inconvenient size, distribution, cardinality, or filtering characteristics. A tiny dimension table may be shuffled unnecessarily. A single hot key may place most of the work in one partition. A many-to-many relationship may multiply rows unexpectedly. A full outer join may prevent optimizations that work for an inner join. Tuning begins by understanding that shape.
Within Databricks Lakehouse Engineering, join performance is a bridge between Spark execution and table design. The same query can behave very differently depending on statistics, clustering, file layout, filters, Adaptive Query Execution, and whether one relation can be broadcast. The right question is not ‘Which join hint is fastest?’ but ‘What does the optimizer know, and what data movement does this relationship require?’
Estimate relation size after filtering
Join strategy should be based on the data that reaches the join, not the raw size of the source table. Select only needed columns and apply selective filters before the join when semantics allow it. A dimension table that is large on disk may become small enough to broadcast after filtering to one region or active date range.
Use broadcast when one side is genuinely small
A broadcast hash join avoids shuffling the larger side by sending the smaller relation to executors. Spark can choose broadcast automatically from statistics and thresholds, and Adaptive Query Execution can switch to broadcast at runtime when actual sizes permit. A deliberate broadcast hint can still be useful when the engineer knows the filtered relation is safely small.
Keep Adaptive Query Execution enabled
Databricks recommends Adaptive Query Execution because it can coalesce shuffle partitions, change eligible sort-merge joins to broadcast hash joins, handle skewed shuffle partitions, and propagate empty relations during execution. Those decisions use runtime evidence that was unavailable when the initial plan was created.
Recognize skew before adding partitions
Data skew appears when a small number of join keys contain a disproportionate share of rows. Increasing the global partition count may create more tiny partitions while the hot key remains concentrated in one or a few expensive tasks. Stage metrics often reveal a long tail where most tasks finish quickly and a few continue for much longer.