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

Tuning Reference

Tuning controls how Rivet queries the source database: batch sizes, timeouts, throttling, and retries.

MongoDB sources don’t use the SQL tuning on this page — a document store has no chunked mode or chunk_size, though tuning.batch_size (per-batch row cap) and max_batch_memory_mb (per-batch byte cap) ARE honored. Mongo’s tuning levers are the driver connection pool and parallel: N _id-range fan-out; see MongoDB → Connection pool & parallel tuning.

Where to place tuning

Tuning can be set at two levels:

  1. Global (source.tuning) – applies to all exports
  2. Per-export (exports[].tuning) – overrides global for that export

Per-export values take precedence. Unset per-export fields fall back to the global value.

source:
  type: postgres
  url_env: DATABASE_URL
  tuning:
    profile: balanced               # global default
    batch_size: 10000

exports:
  - name: small_table
    query: "SELECT * FROM users"
    format: parquet
    destination: { type: local, path: ./out }
    # inherits global tuning (balanced, batch_size=10000)

  - name: huge_table
    query: "SELECT * FROM events"
    format: parquet
    destination: { type: local, path: ./out }
    tuning:
      profile: safe                 # override for this export only
      batch_size: 2000

Common mistake: placing batch_size directly under source: or in the export root instead of under tuning:. Rivet will reject such configs with a clear error message.

Profiles

A profile sets sensible defaults for all tuning parameters. Individual fields override the profile.

Parameterfastbalanced (default)safe
batch_sizeadaptive: 64 MB/flush¹adaptive: 32 MB/flush¹2,000 (static)
throttle_ms050500
statement_timeout_s0 (none)300120
max_retries1310
retry_backoff_ms1,0002,0005,000
lock_timeout_s0 (none)3010
memory_threshold_mb0 (none)4,0962,048

¹ fast and balanced size the batch from memory, not a row count: batch = target_mb / estimated_row_bytes, clamped to 1,000–150,000 rows. A ~320-byte row therefore batches at ~100k rows under balanced; a 4 KB row at ~8k. The static bases (50,000 / 10,000) apply only when the schema is not yet known (e.g. plan before resolve) and as the advisory base in reports. An explicit batch_size: disables adaptive sizing. Either way the batch is a CLIENT-side fetch window — it never enters the source SQL (page size is chunk_size), so a larger batch shortens cursor hold-time without adding server work.

What stays open on the server while a chunk drains (per-engine hold model): PostgreSQL reads through a server cursor — between FETCH N calls nothing executes, but the snapshot transaction stays open (vacuum-horizon cost); MySQL and SQL Server hold one streaming SELECT per chunk, drained under socket flow control — the query stays visible (Sending data) for the chunk’s drain duration and holds the MVCC read view. Consequence: throttle_ms lowers burst IO but lengthens per-chunk hold time on MySQL/MSSQL and the snapshot window on PostgreSQL. For hold-time-sensitive primaries prefer fast in an off-peak window or a replica; reserve safe throttling for replicas and IO-sensitive hosts. chunk_size bounds the worst case held by any single query.

When to use each profile

ProfileUse case
fastDedicated read replica, off-peak hours, small tables
balancedGeneral purpose, shared database, production reads
safeBusy production database, OLTP systems, wide tables with large rows

All tuning parameters

FieldTypeDefaultDescription
profilefast | balanced | safebalancedBase profile (sets defaults for all other fields)
batch_sizeintegerprofile defaultRows fetched per query batch. Explicit value disables the profile’s adaptive (memory-based) sizing — see footnote ¹ above
batch_size_memory_mbinteger—Target memory per batch in MB (adaptive sizing; mutually exclusive with batch_size)
throttle_msintegerprofile defaultDelay in ms between batches (reduces source load)
statement_timeout_sintegerprofile defaultDatabase statement timeout in seconds (0 = no timeout)
max_retriesintegerprofile defaultMax retry attempts for transient errors
retry_backoff_msintegerprofile defaultBase delay between retries in ms (exponential backoff)
lock_timeout_sintegerprofile defaultDatabase lock timeout in seconds (0 = no timeout)
memory_threshold_mbintegerprofile defaultRSS threshold in MB; pauses fetching if exceeded (0 = disabled). balanced defaults to 4096, safe to 2048, fast to 0 (no limit).
max_batch_memory_mbinteger—Hard cap on a single Arrow batch in MB. When exceeded, on_batch_memory_exceeded determines the response.
on_batch_memory_exceededwarn | fail | auto_shrinkwarnPolicy applied when a batch exceeds max_batch_memory_mb.
max_value_mbinteger256Hard ceiling on a single cell (text/JSON/blob) in MB. A value larger than this aborts the run with RIVET_VALUE_TOO_LARGE. Guards against one giant cell OOM-ing the process — the batch cap is average-based and can’t bound a lone outlier. Set 0 to disable. See Per-value ceiling.
adaptivebooleanfalseSample source write-pressure at runtime and react: shrink/restore the fetch batch size, and — on a parallel > 1 chunked or keyset export — drive the concurrency governor (see that section for the exact per-runner coverage).
min_parallelinteger1Floor for the concurrency governor: the fewest workers it will back down to under pressure. Ceiling is the export’s parallel. Only consulted when adaptive is on and parallel > 1.

