In Snowflake's Query Profile, which statistics tell you why a query is slow, and what fix does each imply?
answer
- look at the most expensive node before anything else
- one ratio tells you the filter did nothing
- spilling has two tiers and the second one is bad
- rows out of a join can exceed rows in
- a percentage tells you whether compute was warm
basics
~20 sRead partitions scanned versus total for pruning, bytes spilled to local and remote storage for memory pressure, percentage scanned from cache for I/O locality, and the operator tree's row counts for exploding joins. Each points at a different fix: filter shape, query shape, cache warmth, join keys.
solid answer
~50 sStart with **Most Expensive Nodes** and the execution-time breakdown (processing, local disk I/O, remote disk I/O, network, synchronization, initialization) to find where the time actually went. Then read the statistics. **Partitions scanned vs partitions total** is pruning: scanning nearly everything means the filter is not eliminating micro-partitions, so fix the predicate shape or the table's data layout. **Bytes spilled to local storage** means the operation exceeded memory and went to the cluster's SSD; **bytes spilled to remote storage** means it exceeded local disk too and is now an order of magnitude worse — reduce the data entering the sort/join/aggregate before reaching for more compute. **Percentage scanned from cache** tells you whether the warehouse was warm. In the operator tree, a join whose output rows vastly exceed its inputs is an **exploding join** — usually a missing predicate or duplicated keys on one side. A `UNION` where `UNION ALL` would do adds a whole deduplicating aggregate for nothing.
code
text · 13 linesMost Expensive Nodes
Join [3] 61.4%
TableScan [7] 28.1%
IO
Bytes scanned 1.42 TB
Percentage scanned from cache 2%
Bytes spilled to local storage 210 GB
Bytes spilled to remote storage 38 GB
Pruning
Partitions scanned 184,320
Partitions total 184,320go deeper
Know that Snowflake exposes a per-query profile with an operator tree and statistics, and that the first thing to look at is which operator consumed most of the time.
Explain what each headline statistic means — pruning ratio, local versus remote spilling, percentage scanned from cache — and the mechanism behind each number.
Walk a profile end to end and name the specific fix each symptom implies, resisting the reflex to size up. Be ready to spot an exploding join or a needless UNION deduplication from row counts alone.
Own the standard: what a team must show in a profile before it is allowed to raise a warehouse's size, and how you turn recurring profile symptoms into layout or modelling changes rather than per-query patches.
## What Query Profile shows you Query Profile renders the executed plan as a tree of operators — `TableScan`, `Filter`, `Join`, `Aggregate`, `Sort`, `WindowFunction`, `Result` — with per-operator time percentages, plus a statistics panel and a **Most Expensive Nodes** list. The discipline is: find the dominant node first, then read the statistic that explains it. "Use a bigger warehouse" is the wrong first answer precisely because the profile usually names a cheaper fix. ## Pruning: partitions scanned vs partitions total Snowflake stores tables as micro-partitions and keeps min/max metadata per column per partition. If `partitions scanned` is close to `partitions total`, the query read the whole table and the filter bought nothing. Common causes are predicate shapes that defeat the metadata comparison: wrapping the filtered column in a function (`WHERE TO_DATE(event_ts) = ...`), an implicit or explicit cast on the column side, a leading-wildcard `LIKE`, or filtering on a column whose values are scattered across every partition because rows arrived in unrelated order. The fixes live in the predicate (compare the raw column against a computed constant instead of computing on the column) or in how the data is laid out. ## Spilling: bytes spilled to local and remote storage Sorts, hash joins, and aggregations build state in memory. When that state exceeds the warehouse's memory, Snowflake spills to the cluster's local SSD; when it exceeds local disk, it spills further to remote storage. Remote spilling is dramatically slower than local spilling and is a reliable sign a query has fallen off a cliff rather than just being big. The instinct is to size up, and more memory does help. But the profile usually shows a cheaper fix first: project fewer columns so less width flows into the sort, filter earlier so fewer rows reach the join, aggregate before joining rather than after, replace a global `ORDER BY` over the full result with a top-N, and check for the duplicate-key explosion that quietly multiplied the join's build side. ## Cache warmth: percentage scanned from cache This IO statistic is the share of scanned data served from the warehouse's local SSD cache rather than remote storage. A cold number on a query you expected to be warm tells you the warehouse suspended between runs, was resized, or that the workload is spread over so many warehouses that none stays warm. It explains a latency difference between two runs of the same query without any plan change. ## Exploding joins In the operator tree, compare a join's input row counts with its output. Output far larger than either input means the join key is not as unique as assumed — duplicates on the supposedly-one side, or a join condition missing a component so the engine produced a near-cross product. This is one of the classic problems Query Profile is documented to help identify, alongside `UNION` without `ALL`, queries too large to fit in memory, and inefficient pruning. The fix is a corrected join key or a deduplicated build side, never a bigger warehouse. ## UNION versus UNION ALL `UNION` removes duplicates, which means an extra full aggregation over the combined result. When the inputs are known disjoint — the usual case when someone is stitching together date ranges or source systems — `UNION ALL` deletes an entire expensive operator from the plan. ## The time breakdown The execution-time categories tell you which resource dominated: high **processing** points at CPU-bound expression work or a genuinely large aggregation; high **remote disk I/O** points at cold scans or remote spilling; high **network communication** points at large redistributions between nodes; high **synchronization** points at skew or waiting; high **initialization** on a fast query points at compilation overhead, which matters for very short statements executed at high rates. ## Working the profile in order A reliable routine: (1) Most Expensive Nodes — where did time go; (2) pruning ratio on the dominant `TableScan`; (3) spilling numbers, local then remote; (4) row counts across joins for explosions; (5) cache percentage to explain run-to-run variance. Only after all five come back clean is "this query needs more compute" the honest conclusion — and at that point the right lever belongs to warehouse configuration rather than to the query.
- Bytes spilled to remote storage is large. Why not just size up the warehouse and move on?More memory does remove the spill, and sometimes that is the right call. But it raises the credit rate for every query on that warehouse, and it hides the cause. Check first whether fewer columns, an earlier filter, aggregating before joining, or a top-N instead of a full sort shrinks the state enough. Size up when the data genuinely is that large.
- Partitions scanned equals partitions total, but the filter looks selective. What are the usual causes?Either the predicate defeats metadata comparison — a function or cast applied to the column, a leading-wildcard LIKE, a join-derived value not known at prune time — or the filtered column's values are spread across every micro-partition because the rows were never ordered by it. The first is a rewrite; the second is a data-layout problem.
- A query shows high initialization time but tiny execution time. What does that suggest?Compilation and setup dominate, which happens with very short statements, extremely wide or complex SQL, or high-frequency single-row work. The fix is batching many tiny statements into fewer larger ones, or simplifying generated SQL, rather than anything to do with scan or join performance.
saying these in an interview costs you the question
- Reaches for a bigger warehouse before reading the profile
- Treats local and remote spilling as equally serious
- Ignores the partitions scanned versus total ratio
- Cannot explain what an exploding join looks like in the tree
- Assumes high percentage scanned from cache means the query is optimal