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

Chunked Export Mode

When to use

Use mode: chunked when the table is too large for a single full export. Rivet splits the data into ranges by a numeric ID column and can process multiple chunks in parallel. Best for:

  • Tables with millions or billions of rows
  • Tables with a numeric primary key (BIGINT, SERIAL)
  • When you need parallel extraction to save time
  • Initial loads of large tables

Required fields

  • chunk_column – a numeric, date, or timestamp column to partition by (typically the primary key). Auto-resolved from the single-integer primary key if you use the table: schema.name shortcut (works on Postgres, MySQL, and SQL Server) — warn-level log line: export 'orders_chunked': chunk_column not set — auto-resolved to 'id' from the single-integer primary key on public.orders. Set `chunk_column:` explicitly to pin the choice and silence this warning.

Chunking strategies — pick one

Four ways to slice the table. They differ in how chunk boundaries are computed; everything below the strategy line (parallel, checkpoint, retry) is orthogonal and combines with any of them.

StrategyYAMLHow boundaries are computedWhen to useMutually exclusive with
Fixed size (default)chunk_size: 100000SELECT MIN, MAX → [min..min+N), [min+N..min+2N), … — N rows of range per chunkDense numeric PK, predictable size budget per chunkchunk_size_memory_mb
Fixed countchunk_count: 16Range divided into exactly N equal slices; per-chunk size derived dynamicallyYou want exactly N workers / files (e.g. = CPU cores)chunk_by_days
Date-nativechunk_by_days: 365chunk_column must be DATE / TIMESTAMP / TIMESTAMPTZ; windows of N days with >= AND < (open-end) semanticsTime-series, event logs, historical backfills by periodchunk_count
Memory-targetchunk_size_memory_mb: 256Auto-computes chunk_size from the engine’s row-size estimate (PG pg_class/reltuples, MySQL information_schema avg row length; SQL Server has no estimate and falls back to 512 B/row); clamped to [10_000, 5_000_000] rows. Requires table: shortcut. Works on Postgres, MySQL, and SQL ServerYou want to budget by megabytes, not rows; wide tables where row-width is hard to guessexplicit chunk_size
Keyset (seek)chunk_by_key: uidPages with WHERE key > last ORDER BY key LIMIT chunk_size on a unique index — sequential by default; parallel: N fans it into N disjoint key ranges — each page is one part fileMySQL tables with no single-integer PK (UUID / string / composite PK) — the only bounded shape without a server cursor. See Keyset pagination belowchunk_column, chunk_by_days, chunk_count

Orthogonal options that combine with any strategy:

FieldEffect
parallel: NUp to N chunks execute concurrently (separate DB connections). Default 1. rivet init scaffolds a row-scaled value (≤500 K → 1, <5 M → 2, ≥5 M → 4)
chunk_checkpoint: truePer-chunk row in state DB → after a crash, the next run (plain or --resume) skips completed chunks
chunk_max_attempts: 3Requires chunk_checkpoint: true. Total attempt budget per chunk (first attempt + retries): 3 means each failed chunk is retried up to 2 times before the run bails. The budget is stored on the checkpoint run and enforced when a chunk task is claimed, so without chunk_checkpoint it has no effect — the non-checkpointed runners have no per-chunk retry and a failed chunk fails the run. Defaults to tuning.max_retries + 1

Picking parallel. Extraction is I/O-bound, so the win comes from overlapping FETCH round-trips, not from CPU. Measured on a 10-core host, 2 M-row tables: a narrow table scales near-linearly (1→4 ≈ 4.2× faster), while a wide table (many large columns) saturates the wire early and plateaus at parallel: 2 (2→4 buys almost nothing for +60 % RSS). Each worker holds its own chunk buffer, so RSS grows roughly linearly with N — raise it for narrow tables, keep it at 2 for wide ones. The init heuristic is a good starting point; tune from there if memory or source connection count is constrained.

Each integer-range chunk runs: SELECT * FROM (<base_query>) AS _rivet WHERE <chunk_column> BETWEEN <lo> AND <hi> — inclusive bounds (hi = lo + chunk_size - 1) inlined as literals, not bind parameters. The subquery wrap applies to query: exports; a table: shortcut renders the unwrapped SELECT * FROM <table> WHERE <chunk_column> BETWEEN <lo> AND <hi>. The date variant (chunk_by_days) uses half-open WHERE col >= '<start>' AND col < '<end>'.

Minimal config

source:
  type: postgres
  url: "postgresql://user:pass@host:5432/dbname"