Batch memory cap (max_batch_memory_mb)

memory_threshold_mb is a process-level RSS guard — it fires after the OS has already committed memory. max_batch_memory_mb is an earlier, batch-level guard: it measures the actual Arrow buffer footprint of each batch before it is written.

tuning:
  max_batch_memory_mb: 128
  on_batch_memory_exceeded: warn   # warn | fail | auto_shrink
PolicyBehaviour
warn(default) Log a warning with the actual size, the limit, and a suggested batch_size. Continue the export.
failReturn an error immediately. The export stops. Use in strict pipelines where oversized batches indicate a configuration problem.
auto_shrinkSplit the oversized batch in half recursively until each sub-batch fits within the limit, then write the sub-batches individually. Transparent to the rest of the pipeline — total row count and output are identical.

The warning and error messages include a suggested batch_size:

batch memory 184 MB exceeds max_batch_memory_mb=128 MB (5000 rows).
Consider lowering batch_size to ~3478.

Use auto_shrink when you want protection against accidental wide-table OOM without needing to tune batch_size manually. Use fail in CI pipelines where any oversized batch should block the run.

Per-value ceiling (max_value_mb)

max_batch_memory_mb and the adaptive byte budget are average-based — they size a batch from its mean row width. Neither bounds a single pathological cell: one 300 MB JSONB document or bytea blob among otherwise-small rows still lands whole in memory and can OOM the process (and the auto_shrink splitter can’t divide a single oversized value).

max_value_mb is a hard per-value ceiling. Before a batch is split or encoded, Rivet checks every variable-length cell (text / JSON / binary — fixed-width types can’t be individually huge); a value over the limit aborts the run:

RIVET_VALUE_TOO_LARGE: column 'body' has a single value of 301.2 MB, exceeding the
per-value ceiling of 256 MB. ...Raise `tuning.max_value_mb` (or set it to 0 to
disable the guard) if this value is expected.

It is on by default at 256 MB — high enough to never trip on realistic data, low enough to catch a runaway cell before it OOMs. Raise it for tables that legitimately store large blobs, or set max_value_mb: 0 to disable the guard entirely.

Choosing batch_size

batch_size is the most impactful parameter for both performance and memory usage.

batch_sizeMemory per batch (narrow table)Memory per batch (wide table)Best for
1,000~1-5 MB~20-100 MBWide tables, low-memory environments
5,000~5-25 MB~100-500 MBMedium tables, shared databases
10,000~10-50 MB~200 MB - 1 GBGeneral purpose (default balanced)
50,000~50-250 MB~1-5 GBRead replicas, fast profile

For wide tables (50+ columns, TEXT/JSONB fields), start with batch_size: 1000-2000.

Adaptive batch sizing

Instead of a fixed row count, let Rivet adjust batch size based on memory:

tuning:
  batch_size_memory_mb: 64          # target ~64 MB per batch

Rivet samples the first batch to estimate row size, then adjusts subsequent batches. Cannot be used together with batch_size.

Choosing chunk_size (and bounding statement duration)

chunk_size is a different lever from batch_size. batch_size is internal — how many rows Rivet buffers in Arrow memory at a time (RSS only). chunk_size is the unit of work and output: in chunked mode it is the size of one WHERE key BETWEEN … (or keyset … LIMIT n) window, which is one SQL statement and one output part file.

That makes chunk_size the knob for the longest single query the source sees — the thing a DBA’s statement_timeout, long-running-query alert, or lock-duration monitor reacts to. On a wide table, one chunk statement transfers chunk_size × row_width bytes and stays active on the server for that whole duration. Measured on MySQL content_items (wide ~4 KB rows; one chunk statement, wall):

chunk_sizeone chunk statementoutput files (for 1 M rows)
1,000~0.4 s1,000 small files
10,000~0.6 s100 files
100,000 (default)~3.4 s (≈9 s at ~12 KB rows)10 large files

