Value-Based Output Partitioning (partition_by)
When to use
Set partition_by to split one export’s rows into one destination sub-folder
per value bucket of a column — the Hive-style col=value/ layout that
warehouses and query engines (Snowflake external tables, BigQuery external /
Hive-partitioned tables, Spark, DuckDB, Athena) discover automatically. Best for:
- Daily/monthly snapshots that land each day’s rows under its own prefix
- Backfilling history into a partitioned lake layout in one command
- Feeding a warehouse that prunes partitions by a date column
partition_by is orthogonal to mode: each partition runs the export’s own
mode, so mode: chunked chunks within a partition.
Required fields
partition_by— the column whose value buckets the rows (aDATE/TIMESTAMP/TIMESTAMPTZcolumn).- a
{partition}token indestination.pathordestination.prefix— Rivet refuses the run without it, because every partition would otherwise overwrite the same prefix.
Optional fields
partition_granularity—day(default),month, oryear.
Minimal config
source:
type: postgres
url: "postgresql://user:pass@host:5432/dbname"
exports:
- name: events
table: events
partition_by: created_at
partition_granularity: day
format: parquet
destination:
type: s3
bucket: my-bucket
prefix: "events/{partition}/" # → events/created_at=2023-01-01/
Run it
rivet run --config events.yaml --validate
What happens
- Rivet reads the
[min, max]span ofpartition_byfrom the source (SELECT min(col),SELECT max(col)over your query). - It generates one contiguous bucket per day/month/year across that span.
- Each bucket becomes its own export: the query is wrapped as
SELECT * FROM (<your query>) WHERE col >= '<lo>' AND col < '<hi>'(half-open:loinclusive,hiexclusive — no row counted twice), and{partition}resolves tocol=value. - Rows whose
partition_byvalue is NULL land incol=__HIVE_DEFAULT_PARTITION__/(the Hive default-partition convention) so no row is ever silently dropped.
Each partition is an independent, complete output prefix — its own
manifest.json and _SUCCESS — so it can be validated and consumed on its own.
Example output layout
events/
created_at=2023-01-01/ manifest.json _SUCCESS events__2023-01-01_*.parquet
created_at=2023-01-02/ manifest.json _SUCCESS events__2023-01-02_*.parquet
...
created_at=__HIVE_DEFAULT_PARTITION__/ manifest.json _SUCCESS ...
DuckDB reads the whole tree as one partitioned dataset and recovers the partition column from the path:
SELECT created_at, count(*)
FROM read_parquet('events/**/*.parquet', hive_partitioning = true)
GROUP BY 1;
Granularity
partition_granularity | Path segment | Bucket bounds |
|---|---|---|
day (default) | col=2023-01-01 | [day, day+1) |
month | col=2023-01 | [month-start, next-month-start) |
year | col=2023 | [year-start, next-year-start) |
Notes and limits
- Cloud prefixes: put a
/after{partition}. Object stores have no directories — the part filename is appended to the resolved prefix verbatim.prefix: "events/{partition}/"yieldsevents/created_at=2023-01-01/part.parquet; omitting the trailing slash concatenates into…created_at=2023-01-01part.parquet. (Localpath:joins as directories, so the slash is optional there.) - Index the partition column. Expansion runs three probes over the base
query (
min,max, NULL-count). Without an index onpartition_bythose are up to three full scans before the first row is exported. - Time zones. Bucket bounds are emitted as
YYYY-MM-DDliterals on PostgreSQL and MySQL; SQL Server gets the unseparatedYYYYMMDDform, which T-SQL parses as ISO regardless of the session’sDATEFORMAT/language. For aTIMESTAMPTZcolumn the comparison happens at the session time zone — pin it (e.g. MySQLSET time_zone = '+00:00') when the exact day boundary matters. - Not compatible with
mode: time_window(time_window already filters by a rolling window). Usepartition_bywithfull,chunked, orincremental. - Not compatible with
mode: cdc(CDC reads the log, not a query), aload:block (per-export or top-level: the loader would load a single partition’s manifest), or a MongoDB source (the bucket probes are SQL). - Not compatible with
chunk_by_key(keyset). Each partition reads its bucket as a subquery around the table, and keyset seek pagination over that shape is not supported — Rivet rejects the combination up front. Atable:export keeps its table per partition, so a rangechunk_columnis type-checked and an unset one auto-resolves to the primary key as usual. --parallel-export-processesis disabled while partitioning is active (child processes re-load the config and can’t see the synthesised partitions); the run executes in-process.plan/checkdo not expand partitions yet. They report the parent export as one un-partitioned job (its{partition}token stays literal in the shown path, the row estimate is the whole span, strategy is the base mode). Treat their output as the per-partition shape, not the campaign.- Validating a single partition today: point
validateat the concrete prefix —rivet validate --config c.yaml --export events --prefix events/created_at=2023-01-01. Validating every partition of an export by its parent name in one command is not yet wired.
Choosing a chunking strategy per partition (large partitions)
partition_by is orthogonal to mode, so each partition runs the export’s mode
inside itself. The right choice depends on how big a single partition is and
whether the chunk key is dense within it. Consider a date that holds 100 M
rows:
| Setup | Source load | Use when |
|---|---|---|
mode: chunked, range chunk_column, key dense/correlated within the partition (e.g. the day ≈ the whole table) | One logical pass: ceil(key_span / chunk_size) indexed range scans, bounded memory. Good. | The common big-partition case. Keep chunk_size ≈ 100 k+. |
mode: chunked, range chunk_column, key sparse within the partition (rows interleaved with other days across the key range) | Chunk windows are computed over the partition’s [min,max] of the key, not its row count → many windows read key-range rows then discard out-of-partition ones. Query + I/O amplification. | Avoid — pick a key correlated with the partition, or a finer partition_granularity so each partition’s key range tightens. |
mode: full (no chunking) | One streaming SELECT. Postgres: server-side cursor, bounded memory, but a single long-running transaction (watch vacuum / replication lag). MySQL: one streaming SELECT over the wire (exec_iter); rivet accumulates only one adaptive batch at a time, so memory stays bounded — the cost is the query staying visible (Sending data) for the whole drain. SQL Server: streams one batch at a time server-side, bounded memory (like Postgres). MongoDB (full/cdc only): native cursor, keyset-paged on _id, bounded memory. | PG / SQL Server partitions that fit a long read. MySQL is memory-safe too, but the single query stays visible for the whole drain — prefer chunked there so each query is short. |
Row counts are always exact regardless of the choice — these trade-offs are about
source load and file fan-out, not correctness. There is no keyset/seek option
for a partition today (see the chunk_by_key limit above), so for a genuinely
huge partition with a sparse key the practical levers are a correlated range key
or a finer granularity.