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 thetable: schema.nameshortcut (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.
| Strategy | YAML | How boundaries are computed | When to use | Mutually exclusive with |
|---|---|---|---|---|
| Fixed size (default) | chunk_size: 100000 | SELECT MIN, MAX → [min..min+N), [min+N..min+2N), … — N rows of range per chunk | Dense numeric PK, predictable size budget per chunk | chunk_size_memory_mb |
| Fixed count | chunk_count: 16 | Range divided into exactly N equal slices; per-chunk size derived dynamically | You want exactly N workers / files (e.g. = CPU cores) | chunk_by_days |
| Date-native | chunk_by_days: 365 | chunk_column must be DATE / TIMESTAMP / TIMESTAMPTZ; windows of N days with >= AND < (open-end) semantics | Time-series, event logs, historical backfills by period | chunk_count |
| Memory-target | chunk_size_memory_mb: 256 | Auto-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 Server | You want to budget by megabytes, not rows; wide tables where row-width is hard to guess | explicit chunk_size |
| Keyset (seek) | chunk_by_key: uid | Pages 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 file | MySQL tables with no single-integer PK (UUID / string / composite PK) — the only bounded shape without a server cursor. See Keyset pagination below | chunk_column, chunk_by_days, chunk_count |
Orthogonal options that combine with any strategy:
| Field | Effect |
|---|---|
parallel: N | Up 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: true | Per-chunk row in state DB → after a crash, the next run (plain or --resume) skips completed chunks |
chunk_max_attempts: 3 | Requires 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 atparallel: 2(2→4buys almost nothing for +60 % RSS). Each worker holds its own chunk buffer, so RSS grows roughly linearly withN— raise it for narrow tables, keep it at 2 for wide ones. Theinitheuristic 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
- Rivet queries
SELECT MIN(id), MAX(id) FROM ordersto determine the range - Splits into chunks:
[min..min+chunk_size),[min+chunk_size..min+2*chunk_size), … - Each chunk runs independently:
SELECT ... WHERE id BETWEEN <lo> AND <hi>(inclusive, hi = lo + chunk_size - 1) - With
parallel: 4, up to 4 chunks execute concurrently - 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.
| Invocation | Behaviour |
|---|---|
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 size | Suggested chunk_size | parallel |
|---|---|---|
| 1M rows | 100,000 | 2 |
| 10M rows | 100,000 | 4 |
| 100M+ rows | 200,000-500,000 | 4-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_densewas removed. It paged byROW_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 setschunk_dense: trueis 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 nochunkedmode):source.mongo.page_sizeenables keyset (seek) paging on_id,parallel: Nfans out over disjoint_idranges, 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: 4on an indexed table, little beyond 4. It is tuned for indexed tables up to ~10 M rows. Past that, the range-boundary sampler (an indexOFFSETskip to find each percentile cut) grows costly at setup — for very large tables prefer a rangechunk_column(integer-PK), which slicesMIN..MAXarithmetically 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 andpartition_byare rejected for ALL keyset exports, sequential included: a decimal key is refused at plan time because the keyset cursor cannot read/advance past it, andpartition_byis incompatible withchunk_by_keyat 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.