Introduction
Offline upserts allow an OFFLINE Pinot table to upsert records by primary key across multiple segments, presenting query results as if only the most-recent version of each record exists. In a standard OFFLINE table every segment is independent — pushing two segments that each contain a row for the same primary key means queries will return both rows. Offline upserts fix this: when a new segment is pushed, the server resolves conflicts with all previously loaded segments and surfaces only the winning record per primary key.Record precedence ordering
The row with the greater comparison value wins. The comparison value comes fromupsertConfig.comparisonColumns when it is set, and from the table’s timeColumnName otherwise. A table needs one of the two, or the server refuses to load it.
Segment push time does not pick the winner. It only breaks a tie between two rows with the same comparison value, and then the row from the segment pushed most recently wins. If you want the most recent push to win regardless of the data, add a column that increases with every push, such as a batch timestamp, and set it as the comparison column. The full set of tie breakers is described under Record Comparison and Tie Breakers.
Use Cases
Offline upserts are a good fit when the number of rows being updated is a small-to-moderate fraction of the total table size (roughly less than ~40%). For larger fractions — where you are effectively replacing the whole table — consider Atomic Sync (see below).
How to Configure Offline Upserts
Enabling offline upserts requires changes to three places: the schema, the table config, and (optionally) the ingestion task config.1. Schema — declare primary key columns
primaryKeyColumns can contain multiple columns — Pinot hashes them together as a composite key.
2. Table config
Three sections of the table config must be set together:upsertConfig, segmentPartitionConfig, and routing.
upsertConfig.modemust be"FULL". Partial upsert mode is not supported for offline tables (see Limitations).upsertConfig.metadataManagerClassmust point to the StarTree RocksDB implementation. If omitted the in-heap OSS default is used, which does not persist metadata across server restarts and is not recommended for production.segmentPartitionConfigon the primary key column is recommended. The partition function and count define how the data is partitioned and also governs the overall scalability. Every segment for a given primary key must land on the same server — without this the server cannot see all versions of a key and deduplication will be incorrect. In general, the recommendation is to use a high partition count (eg: 128) to account for organic growth. Note that you cannot change this post table creation.instanceSelectorType: strictReplicaGroupensures that a query is routed to exactly one replica group. This is required so that the server’s local RocksDB view (which is per-partition) is authoritative for the keys it owns.
comparisonColumns is set, the row with the highest value in that column wins. When it is omitted, the table’s timeColumnName is the comparison column. Segment push time only breaks ties between rows with the same comparison value.
3. Partitioning upstream data with FileIngestionTask
For offline upserts to work correctly, every segment must be partitioned on the primary key: all rows with the same primary key value must reside in the same server.FileIngestionTask handles repartitioning automatically, even when the upstream S3 data is not pre-partitioned by primary key. On each execution the task executor reads segmentPartitionConfig from the table config and passes it to SegmentProcessorFramework, which physically sorts all ingested rows into per-partition buckets by applying the configured partition function (e.g., Murmur) across the configured number of partitions. Each output segment contains only rows that hash to a single partition bucket, satisfying the upsert requirement without any extra config.
The segmentPartitionConfig in the table config acts as the contract: the task executor reads it to know how many partitions to create and which hash function to use, producing pre-partitioned segments that satisfy the upsert requirement automatically.
Offline Upserts vs. Atomic Sync
StarTree also offers Atomic Sync (viaconsistentPushSwapEnabled=true on FileIngestionTask), which replaces the entire table atomically. Both features update an offline table from S3, but they are designed for different scenarios.
Rule of thumb: if less than ~40% of the table’s rows are changing in a given push cycle, offline upserts are the better fit — they avoid re-ingesting unchanged data and are cheaper to operate. Above ~40%, or when you need atomic visibility guarantees, use Atomic Sync.
Limitations
No partial upsert support. Onlymode: FULL is supported for offline tables. FULL mode means the entire winning row replaces the entire losing row. Column-level merge strategies available in PARTIAL mode for realtime upsert tables are not available. Note that this is generally ok since column level update pattern is typically found in CDC/streaming use cases not batch.
No StarTree index. Columns in an offline upsert table cannot have a StarTree (multi-dimensional pre-aggregation) index configured. The controller will reject such a table config at creation time.
Segment push time breaks ties, and push order cannot be undone. When two rows for the same primary key have the same comparison value, the row from the segment pushed most recently wins. Pushing an older segment after a newer one makes the older data win for those tied rows. Use a comparison column that reflects the real order of your data, such as a business event timestamp, if push order cannot be guaranteed.
