12 KiB
DuckLake-native materialization
Design sketch for "managed, versioned, incremental" assets built on the
DuckLake substrate. This is a companion to pipelines-vs-dbt.md
and extends its Path C (hybrid, partition-first) recommendation. The new
contribution here is leveraging DuckLake's snapshot/time-travel layer, which
the earlier doc's incremental deep-dive did not use. The annotation grammar is
reconciled with that doc — // partitioned + // unique_key + // append
stay canonical; nothing here forks a competing vocabulary.
The core reframe
dbt had to build a materialization engine (compile SQL → CREATE TABLE AS /
incremental MERGE / SCD2 snapshots) because the warehouse gives it nothing
but raw SQL over a mutable table. DuckLake hands us, at the storage layer, the
four things that engine exists to provide:
| Capability | Source |
|---|---|
| ACID multi-statement transactions | DuckLake |
A snapshot per commit + time-travel (AT (VERSION => n) / AT (TIMESTAMP => ...)) |
DuckLake |
Physical partitioning + pruning (ALTER TABLE … SET PARTITIONED BY (…)) |
DuckLake |
| Schema evolution tracked in the catalog DB | DuckLake |
Windmill already attaches DuckLake fully — see transform_attach_ducklake
(backend/windmill-worker/src/duckdb_executor.rs:661), which rewrites a user's
ATTACH 'ducklake://name' AS dl into the real
ATTACH 'ducklake:postgres:…' AS dl (DATA_PATH 's3://…', OVERRIDE_DATA_PATH TRUE, AUTOMATIC_MIGRATION TRUE) at duckdb_executor.rs:730. But today every
write is a destructive overwrite and all four capabilities above are thrown
away — snapshots are never surfaced, partitioning is purely an orchestration
concept disconnected from the physical layout.
So the materialization engine we need is a thin layer — a write-strategy
wrapper + snapshot capture — not a dbt rebuild. This is exactly the
"buy A's 80% without B's dialect-rewriting tax" tradeoff pipelines-vs-dbt.md
argued for; DuckLake is what makes the remaining 20% (versioning, reproducible
reads, materialization history) nearly free instead of a second project.
Annotation grammar (final)
One self-documenting line; managed-by-default. Strategy options live on the
materialize line (they have no meaning without it), while // partitioned
stays separate because it is cross-cutting (cascade + scheduling + materialize).
// materialize ducklake://analytics/orders_daily → managed, replace (default)
// materialize ducklake://analytics/orders_daily key=order_id → managed, merge
// materialize ducklake://analytics/orders_daily append → managed, append
// materialize manual ducklake://analytics/orders_daily → track-only escape hatch
- managed (default) — the script is setup + one trailing
SELECT; Windmill generates the write DDL, captures the DuckLake snapshot, and records state. DuckDB-only; validated at deploy (a non-SELECT script is rejected with a clear error pointing to thewmll.ducklakehelpers). manual— escape hatch: the script writes its own DDL; Windmill only records state (no snapshot capture, no idempotency guarantee). Rare; explicit.key=<col>→ MERGE (dedup within slice);append→ INSERT-only; neither → DELETE-by-partition + INSERT (replace).appendwins overkeyif both are given (deploy warning).// partitioned <kind>— unit of work + state + backfill (separate; cross-cutting). Polyglot / multi-statement writes use thewmll.ducklakehelpers instead of// materialize.
There is no wrap keyword — materialize is "manage the write," so it was
redundant; the only reason for it was to carve out the weak track-only mode,
which is now the explicit manual opt-out.
DuckLake snapshots are orthogonal to all of the above — they apply to every strategy automatically because every write is a DuckLake commit. The user never annotates for versioning; they get it.
Executor codegen
The seam
run_duckdb already splits the script into statement blocks and rewrites
custom ATTACH blocks in a single pass before execution
(duckdb_executor.rs:114-160):
let query_block_list = parse_sql_blocks(&query, true);
// each block: remove_comments → if ducklake/datatable ATTACH, expand; else passthrough
All blocks run in order on one DuckDB connection. Materialize has two modes:
-
Managed (default). The user writes setup + one trailing
SELECT. Windmill replaces that SELECT with generated statements, wrapped in an explicit DuckLake transaction it controls — never textualBEGIN/COMMITinjected around the user's other statements (fragile across their ownATTACHs and multi-statement SQL). For a partitionedreplacethe SELECT block expands to:-- generated for: // partitioned daily ; target = dl.orders_daily ; partition = '2026-06-19' CREATE TABLE IF NOT EXISTS dl.orders_daily AS SELECT *, CAST(NULL AS VARCHAR) AS _wm_partition FROM (<user_select>) WHERE false; -- first-run bootstrap ALTER TABLE dl.orders_daily SET PARTITIONED BY (_wm_partition); BEGIN TRANSACTION; DELETE FROM dl.orders_daily WHERE _wm_partition = '2026-06-19'; INSERT INTO dl.orders_daily SELECT *, '2026-06-19' AS _wm_partition FROM (<user_select>); COMMIT;The strategy variants are all DELETE+INSERT-shaped — no
MERGE INTO, which DuckLake can't reliably run on a fresh partition (it 404s writing the first rows):- whole-table replace (no
// partitioned) → a singleCREATE OR REPLACE TABLE … AS <user_select>(handles schema changes, still snapshots). key=<col>→ `DELETE FROM … WHERE [ AND] IN (SELECT FROM ())` then `INSERT` (upsert within the slice).append→ theDELETEis dropped (insert-only).
- whole-table replace (no
-
manual. The user writes their own DDL inside their ownBEGIN … COMMIT; Windmill injects nothing into the body and only records state (no snapshot capture, no idempotency guarantee).
The _wm_partition column is the physical link the orchestration layer lacks: on
first materialize Windmill runs ALTER TABLE … SET PARTITIONED BY (_wm_partition)
so DuckLake prunes on read and DELETE-by-partition rewrites only that partition's
Parquet files.
Storage: DuckLake writes go to
s3://_default_/through the windmill S3 proxy (/api/w/{ws}/s3_proxy, gated behind theparquet+privatefeatures). The proxy must sign the SigV4 canonical URI with single percent-encoding — the SigV4 default (Double) 401s Hive-partition keys like_wm_partition=2026-06-19(the=double-encodes to%253Dvs the client's%3D).
Run summary capture
After the generated blocks, Windmill appends one read block — it is both the job's result (a useful preview rendered as the materialized table) and the row it records:
SELECT 'ducklake://<name>/<table>' AS materialized,
'<partition>' AS partition, -- only when partitioned
(SELECT count(*) FROM <target> [WHERE _wm_partition = '<partition>']) AS rows,
(SELECT max(snapshot_id) FROM ducklake_snapshots('<target>')) AS snapshot_id;
The snapshot_id and rows are persisted as materialized_partition metadata.
One extra round-trip per materialization, no new infra.
Metadata schema
Extends the materialized_partitions table proposed in pipelines-vs-dbt.md
§"First implementation slice" with the DuckLake snapshot id:
materialized_partition (
workspace_id TEXT,
asset_kind TEXT, -- 'ducklake'
asset_path TEXT, -- 'analytics/orders_daily'
partition TEXT, -- '2026-06-19' (NULL for unpartitioned)
snapshot_id BIGINT, -- DuckLake snapshot produced by this materialize
row_count BIGINT,
job_id UUID,
materialized_at TIMESTAMPTZ,
PRIMARY KEY (workspace_id, asset_kind, asset_path, partition)
)
This one table drives four things at once:
- Observability — "last materialized: snapshot 42, 1.2M rows, 09:14" per asset node (closes the Dagster-catalog gap from the v1-readiness review).
- Run-stale / gap detection — which partitions exist, which are missing.
- Backfill — the missing/failed set is the backfill worklist.
- Snapshot pinning — see below.
Reproducibility — the beyond-dbt part
Because every materialization records the snapshot it produced, a downstream consumer can read the exact upstream snapshot its run saw:
FROM dl.orders_daily AT (VERSION => $WM_UPSTREAM_SNAPSHOT)
The cascade already threads a trigger blob (producer path, partition) to each
subscriber; add the producer's captured snapshot_id to it, and a consumer's
read is pinned to the upstream state at dispatch time. That makes the whole
pipeline reproducible and time-travelable — something dbt has no native answer
for (dbt models are always "whatever's in the warehouse now"). It also gives
rollback (re-point an asset to snapshot N) and "what did this table look like at
the failing run" debugging, for free off the same captured ids.
This is the differentiator worth leaning on. It is not catch-up to dbt; it is a capability dbt structurally cannot offer, and DuckLake gives it to us at the cost of recording one integer per run.
It also means we do not build SCD2 snapshots (gap #4 in pipelines-vs-dbt.md):
DuckLake time-travel is a strictly better answer for most of what dbt's
{% snapshot %} is used for. One fewer engine to write.
Scoping decision: DuckLake vs DataTable
Make DuckLake the materialization/versioning substrate; keep DataTable as the
live operational table with no versioning. DataTable is plain Postgres
(transform_attach_datatable, duckdb_executor.rs:742) — no native snapshots
or time-travel — so giving it the versioned/incremental story means building
MVCC-on-top ourselves (history tables, SCD2), precisely the complexity this
DuckLake approach exists to avoid. Clean split:
ducklake://→ analytics, versioned, reproducible, backfillable.datatable://→ mutable app/operational state; partition idempotency via DELETE+INSERT still works, but no snapshot/time-travel layer.
Don't try to give both the full treatment for v1.
v1 slice (smallest viable)
- Partition runtime context — resolve
(value, start, end)and surface asWM_PARTITION*bind/env (Path C step 1; partly built per the pipeline-partition-runtime work). - Physical partition wiring —
_wm_partitioncolumn +SET PARTITIONED BYon first materialize forducklake://targets. - Strategy templates — DELETE+INSERT default (
CREATE OR REPLACEfor the whole table); delete-by-key + insert whenkey=<col>; INSERT-only whenappend. Managed// materializewraps a single-SELECT DuckDB script behind these templates;// materialize manualopts out. - Snapshot + metadata capture — append
ducklake_snapshotsread, persistmaterialized_partitionrows. - Surface it — last-materialized/snapshot/row-count on the asset node; missing-partition set feeds the backfill UI.
- v1.x — snapshot pinning across the cascade (
$WM_UPSTREAM_SNAPSHOT), rollback, time-travel read helper.
Steps 1–5 are a thin annotation+template layer plus one metadata table and one extra read per run. They deliver managed/incremental/versioned assets, idempotent partitioned materialization, the backfill substrate, and materialization observability together — and stay recognizably Windmill-shaped.
Open decisions
These ride on top of the six in pipelines-vs-dbt.md §"Decisions either path
forces"; DuckLake-specific:
- Bootstrap of
SET PARTITIONED BY. First-materialize detection — table absent vs. present-but-unpartitioned. Idempotent re-apply. - Snapshot retention / compaction. DuckLake snapshots accumulate; when do we expire old ones, and does pinning hold a snapshot alive past retention?
- Pin scope. Pin only direct producers, or the full transitive upstream set per run? Storage and "stale pin" semantics differ.
- Managed multi-statement. Resolved: managed
// materializeaccepts setup statements (ATTACH/SET/…) followed by exactly one trailing SELECT, and rejects anything else at deploy with a clear error pointing to// materialize manual. The classifier (sql_materialize.rs) is the single source of truth.