In Delta Lake, what does ZORDER BY add to an OPTIMIZE command?
answer
- bin-packing changes file size, this changes file contents
- tighter min/max means more files pruned
- one curve serving several filter columns at once
- only helps columns queries actually filter on
- statistics stop at the 32nd column by default
basics
~20 sZORDER BY makes OPTIMIZE co-locate rows with similar values of the named columns into the same files, so per-file min/max statistics get tighter and the planner skips more files for filters on those columns. Plain OPTIMIZE only bin-packs.
solid answer
~50 sPlain `OPTIMIZE` merges small files; `OPTIMIZE tbl ZORDER BY (user_id, event_type)` additionally **sorts rows across files** using a space-filling curve that interleaves the values of the listed columns, so rows with similar values land together. Each Delta data file carries min/max statistics in its `add` action, and tighter ranges mean the planner can prune whole files for a `WHERE user_id = ...` filter. It pays off for **high-cardinality columns that queries actually filter on**; low-cardinality columns are better served by partitioning. Effectiveness dilutes fast as you add columns — one to three is the practical range. Two catches: statistics are only collected for the first 32 columns by default (`delta.dataSkippingNumIndexedCols`), so Z-ordering an unindexed column buys nothing; and unlike bin-packing, Z-ordering re-clusters an entire partition whenever new data lands there, so its cost scales with partition size, not with the increment.
code
sql · 8 lines-- Cluster one day's files by the two columns dashboards filter on
OPTIMIZE analytics.events
WHERE event_date = '2026-08-20'
ZORDER BY (user_id, device_id);
-- Z-ordering a column past the 32nd needs statistics first
ALTER TABLE analytics.events
SET TBLPROPERTIES (delta.dataSkippingNumIndexedCols = 40);go deeper
Know that Delta can skip whole files using per-file min/max statistics, and that ZORDER BY is the option on OPTIMIZE that arranges rows so those statistics become useful.
Explain the mechanics: interleaved multi-column ordering, tighter per-file ranges, why one to three high-cardinality filter columns is the sweet spot, and why partitioning suits low-cardinality columns better.
Diagnose a Z-order that did not help — unindexed columns beyond the statistics limit, filters the layout does not serve, functions applied to the predicate — and budget the recurring full-partition rewrite cost.
Decide layout policy for tables several teams query with different predicates: whether to serve one access path, adopt liquid clustering for its incremental behaviour, or accept scan cost rather than pay perpetual rewrite amplification.
## The problem: skipping, not scanning Every data file in a Delta table has an `add` action in `_delta_log` that carries per-file statistics: row count, and min, max and null count for each indexed column. When a query says `WHERE user_id = 8812`, the planner reads those statistics and drops any file whose `user_id` range cannot contain 8812. That is **data skipping**, and it is the main reason a Delta query on a big table can finish quickly. Skipping only works if the ranges are narrow. If rows arrive in event-time order, each file holds a near-random spread of `user_id` values, so nearly every file's min/max range spans nearly the whole domain and no file can be excluded. The physical layout, not the statistics mechanism, is the limiting factor. ## What ZORDER BY does about it ```sql OPTIMIZE events WHERE event_date = '2026-08-20' ZORDER BY (user_id, device_id); ``` Z-ordering interleaves the bits of the values of the listed columns to produce a single ordering key — a space-filling curve — and writes rows out in that order. The effect is that rows close in *any* of the Z-ordered dimensions tend to end up in the same file. A one-dimensional sort would make one column's ranges perfect and leave the others as scattered as before; Z-order accepts slightly worse skipping on each column in exchange for useful skipping on all of them. After the rewrite, the new files' statistics are much tighter on the Z-ordered columns, and filters on those columns prune far more files. ## Choosing the columns Three rules decide whether Z-ordering earns its cost: 1. **Pick columns that appear in query filters.** Clustering on a column nobody filters on rewrites the table for nothing. 2. **Pick high-cardinality columns.** A boolean or a five-value status column has nothing to gain from clustering; if you need coarse separation on a low-cardinality column, partition on it instead. Z-order and partitioning compose — you typically partition by date and Z-order within the partition. 3. **Keep the list short.** Each additional column dilutes the curve: the ordering must satisfy more dimensions at once, so per-file ranges widen for all of them. One to three columns is the practical range; five is usually worse than two. ## The statistics trap Delta collects min/max statistics only for the first `delta.dataSkippingNumIndexedCols` columns of the schema — 32 by default. Z-ordering by a column past that boundary produces a beautifully clustered layout with **no statistics to prove it**, so the planner still cannot skip anything. The fixes are to raise the property, or to move the column earlier in the schema (a Delta column-order change), then rewrite. This is a very common "we Z-ordered and nothing improved" cause, and a good thing to mention in an interview. Long string columns are also a poor fit: their statistics are truncated, and high-entropy identifiers dominate the curve. ## Cost, and why re-running is expensive Bin-packing only touches files below the target size, so re-running `OPTIMIZE` is nearly free. Z-ordering is different: to place a new file's rows correctly relative to existing rows, it must re-cluster the whole partition. If nothing changed in a partition since the last run, re-running is a no-op — but as soon as a partition receives new data, the next Z-order run rewrites that partition end to end. Cost therefore scales with **partition size**, not with the size of the increment, which is exactly the wrong shape for a table receiving continuous small writes. That is the operational reason to scope the command with a `WHERE` on partition columns, run it on a cadence that matches how often queries actually benefit, and account for the write amplification (and, for hourly-partitioned tables, the sheer number of partitions). ## Relationship to liquid clustering Newer Delta versions offer liquid clustering, defined with `CLUSTER BY` on the table rather than as an argument to a single command. It clusters incrementally, lets you change the clustering keys without rewriting existing data, and is applied by running plain `OPTIMIZE`. It is mutually exclusive with both `ZORDER BY` and partitioning on the same table. For a table under constant maintenance, that incremental property is precisely what Z-ordering lacks, so on a runtime that supports it, liquid clustering is now the default recommendation and `ZORDER BY` is the older mechanism you will still meet in existing pipelines. ## How to verify Compare files scanned before and after for a representative filter — the query plan reports files read versus files pruned. Do not judge Z-ordering by wall-clock time on a warm cache; judge it by the number of files the planner eliminated.
- Why is Z-ordering on a low-cardinality column like country usually the wrong choice?With a handful of distinct values, files already tend to contain narrow ranges or the filter matches a large share of the data anyway, so clustering buys almost no extra pruning while costing a full rewrite. A low-cardinality column with well-balanced values is a partitioning candidate instead, which physically separates the values into directories the planner eliminates before it even reads statistics.
- You Z-ordered by session_id and query times did not improve — what do you check first?Whether the column has statistics: only the first 32 schema columns are indexed by default, so a column further right has no min/max and cannot be skipped on. Then check that the queries actually filter on it with an equality or range predicate rather than a function or cast, which defeats skipping, and confirm from the query plan how many files were pruned rather than trusting wall-clock time.
- Can you combine ZORDER BY with partitioning?Yes, and it is the classic layout: partition by a coarse, low-cardinality column such as date, then Z-order the high-cardinality filter columns inside each partition. What you cannot combine it with is liquid clustering — a table defined with CLUSTER BY replaces both partitioning and Z-ordering.
Sorting a library by author gets you to one author instantly but scatters every subject. Z-ordering is shelving so that books close in author and subject sit near each other — slightly worse for each single lookup, much better when you search by either.
saying these in an interview costs you the question
- Thinks ZORDER BY is just sorting by the first column
- Adds five or six columns expecting better skipping from each
- Z-orders low-cardinality columns instead of partitioning them
- Assumes re-running Z-order only touches newly added files
- Ignores that unindexed columns have no statistics to skip on