exports:
  - name: orders_chunked
    query: "SELECT id, user_id, product, price, status, ordered_at FROM orders"
    mode: chunked
    chunk_column: id                # numeric column to split ranges on
    chunk_size: 100000              # rows per chunk (default: 100,000)
    parallel: 4                     # concurrent chunk workers
    format: parquet
    destination:
      type: local
      path: ./output

Output files: one per chunk, named {export}_{YYYYMMDD_HHMMSS}_chunk{N}_{16-hex-nonce}.parquet, e.g. orders_chunked_20260406_120000_chunk0_a1b2c3d4e5f60718.parquet. The random nonce makes retried/re-run parts additive (never overwriting); match on *_chunk{N}_*.parquet, not on an exact stem.

Run it

# Preflight — shows chunk plan (how many chunks, range distribution)
rivet check --config large_table.yaml

# Run with validation and reconciliation
rivet run --config large_table.yaml --validate --reconcile

Progress bar (chunked exports)

In mode: chunked, Rivet shows a terminal progress bar while chunks run: export name, current/total chunks, running row count, elapsed time, and ETA. It appears when stderr is an interactive TTY (a normal terminal window). The bar does not depend on RUST_LOG (that variable only controls env_logger text lines). In CI, or when you pipe or redirect stderr, the bar is usually suppressed — then set RUST_LOG=info (or debug) to follow progress in the log instead.

The GIF above was recorded with RUST_LOG=info on a 50,000-row fixture (10 chunks of 5,000) so the per-chunk log line export 'events': chunk N/10 (...) and the final summary are both visible. On a real interactive terminal you would see the progress bar instead; the log lines appear when stderr is captured.

Use a small chunk_size relative to your table if you want many steps on the bar (each finished chunk advances it once). parallel: 1 still updates the bar after each sequential chunk.

Ready-made example in this repo: dev/scenarios/chunked_postgres_bench.yaml includes bench_content_p4_safe: PostgreSQL content_items with parallel: 4 and tuning.profile: safe (good for trying the bar on a wide table without hammering the source). Other exports in the same file cover serial / highly parallel / fatchunk / balanced profiles.

# From repo root; Postgres up + seeded (e.g. docker compose + cargo run --bin seed ...)
mkdir -p dev/output/bench
rivet check --config dev/scenarios/chunked_postgres_bench.yaml
rivet run --config dev/scenarios/chunked_postgres_bench.yaml --export bench_content_p4_safe
# Optional: RUST_LOG=info for more log detail; RUST_LOG=warn to reduce log noise (bar unchanged in a TTY)

What happens

  1. Rivet queries SELECT MIN(id), MAX(id) FROM orders to determine the range
  2. Splits into chunks: [min..min+chunk_size), [min+chunk_size..min+2*chunk_size), …
  3. Each chunk runs independently: SELECT ... WHERE id BETWEEN <lo> AND <hi> (inclusive, hi = lo + chunk_size - 1)
  4. With parallel: 4, up to 4 chunks execute concurrently
  5. Each chunk writes a separate output file

Chunk checkpoint (resume after crash)

For very large exports, enable checkpointing so you can resume from where you left off:

exports:
  - name: orders_chunked
    query: "SELECT id, user_id, product, price, ordered_at FROM orders"
    mode: chunked
    chunk_column: id
    chunk_size: 100000
    parallel: 4
    chunk_checkpoint: true          # persist progress per chunk
    chunk_max_attempts: 3           # total attempt budget per chunk (3 attempts = 2 retries)
    format: parquet
    destination:
      type: local
      path: ./output

Resume after a crash:

# Resume only processes incomplete chunks
rivet run --config large_table.yaml --resume

# View checkpoint status
rivet state chunks --config large_table.yaml --export orders_chunked

# Clear checkpoint (to re-export from scratch)
rivet state reset-chunks --config large_table.yaml --export orders_chunked

Clean re-runs are NOT idempotent

Chunked mode is not “extract once, skip on the next clean run”. Two plain rivet run invocations against the same table re-extract every chunk both times — chunk_checkpoint: true only matters after a crashed run, which the next run resumes. Each clean run produces a new file set with a fresh run_id and timestamp suffix.

InvocationBehaviour
rivet run (fresh)extracts all chunks, writes files with run_id A
rivet run (again, no crash)extracts all chunks again, writes files with run_id B
rivet run --resume (after a crash)extracts only the chunks chunk_state says are incomplete

