Recipe: Idempotent Downstream Warehouse Loading
Rivet provides at-least-once file delivery to its destination. After a clean run, the destination prefix carries:
- one or more data parts — names are per-runner (single/incremental
{export}_{YYYYMMDD_HHMMSS_mmm}.parquetfor a single-part run, with a_part0.._partN-1suffix on every part when the run rotates into multiple parts; chunked{export}_{ts}_chunk{idx}_{nonce}.parquet; keyset{export}_{run-tag}_keyset_{seek-tag}.parquet, …) — always take them from the manifest, never pattern-match them, - a
manifest.jsonlisting every committed part withsize_bytesandcontent_fingerprint(xxh3 over the part bytes), - a
_SUCCESSmarker whose body fingerprints the exactmanifest.jsonbytes.
What Rivet does not provide:
- exactly-once delivery to a destination,
- exactly-once load semantics into a downstream warehouse,
- transactional coupling between the export run and a downstream
MERGE/COPY INTO.
A downstream loader that ignores the manifest and processes “every file
under this prefix” will eventually double-load — chunked retries write
new files alongside originals (RR5), and rivet repair --execute
explicitly does so. Treat the manifest as the source of truth: after a
repair it lists the chunk’s original part(s) as superseded and only the
replacement as committed, so a reader of the committed parts sees each
row once. A warehouse that already loaded the original part before the
repair still holds those rows; without a primary-key dedup it keeps them.
The built-in path: rivet load (BigQuery, Snowflake, ClickHouse)
For BigQuery, Snowflake and ClickHouse you do not build any of this — rivet load is
idempotent by construction:
- Count gate before cleanup. The load refuses to finish — and refuses to
clean the source — unless the warehouse
COUNT(*)equals the summed manifestrow_count. A partial or double load fails loudly instead of silently corrupting. - Manifest-driven, not a glob. It reads the run manifests and loads exactly their committed parts, never “every file under the prefix” — so a repair retry’s extra files or a half-finished run can’t double-load.
mode: fullOVERWRITEs. Re-running a full load re-materialises the latest snapshot; the table lands identical, not doubled (live-verified: two loads of a 3-row table → 3 rows, not 6).mode: incremental/mode: cdcappend + dedup. For mutable sources the load appends to<table>__changesand exposes a current-state view (latest-per-PK, deletes flagged) — the built-in equivalent of the manualMERGEbelow, no staging table or upsert SQL to write. On ClickHouse the log is aReplacingMergeTreeread through aFINALview (ClickHouse load).
load:
target: bigquery # or: snowflake / clickhouse (+ that target's connection keys)
project: my-proj
dataset: analytics
pk: [id] # incremental/cdc dedup key; default: the source primary key
cleanup_source: true # wipe staging only after the count gate passes
The rest of this recipe is the manual pattern — what to do for a warehouse
rivet load does not target (Redshift, Trino, Databricks, dbt), or to understand
the contract rivet load itself builds on.
The contract you build on
Every committed part is recorded in manifest.json:
{
"manifest_version": 1,
"run_id": "orders_20260523T120000.123",
"schema_fingerprint": "xxh3:…",
"row_count": 200000,
"part_count": 2,
"parts": [
{"part_id": 0, "path": "orders_20260523_120000_123_part0.parquet", "rows": 100000,
"size_bytes": 4194304, "content_fingerprint": "xxh3:…", "content_md5": "…", "status": "committed"},
{"part_id": 1, "path": "orders_20260523_120000_123_part1.parquet", "rows": 100000,
"size_bytes": 4198400, "content_fingerprint": "xxh3:…", "content_md5": "…", "status": "committed"}
]
}
(Abridged — the real manifest also records export_name, status, timestamps,
source/destination blocks, format/compression, and optional per-column
checksums. Each part’s object name is its path field.)
_SUCCESS is a single line: xxh3:<16-hex> where the hex is the
fingerprint of the exact manifest.json bytes. See
ADR-0012 for the formal
invariants (M1–M9).
The manifest gives a downstream loader two things it cannot easily recover from raw object listing:
- The exact set of parts that were committed in this run (vs.
parts left over from earlier interrupted runs, parts a repair
superseded (
status: superseded), or external writes). - A content-addressable identity per part (
content_fingerprint) that survives object renames, lifecycle migrations, and CDN copies.
content_fingerprint is the supported dedup key: it is an xxh3 over the
exact part bytes, computed deterministically. Because rivet pins the Parquet
created_by to a version-independent constant, identical rows produce
identical bytes — and therefore the same content_fingerprint — across rivet
releases, not just within one build. So a re-extraction of the same window
(e.g. a crash + --resume, or a deliberate re-run) yields parts a downstream
MERGE can dedup on by fingerprint alone. Two parts with the same
content_fingerprint are byte-identical and interchangeable; drop one.
(The fingerprint is over file bytes, not logical rows — so it is stable for
the same rows + schema + compression settings. Changing compression: or the
column projection changes the bytes, hence the fingerprint.)
Manual pattern — warehouses rivet load doesn’t target
For a warehouse rivet load doesn’t reach (Redshift / Trino / Databricks / dbt),
the pattern is the same across targets:
- Read
manifest.jsonfrom the resolved destination prefix. - Verify
_SUCCESSmatches. If it does not, abort — the export is not complete. - Load only the parts listed in
manifest.jsonwithstatus: committed. Do not glob the prefix, and skipsuperseded/quarantinedentries. - Record the manifest identity (
run_id+schema_fingerprint+_SUCCESSbody) in a downstream control table. - Deduplicate by primary key (or natural key) when merging into the final table.
- Commit the warehouse table only after the load succeeds end-to-end.
If step 4 records the run identity before step 3 starts and after step 6 finishes (an “intent + commit” pair), the load is restartable: on retry, skip any manifest already marked committed.
BigQuery pattern
-- 1. Stage parts referenced by manifest.json into a per-run staging table.
LOAD DATA INTO project.dataset.orders_stage_<run_id>
FROM FILES (
format = 'PARQUET',
-- list the exact `parts[].path` values from manifest.json — never a glob
uris = ['gs://my-bucket/exports/2026-05-23/orders/orders_20260523_120000_123_part0.parquet',
'gs://my-bucket/exports/2026-05-23/orders/orders_20260523_120000_123_part1.parquet']
);
-- 2. Tag every staged row with the run identity.
ALTER TABLE project.dataset.orders_stage_<run_id>
ADD COLUMN _rivet_run_id STRING,
ADD COLUMN _rivet_manifest_fingerprint STRING;
UPDATE project.dataset.orders_stage_<run_id>
SET _rivet_run_id = '<run_id>', _rivet_manifest_fingerprint = '<xxh3>'
WHERE _rivet_run_id IS NULL;
-- 3. Merge into final table on primary key.
MERGE project.dataset.orders AS target
USING project.dataset.orders_stage_<run_id> AS src
ON target.id = src.id
WHEN MATCHED THEN UPDATE SET ... -- or DO NOTHING for append-only sources
WHEN NOT MATCHED THEN INSERT ROW;
-- 4. Drop the staging table after the merge commits.
DROP TABLE project.dataset.orders_stage_<run_id>;
For LOAD DATA, prefer the explicit URI list from manifest.parts
over a wildcard. The wildcard form is fine when you trust the prefix
contains exactly the manifest’s parts (i.e. there is no concurrent
write into the same prefix), but the explicit list is what lets you
prove which bytes were loaded.
Native types. Bare autoload degrades several columns:
json/uuidload asBYTES, a naivetimestampasTIMESTAMP(an instant, not wall-clockDATETIME), and arrays as a nestedRECORD. BigQuery will not coerce these on load — declaring native types in a load schema is rejected — so recover them with a post-loadCREATE TABLE … AS SELECTover the staging table. Load the staging table with--parquet_enable_list_inference(so arrays flatten withUNNEST), then run the recovery SQL thatrivet check --type-report --target bigqueryprints per export. Full table: type-mapping.md § BigQuery autoload & recovery.
Wide decimals (
NUMERIC/DECIMALprecision > 38). The load’s default decimal target isNUMERIC(max precision 38, scale 9) and fails on wider columns; pass--decimal_target_types=NUMERIC,BIGNUMERIC,STRING(repeated flag inbq: one value per flag).BIGNUMERICstill caps at ~5.79×10³⁸ — 38 integer digits — so aDECIMAL(50,10)column loads but a value with 39+ integer digits does not fit ANY BigQuery numeric type; withSTRINGin the list such columns load losslessly as text (verified live). Caveat for verification tooling: DuckDB misreads Parquet fixed-len-byte-array decimals wider than 16 bytes as a garbageDOUBLE— cross-check wide-decimal columns with ClickHouse or BigQuery, not DuckDB.
Snowflake pattern
-- 1. COPY INTO a staging table, listing the exact files from manifest.json.
COPY INTO @my_stage/orders/orders_stage_<run_id>
FROM ('@my_stage/exports/2026-05-23/orders/orders_20260523_120000_123_part0.parquet',
'@my_stage/exports/2026-05-23/orders/orders_20260523_120000_123_part1.parquet',
...)
FILE_FORMAT = (TYPE = PARQUET);
-- 2. Snowflake's COPY automatically deduplicates already-loaded files
-- via load history (default 14d). For longer retention or external
-- coordination, record the manifest fingerprint in a control table
-- and gate the COPY on it.
-- 3. Merge into the final table by primary key.
MERGE INTO orders target
USING orders_stage_<run_id> src
ON target.id = src.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT ...;
Snowflake’s per-stage COPY INTO load history gives you a built-in
“don’t load the same file twice” property within the retention window.
That is not a substitute for the manifest fingerprint check on the
client side — load history protects against accidental double-COPY,
not against loading a stale prefix from a half-finished export.
When append-only is safe
Skipping the merge step is acceptable only when all of the following hold:
- The source is immutable for the period in question (event logs, audit trails, time-partitioned analytics tables).
- The export carries a stable primary key that downstream consumers can use to deduplicate at query time.
- The downstream table is partitioned by event date so duplicate
rows from a re-run land in the same partition and a one-time
DELETE WHERE _rivet_run_id NOT IN (current_run_id)clean-up is cheap.
For mutable upstream tables (orders, users, accounts), always use the merge pattern. At-least-once + mutable source + append-only loader = silent data corruption that is invisible until a downstream join breaks.
What Rivet does not do downstream
- Load targets BigQuery, Snowflake and ClickHouse only.
rivet loadcovers those three (see the built-in path above); for Redshift / Trino / Databricks / dbt the operator wires up the load with the manual pattern here. - No transactional coordination. Rivet does not coordinate with a
downstream
MERGE/COMMIT. If the export run succeeds and the warehouse load fails, the operator is responsible for retry logic. - No dead-letter queue for poisoned parts. A part that fails warehouse parse (e.g. a Parquet version bump on the loader’s side) is the loader’s problem; Rivet’s manifest still says the part is committed.
See also
- docs/cloud-destinations.md — the common output contract and the per-backend support matrix.
- docs/semantics.md § Known non-guarantees — what Rivet explicitly does not promise.
- ADR-0012 — Cloud manifest contract (M1–M9) — the formal invariants this recipe builds on.
- docs/recipes/recover-interrupted-run.md — what to do before you load downstream when the export itself was interrupted.