Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

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}.parquet for a single-part run, with a _part0.._partN-1 suffix 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.json listing every committed part with size_bytes and content_fingerprint (xxh3 over the part bytes),
  • a _SUCCESS marker whose body fingerprints the exact manifest.json bytes.

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 manifest row_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: full OVERWRITEs. 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: cdc append + dedup. For mutable sources the load appends to <table>__changes and exposes a current-state view (latest-per-PK, deletes flagged) — the built-in equivalent of the manual MERGE below, no staging table or upsert SQL to write. On ClickHouse the log is a ReplacingMergeTree read through a FINAL view (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:

  1. 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).
  2. 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:

  1. Read manifest.json from the resolved destination prefix.
  2. Verify _SUCCESS matches. If it does not, abort — the export is not complete.
  3. Load only the parts listed in manifest.json with status: committed. Do not glob the prefix, and skip superseded / quarantined entries.
  4. Record the manifest identity (run_id + schema_fingerprint + _SUCCESS body) in a downstream control table.
  5. Deduplicate by primary key (or natural key) when merging into the final table.
  6. 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 / uuid load as BYTES, a naive timestamp as TIMESTAMP (an instant, not wall-clock DATETIME), and arrays as a nested RECORD. BigQuery will not coerce these on load — declaring native types in a load schema is rejected — so recover them with a post-load CREATE TABLE … AS SELECT over the staging table. Load the staging table with --parquet_enable_list_inference (so arrays flatten with UNNEST), then run the recovery SQL that rivet check --type-report --target bigquery prints per export. Full table: type-mapping.md § BigQuery autoload & recovery.

Wide decimals (NUMERIC/DECIMAL precision > 38). The load’s default decimal target is NUMERIC (max precision 38, scale 9) and fails on wider columns; pass --decimal_target_types=NUMERIC,BIGNUMERIC,STRING (repeated flag in bq: one value per flag). BIGNUMERIC still caps at ~5.79×10³⁸ — 38 integer digits — so a DECIMAL(50,10) column loads but a value with 39+ integer digits does not fit ANY BigQuery numeric type; with STRING in 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 garbage DOUBLE — 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 load covers 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