Proposal: managed materializations¶
Status: stage 1 is built and shipped — see Materializations. Stages 2 and 3 below (incremental refresh, seamless rewrite) are still proposed. This supersedes the result-cache idea in Pre-aggregation vs. a result cache.
The idea¶
Precompute an expensive query into a DuckDB table, and serve requests from that table instead of the base data. Two shapes of it, and they are complements rather than alternatives.
Measured on 10,000,000 rows, with a panel expensive enough to be worth the trouble — top five
models by revenue within each region, which is a GROUP BY feeding a windowed rank():
| build | serve the panel | a question it never saw | |
|---|---|---|---|
| base table | – | 39.6 ms | 21.9 ms |
| the panel's result, materialised (30 rows) | 39.8 ms | 0.4 ms — 99x | cannot answer it |
| a dimensional pre-aggregate (18,000 rows) | 55.2 ms | 2.1 ms — 19x | 0.9 ms — 24x |
The rewritten panel returns identical results — verified row by row against the base table, zero differences.
Two things worth noticing:
The narrow one is five times better at the job it was built for. 0.4 ms against 2.1 ms, because 30 rows is less than 18,000. If you know the panel, materialise the panel.
It is still a table, so it is still queryable. Filtering the materialised result a different way
(WHERE region = 'R2' ORDER BY revenue DESC) also costs 0.4 ms. That is the difference between
this and caching the answer on the Java heap: a cached List is opaque and sits in the heap this
project exists to keep empty, while a materialised table is columnar, compressed, off-heap, and can
still be sliced.
What actually needs library support¶
Not the precomputation — that is CREATE TABLE AS SELECT and needs nothing from us. Refreshing it
without breaking readers does.
The obvious refresh is DROP TABLE then CREATE TABLE AS. With four threads reading while fifteen
refreshes ran:
| refresh style | reader failures |
|---|---|
DROP then CREATE TABLE AS |
1,280 — Catalog Error: Table with name mv does not exist! |
| build aside, then swap in one transaction | 0 (7,488 successful reads meanwhile) |
CREATE OR REPLACE TABLE mv_next AS SELECT ...;
BEGIN; DROP TABLE mv; ALTER TABLE mv_next RENAME TO mv; COMMIT;
DuckDB's catalog is transactional, so readers see either the old table or the new one and never the gap between them. That is a non-obvious, entirely mechanical piece of correctness, which is exactly the sort of thing a library should own — anyone writing the obvious version ships the bug.
The shape¶
Materialization topModels = db.materialize("top_models")
.as("SELECT region, make, model, sum(price) AS revenue,"
+ " rank() OVER (PARTITION BY region ORDER BY sum(price) DESC) AS rank"
+ " FROM sale GROUP BY 1,2,3")
.build();
// Just a table. Query it like one - filtered, joined, aggregated further.
db.query("SELECT * FROM top_models WHERE region = ? AND rank <= 5", "R2").records(Row.class);
topModels.refresh(); // atomic; readers never see it missing
topModels.builtAt(); // when it was last refreshed
topModels.rowCount();
And for the dimensional case, where the measures are additive, the incremental refresh already measured in Aggregates — 3.5 ms to fold in 50,000 new rows against 13.8 ms for a rebuild:
db.materialize("sales_rollup")
.as("SELECT region, make, colour, year, count(*) AS n, sum(price) AS total"
+ " FROM sale GROUP BY 1,2,3,4")
.incrementalOn("saleId") // fold in rows past the watermark instead of rebuilding
.build();
Should it be seamless?¶
Automatically rewriting a query against the base table into one against a materialization is what Oracle calls materialised-view rewrite and ClickHouse calls projections. My first answer was that it is a research problem, because deciding whether one SQL statement can be answered from another is query containment, and we would have to parse SQL we did not generate.
That was wrong about the hard part. DuckDB exposes its own parser, in both directions:
SELECT json_serialize_sql('SELECT make, count(*) AS n, sum(price) AS total
FROM sale WHERE region = ''R0'' GROUP BY 1');
-- {"statements":[{"node":{"type":"SELECT_NODE","select_list":[...],
-- "from_table":{"table_name":"sale"},"where_clause":{...},"group_expressions":[...]}}]}
SELECT json_deserialize_sql(json_serialize_sql('SELECT 1')); -- back to SQL, and it runs
The AST hands over exactly the four things a match needs: table_name, select_list with
each function_name, group_expressions, and where_clause. So we never write a
parser — the engine that will execute the query is the same one that parses it for us.
That makes a restricted rewrite tractable. A query is rewritable against materialization M when:
- it reads exactly the base table of M;
- every grouping expression is one of M's dimensions;
- the
WHEREclause references only M's dimensions; - every aggregate is derivable —
count(*)→sum(n),sum(x)→sum(total_x),min/maxdirectly,avg(x)→sum(total_x)/sum(n).
Anything else — a median, an exact count(DISTINCT), a filter on a column M does not carry, a
join, a window over the base grain — falls through to the original query untouched. That
fall-through is the whole safety argument: a bug in matching costs performance, never correctness.
Can a rewrite give a different answer? Yes — I broke it four ways¶
This is the question that decides whether the feature is worth having, so I tried hard to produce a
wrong answer from a rewrite a reasonable matcher would have accepted. 300,000 rows, some NULLs, a
rollup of count(*), sum(price):
| hazard | base table | rewritten | |
|---|---|---|---|
sum of a DOUBLE |
29400.000882016073 | 29400.000881985878 | differs |
count(*) → sum(n) |
300000 | 300000 | same |
count(price) → sum(n) |
294000 | 300000 | differs by 2% |
avg(price) → sum(total)/sum(n) |
0.100000003 | 0.098000003 | differs by 2% |
max(price) → max(total_price) |
0.100000006 | 9800.000294 | catastrophically wrong |
sum of an integer |
592208 | 592208 | same |
Every one of those is a plausible-looking number on a dashboard. Nobody spots them.
But the failures are two different species, ten orders of magnitude apart.
Semantic errors — 2e-2. count(*) counts NULLs and count(x) does not; deriving an average as
sum/count(*) divides by the wrong denominator; a max taken over pre-summed groups is not a max
of anything. These are bugs, entirely preventable, and they are the catastrophic ones.
Re-association — 1e-12. Floating-point addition is not associative, so summing 18,000 partial sums does not give bit-identical results to summing 300,000 values. This one cannot be fixed by a better matcher. It can, however, be avoided by not using floating point:
| measure type | rewritten result |
|---|---|
DOUBLE |
differs, relative error 1.05e-12 |
DECIMAL(18,4) |
exact |
BIGINT |
exact |
count(*) |
exact |
avg over DECIMAL |
exact |
Money should be DECIMAL anyway. With DECIMAL and integer measures a rewrite is bit-for-bit
identical, and the hazard disappears rather than being tolerated.
The rules this imposes¶
- Never infer a measure's aggregate from its column name. The materialization declares
n = count(*),n_price = count(price),total_price = sum(price), and a query is rewritten only against the exact aggregate that was declared. Themaxdisaster above is what name-guessing produces. count(*)andcount(x)are different measures. Conflating them cost 2%.avg(x)becomessum(total_x)/sum(n_x)wheren_xiscount(x)— nevercount(*).- A
DOUBLEmeasure is re-associated. Declare it, do not hide it: results will differ around the twelfth significant digit. Refuse to rewriteDOUBLEmeasures at all under a strict setting.
The condition I would put on it¶
Seamless is only safe if it is verifiable. Because both queries are available, quackjvm can run them both and compare:
DuckDBDatabase db = DuckDBDatabase.builder()
.materializationRewrite(Rewrite.VERIFY) // OFF | ON | VERIFY
.build();
VERIFY runs the rewritten query and the original and fails loudly if they disagree — slow, and
exactly what you want in a test suite or a staging soak. You turn it on against your own panels,
prove the rewrite on your own data, then switch to ON in production.
And the measurements above are what make VERIFY workable. It cannot demand bit equality,
because a DOUBLE measure legitimately differs at 1e-12. It does not have to: the gap between
unavoidable noise (1e-12) and the smallest bug found here (2e-2) is ten orders of magnitude, so
a relative tolerance anywhere around 1e-9 passes every correct rewrite and catches every incorrect
one. Without measuring both species there would have been no way to size that tolerance, and a
verify mode that cannot tell noise from error is worth nothing.
The real risk, which is not correctness¶
json_serialize_sql is an internal debugging facility, not a stable public API. Its shape can
change between DuckDB releases, and quackjvm pins a DuckDB version but users override that. So:
treat any unexpected AST shape as "not rewritable" and fall through; never throw. The failure mode
of a version bump must be that the rewrite quietly stops happening, not that queries break.
So: staged¶
- ~~Routing by name first.~~ Built —
database.materialize(name).as(sql).build(), with the atomic refresh. See Materializations. - Then the restricted rewrite,
OFFby default, withVERIFYfor adoption. - Never a general rewrite. Containment over arbitrary SQL is not worth attempting here.
What this costs¶
- Storage. A materialization is a real table. The dimensional one above is 0.18% of the base
table; a badly chosen one, grouped on something high-cardinality like
modelplusid, would be as large as the data. The row count is the thing to watch. - Staleness. Unchanged from any cache: it is as old as its last refresh. The difference is that a stale materialization still answers every question it covers, slightly behind, while a stale cache entry is simply wrong for the one query it holds.
- Additive measures only, for incremental refresh.
count,sum,min,max, and averages derived from sum and count. A median or an exact distinct count needs a full rebuild — which at 55 ms on ten million rows is not much of a hardship.
Recommendation¶
Build it. It is the first thing in this thread that pays for its API surface:
- Atomic refresh is a real correctness bug that users will otherwise ship — 1,280 failures against zero, measured.
- Lifecycle (create, refresh, staleness, drop on close) is tedious and mechanical, which is what a library is for.
- The speedups are large and hold at scale: 19–99x on ten million rows, and 24x on questions the materialization was never designed for.
- It needs no query parsing, no heap, and no new staleness policy invented on the caller's behalf.
Start with routing by name and manual refresh(). Add incremental refresh for additive measures.
Leave automatic rewrite alone until there is a reason to believe it is needed.
Reproduce: io.quackjvm.cqengine.bench.PreAggregateBenchmark and the materialisation and swap
measurements recorded here.