Reference · incremental view maintenance

What each SQL operator costs per change.

An incremental engine keeps a query’s result current by processing only what changed. Each operator then has two costs: the work one change takes, and the state it must keep to do that work. Both follow from the math, whatever engine runs it. This page lists them, the way Redis lists the complexity of each command.

Why incremental · The model · The variables · Two complexities · The reference · Skew and cold state · Rules of thumb

Change in

3 orders
+1 +1 −1

WHERE status = 'paid'

0read

0held

JOIN customers ON id

3read

2.4Mheld

SUM(amount) BY plan

3read

12held

Change out

3 plans
−old +new

Three changed orders through a model of three operators. The work is six entry reads; the state is 2.4 million rows, of which this tick touched three.

01 · Why incremental

A refresh costs what it reads.

A database answers a query by reading the data it names, so recomputing a result costs the same whether one row changed or a million. An incremental engine reads only what the change touches, and a thousand changed rows out of a billion is a thousand reads, not a billion.

111k1k1M1M1B1B rows changed per refresh → ↑ entries read a million times fewer recompute: 1B, always incremental: as many as changed

recompute from scratch incremental, a linear aggregate

A table of a billion rows, refreshed after Δ of them changed. Every other operator sits between the two lines: a join reads Σ m, the other side's rows under the changed keys, and reaches the top line only when one key holds everything.

02 · The model

Changes are rows with weights.

Incremental engines in the DBSP and differential dataflow family treat a table as a multiset: each row carries a weight, +1 for an insert and −1 for a delete. An operator turns the change to its input into the change to its output, and keeps whatever state it needs to do that.

Change in

order 812 · pro · $40+1

order 377 · pro · $25−1

order 377 · pro · $30+1

SUM(amount) GROUP BY plan

state: pro → total $9,100 · 228 rows

reads 1 entry, writes 1

Change out

pro · $9,100−1

pro · $9,145+1

An update is a delete and an insert. The aggregate reads one total, not the 228 rows behind it, and its output is itself a change, which the next operator takes in the same way.

03 · The variables

Complexity in more letters than N.

A database’s cost is a function of its whole data. An incremental operator’s cost follows the change and the keys it lands on, so it takes a few variables to say, and most of them describe the shape of the data rather than its size.

SymbolIsSet by
Δthe rows changed in this tick: inserts, deletes, and both halves of an updatethe source's write rate and how often the engine ticks
Kthe distinct keys (join keys, groups, partitions) those changes land onhow the changes spread over the keys; at most Δ
Nthe rows the operator's input holds in allthe data
mthe rows held under one touched key; Σ m adds it up over the K keysthe key's cardinality, and skew: N / G on average, up to N for one hot key
Gthe keys (groups, partitions, distinct rows) in allthe query's key
L, Ra join's left and right inputsthe two tables
W, Ea time window's rows, and the rows entering or leaving it in one clock stepthe window's length and the clock's step
ka top-k filter's limitthe query

Reads and writes count state entries, which is what a tick pays for: an entry in a cached block costs a memory read, one in a cold block a block fetch. A lookup in an ordered store is O(log N) and is counted as O(1) here, as Redis counts a hash lookup.

04 · Two complexities

Work per change and state held are different numbers.

Move the sliders. Operators whose work follows Δ or K stay on the left of the chart as the data grows; the ones whose work follows Σ m climb toward the line with the hottest key. The distance below the line is state a tick does not touch, which is what can live on disk.

100M
1k
100k
even spread

Keys touched per tick (K): 995 · rows under them (Σ m): 995k