If a strict statement_timeout on the source trips your chunk queries, or you want to keep each read short and gentle on a busy OLTP source, lower chunk_size (e.g. chunk_size: 10000):

exports:
  - name: orders
    mode: chunked
    chunk_column: id
    chunk_size: 10000     # ~0.5 s per statement instead of ~3-9 s

The trade-off is more, smaller part files and a small (~25%) increase in total wall time (more index seeks / round-trips for the same rows). It is not a throughput win — it trades total speed and query count for shorter individual statements. Pick the point that fits your source’s tolerance.

PostgreSQL is unaffected: it streams each chunk through a server-side cursor (DECLARE … FETCH N, N capped by work_mem), so its per-statement work is already bounded regardless of chunk_size. The lever above matters for MySQL / SQL Server, which run one statement per chunk.

Why not give MySQL the same server-side cursor? Because its read-only cursor works differently: it materialises the whole result into temp tables when the cursor opens, then fetches cheaply — the open itself is the long statement, and it adds tempdb pressure. Measured directly with a libmysqlclient probe: cursor-open 0.8–1.8 s and 3 temp tables created, every run (dev/spikes/mysql_cursor_efficacy.c). So a MySQL cursor would be worse than just lowering chunk_size (short pages, no temp tables). Lowering chunk_size is the right lever; there is no free server-cursor shortcut on MySQL.

Adaptive concurrency governor

On an export with parallel > 1, setting adaptive: true arms a governor that adjusts how many workers (and therefore source connections) run concurrently, in response to source write-pressure. It backs parallelism down when the source is under load and recovers it when the load eases, staying within [min_parallel, parallel].

Which runners it covers, precisely — the governor is per-runner wiring, so this list is the contract, not an approximation:

RunnerGoverned?
mode: chunked, parallel > 1yes — sheds at chunk granularity
mode: chunked + chunk_checkpoint: true, parallel > 1 (the shape rivet init scaffolds)yes — sheds at claimed-task granularity
chunk_by_key (keyset), parallel > 1yes — sheds at page granularity
MongoDB parallel: N (_id-range fan-out)no — that runner has no shared permit ceiling to shrink. adaptive still drives Mongo’s batch-size adaptation; it does not vary worker count.
Anything with parallel: 1 (or unset)no — one worker has nothing to shed.
source:
  type: postgres
  url_env: DATABASE_URL
  tuning:
    adaptive: true        # arm batch-size adaptation + the governor
    min_parallel: 2       # never drop below 2 workers (default 1)

exports:
  - name: orders
    table: public.orders
    mode: chunked
    chunk_column: id
    parallel: 8           # ceiling — governor varies the live count in [2, 8]
    format: parquet
    destination: { type: local, path: ./out }

How it decides. A dedicated monitoring connection polls a source write-pressure counter every ~1.5 s and compares it to the previous reading. A rising counter means pressure is climbing, so the governor sheds one worker; a flat/falling counter lets it recover one. The counter is:

EngineGovernor pressure proxyRead via
PostgreSQL (< 17)pg_stat_bgwriter.checkpoints_reqSELECT checkpoints_req FROM pg_stat_bgwriter
PostgreSQL (17+)pg_stat_checkpointer.num_requestedSELECT num_requested FROM pg_stat_checkpointer
MySQLglobal Innodb_log_waitsSHOW GLOBAL STATUS LIKE 'Innodb_log_waits'
SQL ServerLog Flush Waits/sec, cumulative cntr_value of the _Total rowSELECT cntr_value FROM sys.dm_os_performance_counters WHERE counter_name LIKE 'Log Flush Waits%' AND instance_name = '_Total'

Two per-engine details worth knowing if you correlate rivet’s decisions with your own monitoring:

  • PostgreSQL 17 moved the counter. pg_stat_bgwriter.checkpoints_req was removed in PG 17 and lives on as pg_stat_checkpointer.num_requested. rivet picks the right one at runtime with an existence probe (SELECT to_regclass('pg_catalog.pg_stat_checkpointer') IS NOT NULL) rather than a single CASE statement, because PG plans the whole statement up front — a dead branch referencing the missing column still ERRORs, and that error would abort the export’s own cursor transaction. Each sample is preceded by pg_stat_clear_snapshot().
  • SQL Server reads the _Total row, it does not SUM. The SQLServer:Databases object exposes one row per database plus a _Total row (verified live: _Total equals the sum of the others), so a SUM(cntr_value) over all rows double-counts — and it shrinks when a database is dropped, which the governor would read as “pressure eased”. Use the _Total-filtered single-row read above if you want the series rivet actually sees.

