Query performance optimization

Jones Oscar 40 Reputation points
2026-09-30T08:32:41.27+00:00

Hi Microsoft Community,

I am working with a data warehouse environment and have encountered a query performance issue that I am unable to resolve on my own.

A large analytical join query is running slower than expected and is causing significant bytes spilled to both local storage and remote storage during execution. The workload processes a substantial volume of data, and performance degradation becomes more noticeable as dataset sizes increase.

The specific challenges I am facing:

  1. Large join operations are spilling data to local and remote storage during query execution.
  2. Query performance degrades significantly when processing high-volume analytical workloads.
  3. Determining the most appropriate warehouse size for the workload has been challenging.
  4. Existing clustering keys may not be effectively supporting join and filtering patterns.

I have reviewed the query profile and identified spill activity, but I am unsure whether the primary issue is warehouse sizing, clustering strategy, or a combination of both. Is there a recommended approach for optimizing warehouse size and selecting effective clustering keys to reduce spill operations and improve query performance?

Any guidance from the community or Microsoft experts. Thanks in advance.

Windows for business | Windows 365 Enterprise
0 comments No comments

2 answers

Sort by: Most helpful
  1. HLBui 13,020 Reputation points Independent Advisor
    2026-09-30T09:38:57.9866667+00:00

    Hi Jones Oscar •

    I think this is most likely a combination of warehouse sizing and data layout rather than a clustering issue alone. When a query starts spilling to local and especially remote storage, it indicates that the join operation requires more memory than the warehouse can provide, which leads to a significant performance penalty.

    My first recommendation would be to temporarily scale the warehouse up one or two sizes and rerun the same query. If the spill volume drops substantially or disappears, you've confirmed that memory pressure is a major contributor. At the same time, review the join keys and filtering columns in your query profile. Clustering keys tend to provide the most benefit when they align with frequently used filter predicates and help reduce the amount of data scanned before the join occurs.

    It's also worth checking for data skew. If a small number of join key values account for a large percentage of rows, a larger warehouse alone may not completely solve the problem. In those cases, query rewrites, pre-aggregation, or breaking large joins into smaller stages can often produce better results than simply adding compute resources.

    As a general rule, addressing the spills first by validating warehouse sizing, then evaluating whether the clustering strategy is actually helping prune data for your most common workload patterns. The query profile should tell you which of these factors is having the biggest impact.

    If you find this response helpful, please click "Accept Answer" so it can help other community members facing similar performance issues as well.

    Was this answer helpful?

    0 comments No comments

  2. AI answer

    2026-09-30T08:33:00.9333333+00:00

    For spill-heavy analytical join workloads, the supported guidance points to addressing both memory capacity and data reduction patterns.

    1. Treat spill as a memory pressure signal. In Azure Databricks SQL, DATA_SPILL means 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.
    2. Check whether the warehouse is also under concurrency pressure. If queries are waiting in queue or showing high waiting_at_capacity_duration_ms, increase max_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.
    3. Use query history to separate sizing problems from query-shape problems. For SQL warehouses, queries waiting at capacity indicate max_clusters is 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
    4. 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.
    5. 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 shows SELECTIVE_JOIN, add filters before the join to reduce input rows. Both patterns help lower memory demand and spill.
    6. 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.
    7. 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);.
    8. For Lakehouse-backed Fabric queries, maintain file layout. If querying Lakehouse data through the SQL analytics endpoint, regular table maintenance and using OPTIMIZE to combine many small files can reduce scan overhead.

    A practical optimization order is:

    1. Confirm whether the main symptom is spill, queueing, or both.
    2. Increase warehouse size when spill is the dominant issue.
    3. Increase max_clusters when queue time is also high.
    4. Align clustering keys with real filter predicates.
    5. Push filters before joins and reduce projected columns.
    6. Review whether the join is exploding row counts.
    7. In Fabric, check for single-node operators and use FORCE DISTRIBUTED PLAN only when query semantics allow it.

    AI-generated content may be incorrect. Read our transparency notes for more information.

    Was this answer helpful?

Your answer

Answers can be marked as 'Accepted' by the question author and 'Recommended' by moderators, which helps users know the answer solved the author's problem.