111k1k1M1M1B1B state entries kept → ↑ entries read per tick reads all it holds below the line: state a tick does not touch 123·4·56·8·9710·11·12·1314
DotOperatorRead per tickState keptShare touched
1Filter, project00none held
2Time-window filter010Munder 0.01%
3Distinct1k100k1.0%
4First row per key995100k1.0%
5SUM, COUNT995100k1.0%
6MIN, MAX2k100Munder 0.01%
7Top-10 per key28k100M0.03%
8Running total, appends2k100Munder 0.01%
9Join to a unique key1k100Munder 0.01%
10Join, many rows per key995k200M0.50%
11COUNT(DISTINCT)995k100M1.0%
12Window, whole partition995k100M1.0%
13Recursion, incremental2M200M1.0%
14Recursion, recomputed100M100M100%
Every dot is an operator shape, numbered as in the table; dots on one spot share a number. The arithmetic is the complexity column of the reference with constants of one: real engines differ by constant factors, not by shape.

05 · The reference

Operator by operator.

Each entry gives the work per tick and the state, how the operator turns a change into a change, what the Nth row costs, and how skew and cold state play out. The drawings show eight keys’ state, with a change landing on two of them.

OperatorWork per tickState
Filter, project, unionO(Δ)none
Time-window filterO(Δ) per change · O(E) per clock stepO(W)
DistinctO(Δ)O(G)
First row per key, append-only inputO(K)O(G)
SUM, COUNT, AVGO(K)O(G)
MIN, MAXO(K) in the common case · O(m) when a group loses its extremesup to O(N)
COUNT(DISTINCT), STRING_AGGO(Σ m)O(N)
JoinO(Σ m)O(L + R)
Top-k per partitionO(K · k)O(N)
Window over a whole partitionO(Σ m)O(N)
Running total and bounded framesO(Δ) for appends · O(m) for a change at the startO(N)
Recursive query, incrementalO(Σ m) at each level the change reachesO(N), plus each level's operators' state
Recursive query, recomputedO(N)O(N)

read this tick written held, not touched

Filter, project, union

WHERE amount > 0 · SELECT a, b * 2 · UNION ALL

Work per tick
O(Δ)
State
none

No state

How
Each output row depends on one input row alone, so a change maps straight to a change: an insert passes as an insert, a delete as a delete, each with its weight. Nothing is looked up and nothing is kept.
The Nth row
The same as the first row: one evaluation of the expression, with no N in it.
Skew and cold state
Nothing to be skewed about. These operators hold no data, and an engine fuses a run of them into one pass over the change.

Time-window filter

WHERE ts > now() - interval '7 days'

Work per tick
O(Δ) per change · O(E) per clock step
State
O(W)

reads 3 of 16 entries, writes 1

How
Each row has a window of clock values in which it passes, known from the row alone. Rows are kept sorted by the time they leave, so when the clock moves, the rows due to expire are at the front of that order and nothing else is read. Each one leaves as a retraction, which every operator above it then applies.
The Nth row
One write and no read, whatever W is. A clock step costs the rows entering or leaving (E), never the ones that stay.
Skew and cold state
The reads are at the oldest end and the writes at the newest, so the hot set is two ends of one sorted run; the middle of the window is cold however large it is. This is the one operator whose state shrinks on its own, and it shrinks everything above it.

Distinct

SELECT DISTINCT country, plan

Work per tick
O(Δ)
State
O(G)

reads 2 of 8 entries, writes in place

How
A count per distinct row. A row is emitted when its count rises from zero and retracted when it falls back to zero; in between, a change to its count emits nothing.
The Nth row
One point read and one write, however many rows share its value.
Skew and cold state
A hot value is one entry read often, which stays in memory. Values seen once and never again are the cold majority, and a value seen for the first time costs a membership check (a bloom filter), not a block read.

First row per key, append-only input

row_number() OVER (PARTITION BY id ORDER BY ts) = 1

Work per tick
O(K)
State
O(G)

reads 2 of 8 entries, writes in place

How
When the input only ever grows, a row behind the current first can never lead, so it is not kept at all. Each key keeps one row, replaced only when an earlier one arrives late, and the replaced row is retracted from the output.
The Nth row
One read of its key and at most two writes, whatever N is. The same query over an input that can delete rows is a top-k window with k = 1, which keeps every row.
Skew and cold state
State grows with distinct keys, not with the stream. Removing duplicate message ids reads one small entry per message and keeps one per id; the ids seen once are the cold majority.

SUM, COUNT, AVG

