pattern

Primary-Key-Partitioned Upsert

also called Upsert Partitioning Contract, Key-Routed Upsert

Routing every version of a record to one partition so a real-time store can resolve the latest value locally - which buys sub-second latest-value queries and fixes the partition count and key budget for the table's lifetime.

pinotupsertpartitioningheapclickhouse

A product needs the current state of 400 million orders, queryable in under half a second, from a stream that emits an update every time an order changes. An append-only analytics table answers this by scanning every version and keeping the newest, which is precisely the work that will not fit in the latency budget at concurrency.

The alternative is to make the newest version findable in constant time. If every version of a key is guaranteed to arrive in the same partition, and therefore on the same server, that server can keep a map from primary key to the location of the latest record and resolve queries without coordination. That guarantee is the pattern, and it is a contract on the producer, not a feature of the store.

Why it matters

The requirement runs backwards from query to producer. Apache Pinot's upsert documentation states both halves: the input stream must be partitioned by the primary key so records for a key map to the same server, and each server maintains an in-memory map from primary key to the location of the latest record for that key.

That makes two ordinarily flexible decisions permanent. The partition count of the source topic cannot be changed after the table is created, because changing it re-maps keys to different partitions and the locality the design depends on is gone. And heap becomes a function of distinct keys rather than of rows, so capacity planning shifts from bytes ingested to keys alive.

Implementation patterns

  • Set the key on the producer, explicitly, and treat it as part of the schema. For Kafka this means the record key, not a field in the value.
  • Over-provision partitions at creation, since this is the one dimension you cannot revise.
  • Budget the primary-key index as capacity. At roughly 60 to 100 bytes per entry, 400 million keys is on the order of 25 to 40 GB of heap across the servers holding those partitions.
  • Coarsen the key if it is unbounded. Latest-value-per-key-per-day turns an ever-growing index into a bounded one, and is often what the product actually needs.
  • Decide the deletion story up front, because a key that never retires never leaves the index.
  • Keep a replay path, since upsert metadata can diverge across replicas around segment commits and rebuilding a replica is the remedy.

Industry example

Pinot is the clearest documented case of the contract being written down, and the contrast with ClickHouse decides which resource a table that must stay in production for years will run out of: ReplacingMergeTree deduplicates only when parts merge, offers no guarantee that duplicates are absent at any given moment, and requires the FINAL modifier to get correct results at query time — which the documentation notes is substantially slower. The two stores place the same cost in different places: Pinot pays memory at write time, ClickHouse pays CPU at read time. Neither avoids the cost, and choosing between them is choosing which resource you have spare.

Failure scenarios

  • A partition-count change after the table exists, which silently splits a key's history across partitions and makes the "latest" value depend on which server answers.
  • Heap exhaustion as the key set grows, with no corresponding growth in ingest rate to explain it.
  • Replica divergence around segment commits, showing up as the same query returning two different answers depending on routing.
  • A producer that forgets the key on one code path, so a subset of updates round-robin across partitions and their latest values are wrong permanently.

Trade-offs

Gained: latest-value-per-key at sub-second latency, high concurrency, and no query-time deduplication. Paid: an irreversible partition count, heap proportional to distinct keys, the loss of some aggregation-time optimisations, and a producer contract that every writer must honour forever. The last of these is the one that breaks in practice, because it is enforced by convention rather than by the store.

When not to use it

If the key set is very large or genuinely unbounded, do not pay memory per key: ingest append-only and deduplicate at query time, accepting the slower reads. If the product can tolerate a minute of staleness, a periodic compaction or materialised view over an append-only table is simpler and has no permanent partitioning contract. And if the key set is in the low millions, this is all cheap — the pattern needs no justification and none of these limits will be reached. The pattern earns its constraints only in the band where the key set is large enough to cost memory and small enough to fit it.

Interview question

Q: A team wants latest-value-per-key queries over 400 million orders at p95 under 500 ms from a stream of updates. Walk me through what you require of the producer, what you can never change afterwards, and what you would do if the key count were 40 billion instead.

What a strong answer covers: partition the input by primary key so a server owns a key's whole history · the in-memory key-to-location index and its capacity implication (tens of GB at 400 million keys) · the irreversible partition count and the practice of over-provisioning · at 40 billion keys, abandon in-store upsert for append-only plus query-time deduplication, or coarsen the key by time bucket · and the producer-side failure where one code path omits the key.

Quick check

Quiz: Why can a real-time store's upsert table not have its source partition count changed after creation? Records for a key must always land on the server that holds that key's index entry; re-partitioning re-maps keys and splits a key's history across servers.

Flashcard: What does an upsert table cost that an append-only table does not? — Heap proportional to distinct keys for the primary-key index, plus a permanent contract that producers key every record and the partition count never changes.