skip to content

In a search-indexing pipeline consuming product change events, how do you stop a stale event from overwriting a newer document?

level: middleimportance: should knowfreq 46%

answer

  1. last write wins, whichever arrives
  2. partition by product ID
  3. version from the source of truth
  4. deletes need a marker

basics

~20 s

Give every event a monotonically increasing version from the source of truth, such as a row version or commit-log position, and write to the index only if that version is newer than the stored one. Partitioning events by product ID keeps each product's changes in order.

solid answer

~40 s

Events reach an indexer out of order because of retries, redelivery, parallel consumers and backfills running beside the live stream. Two defences work together. First, partition the stream by product ID so one product's events are processed in sequence. Second, carry a **version** assigned by the database, such as a per-row counter or the commit-log position, and make every index write **conditional**: apply it only if the incoming version is greater than the one stored with the document, and skip it otherwise. That also makes duplicates harmless. Deletes need a **versioned tombstone** kept for a while, otherwise a late update recreates a deleted product. Avoid application wall-clock timestamps as versions, because clock skew between servers can make a newer change look older.

code

pseudocode · 7 lines
pseudocode
on_event(e):
  doc = (e.op == DELETE) ? TOMBSTONE : e.document
  // compare and write in one atomic step inside the index
  result = index.write_if_newer(e.product_id, e.version, doc)
  if result == REJECTED_STALE:
    metrics.increment("stale_events_skipped")
  ack(e)

go deeper

for a junior

Recall that events can arrive out of order, and that attaching a version from the database lets the index ignore the older ones.

for a middle

Explain why ordering by key is not enough on its own, how a conditional write compares versions atomically, and why deletes need tombstones.

for a senior

Show how you would operate it: the version source you would trust, tombstone retention tied to replay windows, and a metric for skipped stale events.

for a principal

Discuss whether every consumer of the change stream should share one versioning contract, and what it costs to add versions to an existing schema.

## Why change events arrive out of order An indexing pipeline reads change events for products and applies them to a search index. Even when the database committed the changes in a clear order, the events can reach the index in a different one: - **Retries**: a write that failed is retried after a newer write for the same product has already succeeded. - **Redelivery**: queues usually deliver **at least once**, so an old event can arrive again after an indexer restart. - **Parallel consumers**: two workers handling the same product at the same time finish in either order. - **Backfills**: a bulk reload of old data runs beside the live stream and writes older state on top of newer state. An index write is normally a plain **upsert**, so whichever write lands last wins. Without extra information the index cannot tell a late old write from a new one. ## Defence one: keep each product's events together Partitioning the event stream by **product ID** sends every change for a product to the same partition, and a single consumer processes each partition in order. This removes most races cheaply. It does not remove them all: retries inside a consumer, redelivery after a rebalance, and backfills can still reorder writes. So ordering by key is necessary but not sufficient. ## Defence two: version-guarded writes Each event carries a **version** from the source of truth, and the index stores the version with the document. A write is applied only if its version is greater than the stored one. ```pseudocode on_event(e): doc = (e.op == DELETE) ? TOMBSTONE : e.document // compare and write in one atomic step inside the index result = index.write_if_newer(e.product_id, e.version, doc) if result == REJECTED_STALE: metrics.increment("stale_events_skipped") ack(e) ``` The comparison must be **atomic** in the index. Reading the stored version, deciding in the indexer, and then writing leaves a gap in which another worker can write. Many engines support a conditional write keyed on an external version; where one does not, serialize writes per product. A duplicate event carries a version equal to the stored one and is skipped, which makes the write **idempotent**. ## Where the version comes from | Source | Monotonic per product? | Notes | |---|---|---| | Row version counter incremented on each update | Yes | Simple; must be bumped in the same transaction as the change | | Commit-log position of the change | Yes, across the whole database | Comes for free with log-based change capture | | Application server wall-clock timestamp | No | Clock skew between servers can give a newer change an older timestamp | | Indexer arrival time | No | Records when the event arrived, which is exactly what is unreliable | ## Deletes and tombstones Deletes are the classic trap. If a delete simply removes the document, the stored version disappears with it. A delayed update for that product then finds nothing to compare against and recreates the product, so a delisted item reappears in search. The fix is a **tombstone**: keep a marker with the delete's version for longer than the longest expected delay, reject older writes against it, and purge tombstones later. ## The refetch alternative Some pipelines treat an event only as a **trigger**: the indexer reads the current row from the database and indexes that. This is self-healing, because every read returns the latest committed state. Two cautions remain: 1. Two workers can read the row at different moments and write their results in the wrong order, so the row's version should still guard the write. 2. Every event costs a database read, which matters at high change rates and during bulk updates. ## What to monitor - The count of stale events skipped, which should be small and stable. - A sudden rise in skips, which often points to a replay or backfill running unexpectedly. - Reconciliation mismatches, which show whether the guards are actually holding.

  • Why is read-the-version-then-write in the indexer not enough?
    Between reading the stored version and writing, another worker can write a newer version for the same product. The first worker then overwrites it with older data, which is exactly the bug being prevented. The comparison and the write must happen as one atomic step inside the index, or writes for each product must be serialized so only one worker touches a product at a time.
  • How long should tombstones for deleted products be kept?
    Longer than the longest delay a stale event can plausibly have: queue retention, the maximum retry backoff and the duration of any replay or backfill. Keeping them for too short a time lets a late update resurrect a deleted product; keeping them forever wastes space and can skew counts. A cleanup job purges tombstones older than that window.

saying these in an interview costs you the question

  • Partitioning by product ID alone guarantees no stale write ever lands.
  • Application server timestamps are a safe version for conditional writes.
  • Deleting the document outright is enough to handle deletes.
  • Reading the latest row on each event makes version checks unnecessary.
  • Duplicate events must be removed before the indexer ever sees them.