The governor’s proxy is deliberately NOT the adaptive batch loop’s. On MySQL the batch loop listens to own-extraction pressure (spill/temp counters the export’s own reads inflate — shrinking the batch genuinely shrinks the per-query spill); on PostgreSQL the batch loop shares checkpoints_req with the governor; SQL Server’s batch loop currently has no pressure sampling (batch adaptation is inert there). The governor asks a different question — is someone ELSE straining this server while I run? — so it listens to write/redo counters a read-only export cannot move. Feeding it the batch loop’s spill counters makes it read its own exhaust: a keyset export whose pages spill by design would shed workers 4→3→2→1 and never recover (the counter keeps rising as long as its own pages run) — measured on a production pool run as every keyset export slowing 2–2.7×.

Required privileges (read-only is enough)

The governor needs no elevated privileges. A plain read-only role can run every query it issues — verified against PostgreSQL 16 and MySQL 8:

  • PostgreSQL — a role with only CONNECT + USAGE ON SCHEMA + SELECT ON TABLES can read pg_stat_bgwriter (< PG 17) or pg_stat_checkpointer (PG 17+), run the to_regclass probe that chooses between them, and call pg_stat_clear_snapshot() — all are available to PUBLIC. No pg_read_all_stats, no superuser.

    CREATE ROLE rivet_ro LOGIN PASSWORD '…';
    GRANT CONNECT ON DATABASE mydb TO rivet_ro;
    GRANT USAGE ON SCHEMA public TO rivet_ro;
    GRANT SELECT ON ALL TABLES IN SCHEMA public TO rivet_ro;
    
  • MySQL — a user with only SELECT on the target schema can run SHOW GLOBAL STATUS; it needs no PROCESS or other global privilege.

    CREATE USER 'rivet_ro'@'%' IDENTIFIED BY '…';
    GRANT SELECT ON mydb.* TO 'rivet_ro'@'%';
    

Graceful degradation — a transient miss holds flat, a dead signal fails OPEN. An unreadable pressure sample (locked-down role, unsupported engine view, a statement timeout on a busy catalog view) never fails the run. It degrades in two stages:

  1. Transient miss — fewer than 3 consecutive unreadable samples hold parallelism exactly where it is, and keep the last real reading as the baseline so the next successful sample is still compared against it.
  2. Signal lost — at 3 consecutive unreadable samples (~4.5 s at the default 1.5 s interval) the governor says so once for the episode — at warn when a signal it had been reading died — and then steps parallelism back up one worker per tick until it reaches the export’s parallel ceiling. A signal that cannot be READ is not evidence of pressure, so the governor fails open rather than leaving the run pinned at whatever level the last shed reached for the rest of its hours. If you need a hard cap while blind, lower parallel (the ceiling) — min_parallel is a floor and does not bound the recovery.

A failed monitoring connection (as opposed to a failed sample) logs a warning and disables the governor entirely for that run; parallelism then stays static at parallel. Neither case aborts the export.

Note on richer signals. A future iteration may read lock waits / idle in transaction from pg_stat_activity or SHOW PROCESSLIST. Those do require elevated privileges (pg_read_all_stats on PostgreSQL; the PROCESS privilege on MySQL) to observe sessions other than your own. The current proxy was chosen specifically so the default least-privilege, read-only setup keeps working. When the richer signals land, this section will document the additional grants.

Visibility. Every adjustment is recorded in the run journal as a ParallelismAdjusted event (from, to, reason). The log level is asymmetric on purpose: a shed is a deliberate slowdown of your run, so it must be visible at the default level (an info-level “this will be slower” is functionally silent — a field pool run lost 1h48m to invisible sheds), while a recovery is good news and stays quiet.

EventLevelLine
Governor armedinfoexport 'orders': adaptive concurrency governor active (parallel 2..8)
Shedwarnexport 'orders': governor parallelism 8 → 7 (source pressure rising: backed off) — raise `min_parallel` to floor it, or set `adaptive: false` to disarm
Recoveryinfoexport 'orders': governor parallelism 7 → 8 (source pressure eased: recovered)
Armed but no signal at all (first probe)warnexport 'orders': governor armed, but the source provides no pressure signal … — parallelism stays at 8
Signal died mid-run (see above)warnexport 'orders': governor lost its pressure signal … parallelism was pinned at 3 of 8; stepping back toward 8 …
Monitoring connection failedwarnexport 'orders': governor monitoring connection failed; parallelism stays static at 8: …

