ORDER BY — filtering becomes range-scan

Rows inside a part are physically sorted by ORDER BY, so a WHERE on a leading-prefix column becomes a contiguous range scan; a WHERE on a non-prefix column is a full part scan.

Previously

Merge keeps the part count bounded, but what a query does inside each part — scan everything, or jump to a range — depends on how the rows were sorted when the part was written.

Scene 08

ORDER BY: filtering becomes range-scan

  1. Watch
  2. Try it
  3. Predict
  4. Capture
WHEREWHERE service = 'api'SORT KEYORDER BY (service, ts)ONE PART · 0 rows in physical on-disk orderLEADING-COLUMN CARDINALITYservice: 3 distinct values (low)filling part…Sort key decides what is contiguous on disk; WHERE on a leading prefix is a range scan.
What to watch for

Watch the part fill up with 60 row tiles in physical on-disk order — all the api rows first, then web, then cron, because the part was written with ORDER BY (service, ts). Once the strip lands, WHERE service='api' lights up the first 20 tiles as a contiguous range — the engine reads that slice and nothing else.

Continue unlocks when the animation finishes.
Implementation

Highlighted lines are the ones running in the diagram right now.

Planner.plan_query(part, where)
walk ORDER BY left-to-right, match WHERE, pick range vs scan
def plan_query(part, where):
sort_cols = part.order_by # e.g. ('service', 'ts')
# Find the longest prefix of sort_cols that WHERE
# constrains with an equality / range predicate.
prefix = []
for col in sort_cols:
if where.constrains(col):
prefix.append(col)
else:
break # gap kills everything to the right
if not prefix:
return FullPartScan() # non-prefix WHERE
return RangeScanBySortKey(prefix, where)
RangeScanBySortKey.execute(part, where)
on (service, ts): seek to service='api', read until it flips
def range_scan_by_sort_key(part, where):
# Binary-search the in-memory primary.idx to slice
# the granule range where the prefix could live.
lo = bisect_left(part.sort_columns, where.prefix_lo)
hi = bisect_right(part.sort_columns, where.prefix_hi)
# Seek to first row where service = 'api',
# then sequential read until service != 'api'.
for row in part.rows[lo:hi]:
if where.matches(row):
yield row
# Leftmost-prefix rule
the schema-design law the planner is enforcing
# ORDER BY (A, B):
# WHERE A = ... -> range scan (fast)
# WHERE A = .. AND B = .. -> range scan (fastest)
# WHERE B = ... -> full scan (B's range
# restarts inside every A)
#
# ORDER BY (user_id, ts) with user_id unique per row:
# primary.idx has one entry per granule, each a
# different user_id -> no granule can be skipped,
# every WHERE degenerates to a full part scan.
# This is the canonical ClickHouse anti-pattern.

Where this sits in Build a columnar OLAP store (ClickHouse / Druid style)

Scene 08 of 13, in the Read side act — Sort key, sparse primary index, skip indexes — narrowing what to scan.. Rows inside a part are sorted by ORDER BY; WHERE on a leftmost-prefix column is a contiguous range, anything else is a full scan.

Up next. If rows inside a part are sorted, the engine doesn't need an index entry per row — it can keep one entry per group of 8192 rows and still find a range fast. That gives us an index small enough to keep entirely in memory.

All 13 scenes in Build a columnar OLAP store (ClickHouse / Druid style) · Every curriculum

Built with Arqly
Every scene in Build a columnar OLAP store (ClickHouse / Druid style) builds on the one before it.All 13 Build a columnar OLAP store (ClickHouse / Druid style) scenes