skip to content

Spark SQL and Catalyst

Catalyst turns DataFrame code into an optimized physical plan — resolving names, rewriting the logical plan, choosing join strategies, and generating code. Reading explain() output is a standard interview exercise because it shows whether you can diagnose a slow query instead of guessing at it.

on this pageshow

explore

questions

6

How does Spark's planner choose between BroadcastHashJoin and SortMergeJoin?

level: middleimportance: must knowfreq 64%

answer

  1. hint first, then size, then shuffle
  2. ten megabytes is the default line
  3. the estimate is compressed bytes on disk
  4. which side may be replicated depends on the join type
  5. no equality means nested loops

basics

~20 s

Spark broadcasts when a hint says so or when one side's estimated size fits under spark.sql.autoBroadcastJoinThreshold, which defaults to 10 MB; otherwise an equi-join on sortable keys becomes a SortMergeJoin that shuffles and sorts both sides.

solid answer

~40 s

The `JoinSelection` strategy works down a priority list. A **hint wins first**: `broadcast(df)` or `/*+ BROADCAST(t) */` forces a `BroadcastHashJoin` if the join type allows it. Otherwise, if one side's **estimated** size is below `spark.sql.autoBroadcastJoinThreshold` (default 10 MB), Spark collects that side to the driver, ships it to every executor and hash-joins locally — no shuffle at all. If neither side is small enough and the join is an equi-join on orderable keys, it plans a `SortMergeJoin`: both sides are shuffled by hash of the key, sorted, and merged. `ShuffledHashJoin` is the rarer middle option, preferred only when sorting is unattractive. A non-equality condition cannot use either and falls back to `BroadcastNestedLoopJoin` or `CartesianProduct`. The join type constrains the build side — for a left outer join only the right side may be broadcast.

code

python · 13 lines
python
from pyspark.sql.functions import broadcast

# force the small side, when the planner's estimate cannot see the filter
joined = fact.join(broadcast(dim.filter("active = true")), "dim_id")

# the SQL equivalent
spark.sql("""
    SELECT /*+ BROADCAST(d) */ f.*, d.name
    FROM fact f JOIN dim d ON f.dim_id = d.dim_id
""")

# give the planner real numbers instead of file sizes
spark.sql("ANALYZE TABLE dim COMPUTE STATISTICS FOR COLUMNS dim_id")

go deeper

for a junior

Know the two main join operators by name and that Spark broadcasts a small table to avoid a shuffle. Be able to spot BroadcastHashJoin versus SortMergeJoin in a printed plan.

for a middle

Explain the selection order — hint, then size against the 10 MB default threshold, then sort-merge — and what each strategy does to the data physically.

for a senior

Reason about where the size estimate comes from, why compressed file size misleads, and when to fix statistics rather than pin a hint. Recognize a broadcast timeout or driver OOM as the same root cause.

for a principal

Own the policy: whether the platform raises the threshold globally or relies on targeted hints, how table statistics get refreshed, and how hints are reviewed so they do not silently outlive the data that justified them.