The lines above are quoted to show the level and the shape; grep for governor parallelism (adjustments) and governor (everything else) rather than matching a full line, since the trailing hints get refined between releases.

Write pipelining

For single/snapshot exports (mode: full), Rivet runs the fetch+convert stage and the Parquet encode+compress stage on two threads with a small bounded channel between them, so the database round-trip wait overlaps the compression CPU. It is on by default, FIFO-ordered (byte-identical output), and free — no measurable RSS penalty at the default depth, and the commit- critical finalize still runs on the main thread.

  • Disable it (old synchronous path): RIVET_PIPELINE_WRITES=0.
  • Tune the channel depth (memory ↔ overlap): RIVET_PIPELINE_WRITES=<n>.

The gain scales with how much real work compression does: on diverse data the encoder is busy and the overlap is worth it; on trivially-compressible data (near-zero zstd work) there is little to overlap. Chunked exports already run each chunk on its own worker and are not intra-chunk pipelined.

Memory optimization tips

  1. Reduce batch_size – the single most effective knob
  2. Use safe profile for wide tables on production databases
  3. jemalloc is the default allocator – ordinary builds (cargo build, cargo install rivet-cli) already include it (a default cargo feature), so its 20-40% RSS reduction is in effect out of the box; only a --no-default-features build loses it
  4. Set memory_threshold_mb – Rivet pauses fetching when RSS exceeds this

Examples

Minimal (use defaults)

source:
  type: postgres
  url_env: DATABASE_URL
  # No tuning block → balanced profile with all defaults

Aggressive (read replica)

source:
  type: postgres
  url_env: REPLICA_URL
  tuning:
    profile: fast
    batch_size: 100000
    throttle_ms: 0

Conservative (production OLTP)

source:
  type: postgres
  url_env: DATABASE_URL
  tuning:
    profile: safe
    batch_size: 1000
    throttle_ms: 1000
    statement_timeout_s: 60
    memory_threshold_mb: 512

Capacity and memory planning

Peak RSS formula

peak_rss ≈ batch_size × avg_row_bytes × parallel_workers
         + Σ distinct destination.oneshot_budget_mb   (64 MB when unset; cloud only)

The one-shot term is per rivet process: under parallel_export_processes each child adds its own. Add ~50–150 MB overhead for the Tokio runtime, the source connection pool, jemalloc bookkeeping, and the OS page cache on the temp file.

Rule of thumb by table width

Table typeAvg row bytesRecommended batch_sizeExpected peak RSS
Narrow (IDs, timestamps, small text)~100 B50 000–100 000~50–200 MB
Medium (mixed text, JSON)~1 KB10 000–25 000~50–250 MB
Wide (TEXT/JSONB payloads ≥ 10 KB avg)~10 KB500–2 000~50–200 MB

Use the safe profile for wide tables — it uses a conservative static batch_size of 2 000 and the tightest throttle. For an explicit memory ceiling regardless of profile, set batch_size_memory_mb (memory-driven sizing) or a smaller batch_size directly.

How memory_threshold_mb works

When tuning.memory_threshold_mb is set, the chunked runners sample RSS at each chunk boundary (via mach_task_basic_info on macOS, /proc/self/statm on Linux). What happens above the threshold depends on the path: the parallel chunked runner (without checkpointing) holds the next chunk back, re-polling every 2 s until RSS falls back below the threshold; the sequential and checkpointed chunked paths pause once for a fixed 2–5 s and then proceed with the next chunk even if RSS is still above the threshold. There is no hysteresis band, and no per-batch check. Full, incremental, keyset, and mongo-parallel exports never pause on this knob; there RSS is only recorded for the peak-RSS metric. For a memory ceiling on those paths use batch_size_memory_mb / max_batch_memory_mb instead.

source:
  tuning:
    memory_threshold_mb: 1024   # pause fetching above 1 GB RSS

The RSS syscall costs ~1–2 ms, paid once per chunk. The guard is enabled by default on balanced (4096 MB) and safe (2048 MB) profiles; set memory_threshold_mb: 0 to disable it.

Parallelism and source capacity

Each parallel chunk worker opens its own source connection. Postgres max_connections is typically 100–200 for shared instances and 20–50 for read replicas. rivet check warns when parallel >= max_connections.

Safe upper bound: parallel ≤ max_connections / 4 to leave headroom for application traffic.

Per-export memory isolation

--parallel-export-processes spawns one OS process per export — each export has its own allocator and heap, so peak RSS is per-export rather than aggregate. Use this mode when running many wide-table exports at once on memory-constrained hosts.