SELECT plan, count(*), sum(amount) GROUP BY plan

Work per tick
O(K)
State
O(G)

reads 2 of 8 entries, writes in place

How
These aggregates are linear: a row's contribution is added with its weight, so a delete subtracts it and an update is a subtract and an add. Each group keeps one small vector of totals. A change reads the group's totals, adds the delta's, and emits the group's old row retracted and its new row inserted.
The Nth row
One read and one write of its group's totals, however many rows the group holds. AVG is a sum and a count; a global aggregate is one group.
Skew and cold state
The cheapest stateful operator: one entry per group, whatever the groups hold. Hot groups stay cached, and a group nobody touches costs nothing but its bytes.

MIN, MAX

SELECT customer, max(amount) GROUP BY customer

Work per tick
O(K) in the common case · O(m) when a group loses its extremes
State
up to O(N)

reads 4 of 38 entries, writes 2

How
Not linear: deleting the maximum needs the runner-up, which a running total cannot give. Each group keeps a count per distinct value, sorted so the extreme comes first, and a change reads the front of its group to find the extreme before and after.
The Nth row
A read of the first few values of its group. A delete of the current maximum reads on to the next value still present; a group that loses every value it had sorted first is read further.
Skew and cold state
Reads touch the front of each touched group only, so the long tail of values in a large group stays cold. An append-only input never deletes, so the read is always the front.

COUNT(DISTINCT), STRING_AGG

SELECT plan, count(DISTINCT customer) GROUP BY plan

Work per tick
O(Σ m)
State
O(N)

reads 16 of 38 entries, writes 2

How
Neither can be updated from a running total: a distinct count needs to know whether the value was already there, a string needs its order. Each group keeps its rows (or values), and a change recomputes its touched groups from all of them.
The Nth row
Reads every row of its group. Without a GROUP BY that is the whole input on every change.
Skew and cold state
The hot group is the expensive one: a group of a million rows is read whole on every change to it. A rewrite as two aggregates (a DISTINCT, then a COUNT) turns it into two linear steps.

Join

orders JOIN customers ON orders.customer_id = customers.id

Work per tick
O(Σ m)
State
O(L + R)

reads 16 of 38 entries, writes 2

How
Both inputs are indexed by the join key. A changed row on one side is matched against the other side's rows under its key: ΔL ⋈ R + L ⋈ ΔR + ΔL ⋈ ΔR. Outer, semi and anti joins compute a key's output before and after the change and emit the difference. Both sides are kept whole, since a change on either can change the output of every row under its key on the other.
The Nth row
Reads the other side's rows under its key: one row for a lookup on a primary key, so O(Δ), and every row for a key the other side holds many times.
Skew and cold state
The hot key is the hazard. A key the other side holds a million times reads a million rows on every change on this side, and emits up to as many. A cross join is the extreme: one key holding the whole other side.

Top-k per partition

rank() OVER (PARTITION BY region ORDER BY revenue DESC) <= 10

Work per tick
O(K · k)
State
O(N)

reads 4 of 38 entries, writes 2

How
Every row is kept, in case a row in the top k is deleted and the next one moves up. A change reads only the first rows of its partition in order, k plus a margin for the deletes, and emits the rows whose rank crossed k.
The Nth row
Reads about 2k rows of its partition, whatever the partition's size.
Skew and cold state
Only the front of each touched partition is read. The rest of a large partition stays cold, and is only ever read when deletes eat through the front.

Window over a whole partition

row_number() OVER (PARTITION BY customer ORDER BY ts)

Work per tick
O(Σ m)
State
O(N)

reads 16 of 38 entries, writes 2

How
A row's value depends on every row before it (or after it), so a change reads its partition, computes the functions before and after, and emits the rows whose values moved.
The Nth row
Reads its whole partition. An insert at the start of a partition renumbers every row after it, and emits them all.
Skew and cold state
Cost follows partition size, so a partition per customer is cheap and a partition per region is not. Bound the frame (lag, a row range) or filter on the rank, and the read shrinks to the rows near the change.