If you want skip-on-no-change semantics, use mode: incremental with a cursor_column instead — that mode persists the cursor between runs, and skip_empty: true records a run that found nothing new as skipped.

Chunk sizing guidance

Table sizeSuggested chunk_sizeparallel
1M rows100,0002
10M rows100,0004
100M+ rows200,000-500,0004-8

Larger chunks = fewer queries but more memory per batch. Smaller chunks = more queries but lower peak RSS.

Date-based chunking

When your table’s natural partition boundary is time rather than a numeric ID, use chunk_by_days instead of relying on integer ranges.

exports:
  - name: orders_by_year
    query: "SELECT id, user_id, product, price, ordered_at FROM orders"
    mode: chunked
    chunk_column: ordered_at        # DATE or TIMESTAMP column
    chunk_by_days: 365              # one chunk per ~year
    format: parquet
    destination:
      type: local
      path: ./output

Rivet fetches MIN / MAX of the column as text, parses the dates, then generates non-overlapping windows:

-- each chunk window (open-end exclusive):
WHERE ordered_at >= '2023-01-01' AND ordered_at < '2024-01-01'
WHERE ordered_at >= '2024-01-01' AND ordered_at < '2025-01-01'
...

The open-end < end_date bound is intentional: it correctly captures all TIMESTAMP values within the day, including 23:59:59.999….

When to use date chunking over numeric chunking:

  • The table has no dense numeric PK (UUIDs, composite keys)
  • You want even partitions by time, not by row count
  • The source DB has better statistics / indexes on the timestamp column
  • You want to avoid unix-epoch arithmetic that JDBC tools often get wrong

chunk_by_days can be combined with parallel for concurrent date windows, and supports chunk_checkpoint / --resume like numeric chunked mode.

rivet check will report the strategy as date-chunked(ordered_at, 365d).

Sparse ID ranges

If IDs have large gaps (e.g. UUIDs cast to BIGINT, or deleted rows), many chunks may be empty. On a unique key use keyset (chunk_by_key, below) — it pages by rows, so gaps cost nothing. Otherwise chunk_count: N caps the number of windows.

rivet check will warn you about sparse ranges.

chunk_dense was removed. It paged by ROW_NUMBER() OVER (ORDER BY chunk_column), recomputed per chunk, so concurrent inserts or deletes skipped or duplicated rows even on a unique key. A config that still sets chunk_dense: true is refused at load.

Keyset (seek) pagination — the safe shape without an integer PK

Range chunking needs a single integer PK to slice MIN..MAX. A MySQL table whose PK is a UUID, string, or composite key has no such column, and — unlike PostgreSQL — MySQL has no server-side cursor to bound a mode: full snapshot. That left a real hole: such tables could only be exported as one long-held SELECT * (the exact “don’t hold a long query on prod” risk Rivet exists to avoid).

MongoDB has the same shape, but under mode: full (a document store has no chunked mode): source.mongo.page_size enables keyset (seek) paging on _id, parallel: N fans out over disjoint _id ranges, and both resume on _id. See ../reference/mongodb.md.

Keyset pagination closes it. Rivet pages the table by a unique, NOT NULL, index-backed key:

-- first page
SELECT * FROM (<base>) AS _rivet ORDER BY `uid` LIMIT 1000
-- subsequent pages (cursor = last page's max key)
SELECT * FROM (<base>) AS _rivet WHERE `uid` > ? ORDER BY `uid` LIMIT 1000

Each page is a bounded, index range scan (verified EXPLAIN: type: range on the PK, no filesort, no full scan) and becomes one part file. This bounds both peak RSS (≤ chunk_size rows in flight) and longest-query time.

exports:
  - name: events
    table: app.events            # `table:` shortcut required (index check)
    mode: chunked
    chunk_by_key: event_uuid     # single-column UNIQUE / PRIMARY, NOT NULL
    chunk_size: 1000             # rows per page
    format: parquet
    destination:
      type: local
      path: ./output

Output files: one per page, named {export}_{run_id}_keyset_{tag}.parquet where run_id is {export}_YYYYMMDDTHHMMSS.mmm (filename sanitization maps the . to _) and the tag is start for the first page, then a 16-hex hash of that page’s seek cursor — e.g. events_events_20260529T120000_123_keyset_start.parquet. The run_id/seek-based name makes a crash-resume overwrite its own page idempotently.

Auto-resolution (MySQL). With the table: shortcut and no chunk_by_key, if the table has no single-integer PK but does have a usable single-column unique key, Rivet auto-selects keyset on it and logs a warn naming the key (set chunk_by_key: to pin the choice and silence the warning). On PostgreSQL, auto-resolution stays off — its DECLARE CURSOR snapshot is already bounded, so mode: full is the safe answer there; chunk_by_key: still works if you want per-page files.

The key must be index-backed. This is the load-bearing safety property: an ORDER BY on a non-indexed column degrades to a full-scan + filesort — worse than the snapshot it replaces. Rivet refuses a chunk_by_key that is not a single-column, NOT NULL, UNIQUE/PRIMARY key rather than emit such a query:

chunk_by_key 'payload' is not a usable keyset key on app.events — it must be a
single-column, NOT NULL, UNIQUE or PRIMARY key WHOSE TYPE the keyset cursor can
read (integer / float / string / timestamp / date / uuid). A `decimal`/`numeric`
key is excluded: the cursor cannot advance past it (it would fail mid-run after
a partial write). Without a usable key, `ORDER BY payload LIMIT n` would also
full-scan + filesort. Add a unique index of a supported type, pick another key,
use a range `chunk_column:` (integer), or `mode: full`.

Required privileges: read-only is sufficient — the introspection probe reads information_schema index metadata, no elevated grants needed.

Resumability (chunk_checkpoint). rivet init defaults chunk_checkpoint: true on keyset exports. It is crash-recovery: a run that dies mid-stream resumes from its last committed key on the next run (its in-progress run_id is still open); a run that finished CLEANLY clears that marker, so a plain re-run does a full pass and never silently skips already-exported rows.

Append-only incremental (keyset_incremental). Off by default. When set, a CLEAN re-run continues from the last exported key — pulling ONLY rows past the high-water mark. Correct only for append-only tables: on a mutable table a row whose key already passed is silently never re-read. For a mutable table use mode: incremental on a timestamp cursor instead.

    chunk_by_key: event_uuid
    chunk_checkpoint: true       # crash-recovery (default on for keyset)
    keyset_incremental: true     # append-only ONLY: clean re-run pulls just new keys

Parallel keyset (parallel: N)

By default keyset pages sequentially — each page seeks past the previous page’s last key. Set parallel: N and Rivet instead splits the key into N ROW-based percentile ranges and seeks each range concurrently on its own connection:

    chunk_by_key: event_uuid
    chunk_size: 1000
    parallel: 4                  # 4 workers, each seeks a disjoint key range

The ranges are half-open ([lo, hi)) and adjacent, so together they partition the key — every row is read exactly once (structural parity, proven on all engines). Each worker still pages by seek within its range, so peak RSS stays bounded (≤ N × chunk_size rows in flight) and no worker holds a long query. Per-range crash-recovery is tracked in the state DB (chunk_checkpoint): a run that dies resumes only the unfinished ranges.

Sweet spot. Extraction is I/O-bound, so the speedup plateaus early — ~3.1× at parallel: 4 on an indexed table, little beyond 4. It is tuned for indexed tables up to ~10 M rows. Past that, the range-boundary sampler (an index OFFSET skip to find each percentile cut) grows costly at setup — for very large tables prefer a range chunk_column (integer-PK), which slices MIN..MAX arithmetically with no sampling pass.

Canary first. Parallel keyset is new in 0.23.0. Run it on a canary table and diff row counts against a sequential (parallel: 1) pass before rolling it out across a fleet — the sequential path is unchanged and remains the conservative default.

rivet init scaffolds a row-scaled parallel (≤500 K → 1, <5 M → 2, ≥5 M → 4) on range chunk_column tables (no single-column PK). A keyset (chunk_by_key) table is scaffolded sequential — add parallel: N yourself to opt into parallel keyset. A preflight warns past ~5 M rows that peak RSS scales with N.

Limitations (current):

  • Single-column keys only — composite unique keys are not yet supported.
  • Decimal (numeric) keys and partition_by are rejected for ALL keyset exports, sequential included: a decimal key is refused at plan time because the keyset cursor cannot read/advance past it, and partition_by is incompatible with chunk_by_key at config validation. (Parallel keyset adds no extra key-type restriction beyond these.)

Troubleshooting

Many empty chunks – Your ID column has gaps. Use chunk_by_key on a unique key, or chunk_count: N.

High memory usage with parallel > 1 – Reduce chunk_size or add tuning.profile: safe.

Export fails midway through 1000 chunks – Enable chunk_checkpoint: true; the next run resumes from the completed chunks, whether the last run crashed or failed.