## The decision, in order During physical planning, Catalyst's `JoinSelection` strategy walks a fixed priority order for an equi-join: 1. **A broadcast hint** — `broadcast(smallDf)` in the DataFrame API, or `/*+ BROADCAST(t) */` in SQL — plans a `BroadcastHashJoin` if the join type permits that build side. 2. **Size-based broadcast** — if the planner's estimate for one side is at most `spark.sql.autoBroadcastJoinThreshold` (default `10485760`, i.e. 10 MB), that side becomes the build side of a `BroadcastHashJoin`. Setting the threshold to `-1` disables automatic broadcasting entirely; explicit hints still work. 3. **Shuffled hash join** — both sides are shuffled by key, and one side's partitions are built into hash maps. `spark.sql.join.preferSortMergeJoin` is `true` by default, so this is chosen relatively rarely — typically when one side is much smaller than the other or the keys are not orderable. 4. **Sort-merge join** — the general case: shuffle both sides by hash of the join key, sort each partition by the key, then merge the two sorted streams. In the plan this appears as `SortMergeJoin` with an `Exchange hashpartitioning` and a `Sort` under each branch. 5. **No equality condition** — a range or inequality predicate cannot be hashed or merged, so Spark falls back to `BroadcastNestedLoopJoin` (broadcast one side, loop over pairs) or `CartesianProduct`. These are quadratic and are the plan nodes to fear. ## What "one side is small" actually means The estimate comes from statistics on the relation. For a file-based table with no computed statistics, that is essentially the **on-disk size of the files**. Parquet is compressed and columnar, so a 9 MB file can expand to hundreds of megabytes of deserialized JVM rows in the driver. That gap is the single most common cause of a driver OOM or a `spark.driver.maxResultSize` failure right after a plan showed a friendly `BroadcastHashJoin`. Better numbers come from `ANALYZE TABLE t COMPUTE STATISTICS`, and per-column statistics from `ANALYZE TABLE t COMPUTE STATISTICS FOR COLUMNS ...`, which the cost-based optimizer uses for cardinality estimation. CBO is off by default (`spark.sql.cbo.enabled` is `false`), and cost-based join reordering (`spark.sql.cbo.joinReorder.enabled`) is off as well; without them the planner leans on size estimates propagated through the plan, which degrade badly after filters and joins. The other half of the risk is *time*: broadcasting is a blocking step. The driver builds the relation and pushes it out, and if that exceeds `spark.sql.broadcastTimeout` (default 300 seconds) the query fails with a broadcast timeout, which usually means the "small" side was not small. ## Why broadcast is so much cheaper when it fits A `SortMergeJoin` shuffles **both** inputs: every row of both tables is serialized, written to local disk, fetched over the network and sorted. A `BroadcastHashJoin` shuffles **neither**: the small side is replicated to every executor, and the large side is joined in place, partition by partition, streaming. For a fact-to-dimension join — the most common shape in analytics — this is the difference between minutes and seconds, which is why the first question about a slow join is always whether the dimension could be broadcast. ## Join type constrains the build side You can only broadcast the side whose rows may be probed. For an inner join either side works. For a **left outer** join, every left row must appear in the output, so the left side must stream and only the **right** side can be broadcast; the mirror holds for a right outer join. A full outer join cannot be executed as a broadcast hash join at all. If you hint a side that the join type forbids, Spark ignores the hint and plans something else — which is why a hint that appears to do nothing is usually a join-type problem rather than a bug. ## Hints, and when to reach for one Available join hints include `BROADCAST` (aliases `BROADCASTJOIN`, `MAPJOIN`), `MERGE` (`SHUFFLE_MERGE`), `SHUFFLE_HASH` and `SHUFFLE_REPLICATE_NL`. Use `BROADCAST` when you know a side is small and the estimate does not — a heavily filtered dimension whose post-filter size the planner cannot see, or a table whose statistics are stale. Use `MERGE` to *stop* Spark broadcasting something that turns out to be large at runtime. Treat every hint as a pin that must be revisited: it overrides the planner's view forever, including after the data changes. ## Raising the threshold Raising `spark.sql.autoBroadcastJoinThreshold` is a legitimate tuning move for a cluster with fat drivers and executors — a few hundred megabytes is not unusual on modern hardware — but it is a whole-session setting that applies to every join. The blast radius is the driver's heap and every executor's memory, since each one holds its own copy of the broadcast relation. Prefer a targeted hint on the specific query over a global increase. ## And then AQE gets a vote With Adaptive Query Execution enabled, these are only the *initial* choices. After a shuffle completes, Spark knows the real size of each side and may replan the join with that knowledge. The planning-time decision described here is what happens with the statistics available before any data has been read; the runtime adjustment is a separate mechanism, and the plan printed by `explain()` shows only the first of the two.

  • A 9 MB Parquet dimension was broadcast and the driver ran out of memory. Why?
    The size estimate is the compressed, columnar on-disk footprint. Broadcasting materializes the table as deserialized JVM rows in a hash relation, which can be many times larger — dictionary-encoded strings expand, per-object overhead is added. Either force a sort-merge join with a MERGE hint, lower the threshold, or run ANALYZE TABLE so the planner works from row counts rather than file bytes.
  • Which side can Spark broadcast for a left outer join?
    Only the right side. Every row of the left input must appear in the output, so the left side has to stream through the join while the right side is replicated and probed. The mirror rule applies to a right outer join, and a full outer join cannot be a broadcast hash join at all — a BROADCAST hint on the forbidden side is simply ignored.
  • What happens when the join condition is a range rather than an equality?
    Neither hashing nor merging applies, so the planner falls back to BroadcastNestedLoopJoin if one side is small enough, and CartesianProduct otherwise. Both are quadratic in the input sizes. The usual fixes are to derive an equality key to join on first — a bucketed time key, for example — and apply the range predicate afterwards.
  • When would you set spark.sql.autoBroadcastJoinThreshold to -1?
    To disable automatic broadcasting for a session where estimates are untrustworthy and a surprise broadcast keeps killing the driver — for example over external sources with no statistics. Explicit BROADCAST hints still take effect, so you keep deliberate broadcasts while removing the guesswork. It is a blunt instrument; a targeted MERGE hint on the offending query is usually better.

Handing every clerk a copy of a thin price list is trivial; when both documents are thick, you instead sort each into alphabetical piles and walk them in step.

saying these in an interview costs you the question

  • Says Spark always picks the fastest join with no size input
  • Thinks the broadcast threshold compares in-memory size, not the estimate
  • Believes any side can be broadcast regardless of join type
  • Claims a sort-merge join shuffles only the larger side
  • Expects a broadcast hint to work for an inequality join condition

context

open as a page

In Spark, what does Catalyst do between your DataFrame code and the physical plan?

level: middleimportance: must knowfreq 70%

basics

~10 s

Catalyst runs four stages: analysis resolves names and types against the catalog, rule-based logical optimization rewrites the plan, physical planning picks operators and inserts shuffles, and Tungsten generates JVM code for fused operator chains.

open as a page

In Spark, does writing SQL instead of DataFrame code change how fast a query runs?

level: juniorimportance: should knowfreq 55%

basics

~20 s

No. A SQL string and the equivalent DataFrame code build the same logical plan, and Catalyst analyzes, optimizes and compiles both identically. Pick whichever reads better. The real exceptions are the RDD API and Python UDFs, which Catalyst cannot optimize through.

open as a page

In Spark, how do you read the physical plan that df.explain() prints?

level: middleimportance: should knowfreq 66%

basics

~20 s

Read it bottom-up: leaves are scans, indentation shows the tree. Each Exchange is a shuffle and a stage boundary, a *(n) prefix marks whole-stage code generation, and the scan line's PushedFilters and ReadSchema show what was pushed into the source.

open as a page

Why does a Python UDF in a Spark query defeat Catalyst optimization?

level: seniorimportance: should knowfreq 52%

basics

~20 s

Catalyst cannot see inside a Python UDF, so it cannot push it into the data source, fold it, or compile it into generated code. The plan gains a Python-evaluation node that breaks the fused loop and ships every row into a separate Python process and back.

open as a page

When is a custom Catalyst rule the right fix instead of rewriting the Spark query?

level: principalimportance: nice to knowfreq 15%

basics

~20 s

Only when the behaviour must apply to every query in the session and cannot live in user code — governance predicates, lineage capture, a new SQL grammar. For a single slow query, fix statistics, the predicate or the layout instead.

open as a page