For spill-heavy analytical join workloads, the supported guidance points to addressing both memory capacity and data reduction patterns.
- Treat spill as a memory pressure signal.
In Azure Databricks SQL,
DATA_SPILLmeans data did not fit in memory during query execution. The recommended action is to increase the warehouse size to add memory. Also reduce memory usage by limiting rows, columns, or large column types such as strings, arrays, maps, and structs. - Check whether the warehouse is also under concurrency pressure.
If queries are waiting in queue or showing high
waiting_at_capacity_duration_ms, increasemax_clusters. This addresses queue time, not spill directly, but it helps when performance degradation is caused by workload concurrency in addition to query memory pressure. - Use query history to separate sizing problems from query-shape problems.
For SQL warehouses, queries waiting at capacity indicate
max_clustersis too low, while excessive disk spill indicates the warehouse size is too small or the query needs optimization. That gives a practical decision path:- High queue time -> increase
max_clusters - High spill -> increase warehouse size
- High spill even after resizing -> optimize the query and data layout
- High queue time -> increase
- Revisit clustering and filtering alignment.
If performance insights show
COVERAGE_FILTER_KEYS_CLUSTERING, the table is clustered by keys that are not used in filters during the scan. The recommendation is to add filters on the clustering keys to reduce bytes read. If current clustering keys do not match actual filter patterns, they are not helping scan reduction. - Reduce join input before the join.
If the query shows
EXPLODING_JOIN, the join is producing far more rows than it reads. The recommendation is to tighten the join condition or reduce input rows from both sides. If it showsSELECTIVE_JOIN, add filters before the join to reduce input rows. Both patterns help lower memory demand and spill. - Remove unnecessary data movement in the query shape.
Additional supported optimizations include:
- Remove redundant joins when they do not affect the result.
- Remove redundant aggregations when they do not change the result.
- Avoid wide projections by selecting only needed columns, especially when large column types are present.
- If the workload is in Fabric Data Warehouse, watch for single-node plan limitations.
Queries can run slowly or fail when parts of the plan must execute on a single node, such as
TOP, global sorting, or final result ordering. In that case:- Reduce the filtered dataset size.
- If semantics allow it, try
OPTION (FORCE DISTRIBUTED PLAN);.
- For Lakehouse-backed Fabric queries, maintain file layout.
If querying Lakehouse data through the SQL analytics endpoint, regular table maintenance and using
OPTIMIZEto combine many small files can reduce scan overhead.
A practical optimization order is:
- Confirm whether the main symptom is spill, queueing, or both.
- Increase warehouse size when spill is the dominant issue.
- Increase
max_clusterswhen queue time is also high. - Align clustering keys with real filter predicates.
- Push filters before joins and reduce projected columns.
- Review whether the join is exploding row counts.
- In Fabric, check for single-node operators and use
FORCE DISTRIBUTED PLANonly when query semantics allow it.