Running total and bounded frames

sum(amount) OVER (PARTITION BY account ORDER BY ts) · lag(ts)

Work per tick
O(Δ) for appends · O(m) for a change at the start
State
O(N)

reads 2 of 38 entries, writes 2

How
A bounded frame (lag, lead, a row range) needs only the rows within its reach of the change. A running total keeps a per-partition total beside the rows, so an append reads the tail and the total, and a change earlier in the partition re-emits every row after it.
The Nth row
An appended row reads the last row of its partition and the total. A row inserted at position i changes the running value of the m − i rows after it.
Skew and cold state
Appends land at the tail of each partition, so the hot set is the partitions' last rows. History behind them stays cold.

Recursive query, incremental

WITH RECURSIVE reachable AS (…)

Work per tick
O(Σ m) at each level the change reaches
State
O(N), plus each level's operators' state

reads 4 of 38 entries, writes 2

How
The DBSP form: the fixed point unrolled into levels, level 0 the static term and level i + 1 the step over level i, each level an incremental circuit of its own. A change flows down the levels as a delta and costs what each level's joins and aggregates read for it, which for a walk over a graph is the neighbours of the changed rows at each level.
The Nth row
One changed edge reads its neighbours at each level it reaches, and stops at the first level it does not change.
Skew and cold state
Each level's operators follow their own rows above: a hub's neighbours are read whenever the hub changes. The levels a change never reaches stay cold, and so does everything below the first level it leaves unchanged.

Recursive query, recomputed

WITH RECURSIVE reachable AS (…)

Work per tick
O(N)
State
O(N)

reads 38 of 38 entries, writes 2

How
An engine's choice for a small recursion rather than a bound on the operator: the fixed point is recomputed in memory when an input changes, each round seeing only the rows the round before added, and the difference from the last result is emitted. Cheaper than the levels while the recursion is small, since it keeps no state per level.
The Nth row
Every tick that changes an input reads the inputs and the last result whole.
Skew and cold state
Nothing can stay cold, which is why it is a mode for small recursions and not for large ones.

06 · Skew and cold state

Most state is cold most of the time.

An operator reads only the keys a tick touches. How many distinct blocks that adds up to depends on the keys: an even spread touches new blocks every tick, a few hot keys touch the same few, and increasing keys touch only the end.

Key distribution

Such as: a customer id in a join, where a few large customers make most of the changes.

Working set: 12 of 240 blocks (5%) touched in the last 25 ticks

Each square is a block of state, sorted by key. Each tick a dozen changed rows light the blocks their keys fall in; a block stays warm for 25 ticks, about what a cache would keep. Everything grey can live on disk. The line is the working set over the last 60 ticks.

07 · Rules of thumb

Writing SQL that stays cheap to keep current.

What the complexities above say about how to write an incrementally maintained query, in any engine.

  1. Prefer linear aggregates. SUM, COUNT and AVG cost one entry per touched group. COUNT(DISTINCT) and STRING_AGG read the whole group; a DISTINCT followed by a COUNT is two linear steps.
  2. Join on keys that are unique on one side. A lookup into a dimension reads one row; a many-to-many key reads all of them, every time, and emits as many.
  3. Bound your windows. A rank filter, a row range or lag turns a read of the partition into a read of a few rows. A running total over an append-only stream is cheap; over an input that changes history, it is not.
  4. Put a time window on streams. A filter on now() is the one operator that lets state shrink, and it shrinks the state of everything above it.
  5. Watch the hottest key, not the average. Per-tick cost follows Σ m, and one tenant or customer with a million rows sets it. The explorer's hot-key slider is that tenant.
  6. State nobody touches can live on disk. When work per tick is far below state, the cold majority belongs on cheap storage with only the touched keys in memory. That is most operators, most of the time.

Further reading: Budiu, McSherry, Ryzhyk and Tannen, DBSP: Automatic Incremental View Maintenance for Rich Query Languages, VLDB 2023. The figures on this page come from the operator notes of eddy’s engine and hold for any engine that keeps its state by key.