Δ(Q(R)) = Q(ΔR)scan, projection, filter, UNION ALL, SUM, COUNT
proportional to |Δ|, no extra state
How it works
OpenIVM does not replace your engine. It reads your view definition and emits the SQL that keeps it up to date, so the same approach carries from DuckDB to Spark and beyond.
CREATE MATERIALIZED VIEW with standard SQL: joins, aggregates, CTEs, subqueries.
The logical plan is classified operator by operator and rewritten into delta propagation rules. The output is plain SQL.
An optimizer rule intercepts DML on base tables and writes signed rows to openivm_delta_<table>, or reads DuckLake snapshots.
On demand or on a schedule, consolidated deltas are pushed through the compiled SQL and merged into the view. Pipelines cascade.
The model
Given a table T and a view V = Q(T), find a functionf so that f(ΔT) = ΔV: compute the change to the view from the change to the table alone, instead of re-running Q.
OpenIVM follows DBSP, which gives such an f for arbitrary relational queries using two operators: differentiation D (ΔT = T′ −T) and integration I (T + ΔT = T′). Each operator is replaced by its incremental form, and the plan runs over deltas.
Selection and projection are their own incremental form; a join becomes three joins.
Z-sets: every row carries an integer weight, and Rnew = Rold + ΔR
{ apple → 5, banana → 2 }{ apple → −3, banana → +1 }{ apple → 2, banana → 3 }…and every algebraic step is plain SQL
UNION ALL+ grouped weight sumsJOINweights multiplyMERGE/ INSERT, DELETEDelta tables
A delta table has the same columns as its source, plus a signed weight. An INSERT writes+1, a DELETE writes −1, and anUPDATE writes both (old row −1, new row +1) in one atomic insert. Operators then work on these weights algebraically: aggregates sum weight × value, joins multiply weights.
Several views can share one delta table; each keeps its own cursor and reads only what it has not yet processed. Over DuckLake, there is no delta table at all: changes are derived from table snapshots, and the pre-refresh state is read through time travel.
SELECT * FROM openivm_delta_sales;
| region | product | amount | multiplicity |
|---|---|---|---|
| US | Bolt | 50 | +1 |
| JP | Gear | 300 | +1 |
| EU | Gadget | 200 | −1 |
After two inserts and one delete. Timestamps omitted. Captured from the DuckDB extension.
Architecture
Most IVM compilers are detached prototypes that talk to a database through drivers. OpenIVM lives inside one: DuckDB parses, binds and plans the view, so the compiler gets schema validation, cardinality estimates and a cost model for free, and physical choices such as indexes or what to materialize stay tunable.
CREATE MATERIALIZED VIEW, take the bound logical plan, and rewrite equivalent shapes into forms the delta rules handle.Generated SQL
One view definition becomes several SQL statements: create the deltas, compute ΔV from ΔT, fold ΔV into V and drop rows whose count reached zero, then clear the consumed deltas. Below is what OpenIVM compiles for regional_totals, captured withSET openivm_files_path and lightly reformatted (catalog prefixes and literal timestamps elided). It is ordinary SQL, which is what makes the approach portable.
A grouped aggregate over one table. This is all you write.
CREATE MATERIALIZED VIEW regional_totals AS
SELECT region, SUM(amount) AS total, COUNT(*) AS cnt
FROM sales
GROUP BY region;The view is stored in openivm_data_regional_totals with a hidden helper column (the non-NULL count behind SUM, so a group that loses all its values returns NULL, not 0). A plain view hides it; a delta table is created for sales.
CREATE TABLE openivm_data_regional_totals AS
SELECT region, sum(amount) AS total, count_star() AS cnt,
count(amount) AS openivm_nonnull_sum_count_1
FROM sales
GROUP BY region;
CREATE VIEW regional_totals AS
SELECT * EXCLUDE (openivm_nonnull_sum_count_1) FROM openivm_data_regional_totals;
CREATE TABLE IF NOT EXISTS openivm_delta_sales AS
SELECT *, 1::INTEGER AS openivm_multiplicity,
make_timestamp(epoch_us(now())) AS openivm_timestamp
FROM sales LIMIT 0;
CREATE UNIQUE INDEX openivm_data_regional_totalsopenivm_index
ON openivm_data_regional_totals(region);On refresh, only rows newer than the view's cursor are read. Duplicate deltas are consolidated first (a +1 and a −1 for the same tuple cancel), then the aggregate is applied to the delta itself, keeping the multiplicity as a group key.
WITH
t0_scan (t0_region, t0_amount, t0_openivm_multiplicity) AS (
SELECT region, amount, openivm_multiplicity
FROM openivm_delta_sales
WHERE openivm_timestamp >= $last_update
),
t1_aggregate (t8_region, t8_amount, t9_openivm_multiplicity) AS (
SELECT t0_region, t0_amount, sum(t0_openivm_multiplicity)
FROM t0_scan
GROUP BY t0_region, t0_amount
),
t2_block (t2_region, total, cnt, openivm_nonnull_sum_count_1, t1_openivm_multiplicity) AS (
SELECT t8_region, sum(t8_amount), count_star(), count(t8_amount),
CAST(t9_openivm_multiplicity AS INTEGER)
FROM t1_aggregate
WHERE t9_openivm_multiplicity != 0
GROUP BY t8_region, CAST(t9_openivm_multiplicity AS INTEGER)
)
INSERT INTO openivm_delta_regional_totals
(region, total, cnt, openivm_nonnull_sum_count_1, openivm_multiplicity)
SELECT * FROM t2_block;The view's own delta is folded with SUM(multiplicity × value) and merged on the group key. Groups whose count drops to zero are deleted. Downstream views can consume openivm_delta_regional_totals before it is cleared.
WITH refresh_cte AS (
SELECT region,
sum(openivm_multiplicity * total) AS total,
sum(openivm_multiplicity * cnt) AS cnt,
sum(openivm_multiplicity * openivm_nonnull_sum_count_1) AS openivm_nonnull_sum_count_1
FROM openivm_delta_regional_totals
WHERE openivm_timestamp > $last_update
GROUP BY region
)
MERGE INTO openivm_data_regional_totals v USING refresh_cte d
ON v.region IS NOT DISTINCT FROM d.region
WHEN MATCHED THEN UPDATE SET
total = CASE WHEN COALESCE(v.openivm_nonnull_sum_count_1 + d.openivm_nonnull_sum_count_1,
v.openivm_nonnull_sum_count_1, d.openivm_nonnull_sum_count_1) = 0
THEN NULL
ELSE COALESCE(v.total + d.total, v.total, d.total) END,
cnt = COALESCE(v.cnt + d.cnt, v.cnt, d.cnt),
openivm_nonnull_sum_count_1 = COALESCE(v.openivm_nonnull_sum_count_1 + d.openivm_nonnull_sum_count_1,
v.openivm_nonnull_sum_count_1, d.openivm_nonnull_sum_count_1)
WHEN NOT MATCHED THEN INSERT (region, total, cnt, openivm_nonnull_sum_count_1)
VALUES (d.region, d.total, d.cnt, d.openivm_nonnull_sum_count_1);
DELETE FROM openivm_data_regional_totals WHERE COALESCE(cnt, 0) = 0;Consumed deltas are removed once every dependent view has moved past them, and the cursor advances to MAX(timestamp) + 1µs of what this refresh actually saw, not now(), so a concurrent insert is never skipped or applied twice.
DELETE FROM openivm_delta_regional_totals;
DELETE FROM openivm_delta_sales
WHERE openivm_timestamp < (SELECT MIN(last_update) FROM openivm_delta_tables
WHERE table_name = 'openivm_delta_sales');
UPDATE openivm_delta_tables
SET last_update = COALESCE((SELECT MAX(openivm_timestamp) + INTERVAL '1 microsecond'
FROM openivm_delta_sales),
make_timestamp(epoch_us(now()))),
last_refresh_ts = make_timestamp(epoch_us(now()))
WHERE view_name = 'regional_totals' AND table_name = 'openivm_delta_sales';Delta rules
Each node of the plan carries a rule kind that determines the shape of its delta and what state it needs. It is the same taxonomy DBSP uses (Budiu et al., VLDB 2023), and it tells you, before the first refresh, what a view will cost to maintain.
Δ(Q(R)) = Q(ΔR)scan, projection, filter, UNION ALL, SUM, COUNT
proportional to |Δ|, no extra state
Δ(R ⋈ S) = ΔR ⋈ S + R ⋈ ΔS − ΔR ⋈ ΔSINNER, CROSS, LEFT, RIGHT, FULL OUTER JOIN
2ᴺ − 1 terms for N tables (or N with telescoping)
needs accumulated stateDISTINCT, SEMI / ANTI JOIN, MIN / MAX with deletes, window functions
aux state, or recompute of affected groups / partitions
no delta ruleanything else
detected at CREATE time; full refresh
Inclusion–exclusion, over current state
Δ(R ⋈ S) = ΔR ⋈ Snew + Rnew ⋈ ΔS − ΔR ⋈ ΔSBy refresh time the base tables already contain their deltas. Joining each delta against the current state counts the cross-term twice, so it is subtracted once: a Möbius sign(−1)^(k−1) × ∏ wᵢ for a term using k deltas, giving 2ᴺ − 1 terms for Ntables.
Telescoping, with time travel
Δ(⋈ᵢ Rᵢ) = Σᵢ (⋈j<i Rjnew) ⋈ ΔRᵢ ⋈ (⋈j>i Rjold)When the old state is readable, for example through DuckLake snapshots, each combination of changes lands in exactly one term: N joins instead of 2ᴺ − 1. The same rule covers inner, cross and theta joins.
Outer joins and negation use match counts: a NULL-extended row is retracted when its match count goes from zero to non-zero, and reinstated when it drops back to zero. SEMI,ANTI, EXISTS and NOT EXISTS use the same zero-crossing transition.
Beyond the textbook
Correct delta rules are the floor. Because OpenIVM runs inside an engine, it can also look at the data and the change itself to skip work entirely.
πk(ΔR) ∩ πk(S) = ∅ ⇒ skipNo source changed? Skip the refresh. A filter rejects every delta row? Every downstream term is zero. New orders for customers the join never sees? That term is omitted even though the delta is non-empty.
ΔR̄(t) = Σu = t ΔR(u), keep ΔR̄(t) ≠ 0Equal tuples are summed before anything else runs, so an insert and a delete of the same row, or an update and its revert, cancel out before they cost anything.
ΔR ≥ 0 ∧ Δ-rule ≥ 0 ⇒ appendWhen both the source change and the chosen delta rule are provably non-negative, projections append directly, zero-group deletion is skipped, and MIN/MAX useLEAST/GREATEST instead of rescanning groups.
A = πk(ΔT), N = σk ∈ A(V(Tnew))For operators without a closed-form delta (MIN/MAX with deletes, windows), only the affected keys A are recomputed. A and N are materialized once and reused to delete, replace, and emit the downstream delta.
countnon-NULL = 0 ⇒ SUM = NULLNullable aggregate inputs carry a hidden non-NULL count, so a group that loses its last value returnsNULL, not 0, and lookups use NULL-safe equality, exactly as SQL grouping does.
min ord(ΔP) > max ord(Pold)A cumulative window whose new rows all come after the old ones is extended from its last stored aggregate instead of recomputing the partition; late rows fall back to partition recompute.
Cross-system
Modern pipelines span several systems: transactions in one, analytics in another. Because the maintenance steps are SQL, rendered by LPTS in each engine's dialect, they can be shipped to wherever the data lives instead of copying the data to wherever the IVM engine lives.
Operator coverage
Views that cannot be maintained incrementally are detected at creation time. Withopenivm_refresh_mode = 'auto' they fall back to a full refresh.
| Operator | Strategy |
|---|---|
| Projection, filter, expressions | incremental |
| GROUP BY · SUM, COUNT, AVG, STDDEV, VARIANCE | incremental |
| MIN / MAX | incremental (insert-only) · group recompute |
| HAVING, ungrouped aggregates, LIST | incremental |
| INNER, CROSS, LEFT, RIGHT JOIN | incremental |
| FULL OUTER JOIN | incremental (merge + recompute) |
| SEMI / ANTI JOIN, EXISTS, NOT EXISTS | aux-state incremental |
| UNION ALL, DISTINCT | incremental |
| Window functions | partition-level recompute |
| CTEs, decorrelated & scalar subqueries | incremental when the lowered plan is |
| Anything else | full refresh |
Per-operator derivations live in the operator docs.
Try it
The view below is SELECT region, SUM(amount), COUNT(*) FROM sales GROUP BY region. Changes collect as signed deltas; a refresh folds them in. Insert then delete the same row and the deltas cancel out before any work is done.
| region | amount |
|---|
| ± | region | amount |
|---|
No pending changes.
| region | total | cnt |
|---|
Read the paper
Ilaria Battiston, Kriti Kathuria, Peter Boncz · SIGMOD Companion 2024