Skip to content
Reliable Data Engineering
Lesson
Open in the interactive app

SQL Performance, Execution Plans and Dialects

Senior candidates are expected to say why a query is slow and how the engine executes it, not just produce correct results.


1. How a distributed SQL engine runs your query

flowchart LR
    SQL[SQL text] --> P[Parse → logical plan]
    P --> O[Optimiser<br/>predicate pushdown, column pruning,<br/>join reordering, constant folding]
    O --> PH[Physical plan<br/>join strategies, exchanges]
    PH --> ST[Stages split at shuffles]
    ST --> T[Tasks per partition<br/>on executors]

The expensive parts, in order: reading data (I/O), shuffling (network + disk), sorting, then CPU.

2. Join strategies

StrategyHowWhen chosenCost
Broadcast hash joinSmall table copied to every executor; probe locallyOne side < broadcast threshold (Spark default 10 MB; AQE can switch at runtime)No shuffle of the big side: fastest
Shuffle hash joinBoth sides shuffled by key; hash table on one side per partitionMedium sides, equi-joinShuffle both
Sort-merge joinBoth sides shuffled + sorted by key; mergedLarge-large equi-joins (Spark default)Shuffle + sort both
Nested loop / cartesianEvery row × every rowNon-equi joins without better optionO(n·m): danger
-- Force a broadcast in Spark SQL when you know the dimension is small
SELECT /*+ BROADCAST(d) */ f.*, d.name FROM fact f JOIN dim d ON f.dim_id = d.id;

Non-equi joins (ranges, BETWEEN, interval overlaps) often become nested loops. Mitigate with bucketing tricks (join on day/bucket first, then filter), range join hints (Databricks RANGE_JOIN), or rewriting as window functions.

3. Reading a plan: what to look for

== Physical Plan ==
*(5) HashAggregate(keys=[country], functions=[sum(amount)])
+- Exchange hashpartitioning(country, 200)           ← shuffle
   +- *(4) HashAggregate(keys=[country], functions=[partial_sum(amount)])   ← map-side partial agg
      +- *(4) Project [country, amount]
         +- *(4) BroadcastHashJoin [customer_id], [id], Inner, BuildRight   ← broadcast join, good
            :- *(4) Filter isnotnull(customer_id)
            :  +- *(4) ColumnarToRow
            :     +- FileScan parquet [customer_id, amount] PushedFilters: [IsNotNull(customer_id)],
            :        PartitionFilters: [isnotnull(dt), (dt = 2026-10-01)]    ← partition pruning worked
            +- BroadcastExchange
               +- FileScan parquet [id, country]

Checklist:

4. Query anti-patterns

Anti-patternWhy it’s slowFix
SELECT * on wide columnar tablesReads every columnSelect needed columns
Function on filter column: WHERE DATE(ts) = '2026-10-01'Can block pruning/pushdownWHERE ts >= '2026-10-01' AND ts < '2026-10-02' (or filter on the partition column)
UNION instead of UNION ALLExtra dedup (sort/shuffle)UNION ALL when duplicates impossible or wanted
COUNT(DISTINCT) on huge cardinalityBig shuffleapprox_count_distinct, or pre-aggregate
Join then aggregateJoins explode row countsAggregate first, then join (when semantics allow)
OR across different columns in join conditionsPrevents hash joinSplit into UNION ALL of two joins
Correlated subqueries per rowN executions (engine-dependent)Rewrite as join / window
ORDER BY without LIMIT on huge resultsGlobal sortSort only when needed
NOT IN with NULLable subqueryWrong results and slowNOT EXISTS
Window with no PARTITION BYSingle taskPartition, or accept for small data

5. Data skew in SQL engines

Symptom: one task runs 30 min while others take 30 s. Causes: join or group key with a dominant value (NULLs, “unknown”, one mega-customer).

Fixes:

6. Dialect cheat sheet

TaskPostgresSpark SQL / DatabricksSnowflakeBigQuerySQLite (used in this repo’s runner)
Filter on window resultsubqueryQUALIFYQUALIFYQUALIFYsubquery
Date addd + INTERVAL '7 day'date_add(d, 7)DATEADD(day, 7, d)DATE_ADD(d, INTERVAL 7 DAY)DATE(d, '+7 days')
Date diff (days)d2 - d1datediff(d2, d1)DATEDIFF(day, d1, d2)DATE_DIFF(d2, d1, DAY)julianday(d2) - julianday(d1)
Truncate to monthdate_trunc('month', d)date_trunc('MONTH', d)DATE_TRUNC('month', d)DATE_TRUNC(d, MONTH)strftime('%Y-%m-01', d)
Pick column of max rowDISTINCT ONmax_by(col, ts)MAX_BY(col, ts)ARRAY_AGG(col ORDER BY ts DESC LIMIT 1)[OFFSET(0)]window + filter
String aggstring_agg(x, ',')array_join(collect_list(x), ',')LISTAGG(x, ',')STRING_AGG(x, ',')group_concat(x, ',')
Conditional aggFILTER (WHERE …)count_if, FILTERCOUNT_IFCOUNTIFSUM(CASE …) / FILTER
Medianpercentile_cont(0.5) WITHIN GROUPpercentile_approx / medianMEDIANAPPROX_QUANTILESmanual
Explode arrayunnestexplodeLATERAL FLATTENUNNESTjson_each
UpsertINSERT … ON CONFLICTMERGE INTOMERGEMERGEINSERT … ON CONFLICT
Integer division 5/222.52.52.5 (/ is float)2

7. Interview questions

A query joining a 2 TB fact table with a 50 MB dimension is slow. What do you check?

Whether the dimension is broadcast (plan shows SortMergeJoin with two Exchanges instead of BroadcastHashJoin). Raise the broadcast threshold or add a hint; check stats are up to date. Then check partition pruning on the fact table, skew on the join key (NULL keys), and whether only needed columns are read.

Why can GROUP BY before JOIN be faster, and when is it incorrect?

Pre-aggregating the large side shrinks the data before the shuffle/join. It’s incorrect if the join filters or multiplies rows in a way that changes the aggregate (e.g. the join removes some rows, or a one-to-many join should duplicate values), or if you need columns lost in aggregation.

What's the difference between WHERE and HAVING, and between ON and WHERE in a LEFT JOIN?

WHERE filters rows before grouping, HAVING filters groups after aggregation. In a LEFT JOIN, a condition on the right table in ON keeps unmatched left rows (with NULLs); the same condition in WHERE removes them, effectively turning the LEFT JOIN into an INNER JOIN.