skip to content

In Delta Lake, what does ZORDER BY add to an OPTIMIZE command?

level: middleimportance: should knowfreq 55%

answer

  1. bin-packing changes file size, this changes file contents
  2. tighter min/max means more files pruned
  3. one curve serving several filter columns at once
  4. only helps columns queries actually filter on
  5. statistics stop at the 32nd column by default

basics

~20 s

ZORDER 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 s

Plain `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
sql
-- 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

for a junior

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.

for a middle

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.

for a senior

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.

for a principal

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

context