skip to content

A projection's catch-up subscription redelivers a batch of already-processed events after the consumer crashes and restarts. What must the projection's event handlers do to avoid corrupting the read model on redelivery, and how does at-least-once delivery change how projection code should be written?

level: seniorimportance: should knowfreq 50%

answer

  1. at-least-once, not exactly-once
  2. checkpoint and store write aren't atomic
  3. upsert by stable key beats blind insert
  4. set-value beats increment for idempotency
  5. silent counter drift is a classic symptom

basics

~20 s

The handler must be safe to run twice on the same event without messing up the data - for example, setting a value directly instead of always adding to it - because most event systems can redeliver an event after a crash rather than guaranteeing it's processed exactly once.

solid answer

~40 s

Most event streaming and event-store systems guarantee at-least-once delivery, not exactly-once: after a crash, a subscriber may reprocess events it already applied, because the checkpoint update and the projection-store write aren't a single atomic operation. To stay correct under that, handlers need to be idempotent - reapplying the same event produces the same end state as applying it once. Practical techniques: use upserts keyed by a stable identity derived from the event (aggregate id plus event position) instead of blind inserts, track the last-applied position and skip an incoming event that isn't newer than what's stored, or use 'set to value X' semantics instead of 'increment by X' wherever the event carries the resulting value rather than a delta.

go deeper

for a junior

Should understand the basic idea that an event might be processed twice and the handler needs to not double-count when that happens.

for a middle

Should name at-least-once delivery specifically and describe upsert-by-key or set-not-increment as concrete idempotency techniques.

for a senior

Should explain why the checkpoint and store-write gap causes redelivery, distinguish idempotency from ordering, and describe realistic detection strategies like reconciliation jobs for silent drift.

for a principal

Should discuss this as a system-wide handler-design standard - mandating idempotent-by-construction patterns across all projections in an organization, and evaluating when infrastructure-level exactly-once processing is worth the added complexity versus handler-level idempotency.

## What a delivery guarantee promises Delivery guarantees describe what a messaging or event-store system promises about how many times a consumer will see a given event. | Guarantee | What it actually promises | |---|---| | **Exactly-once** - guaranteed to see each event exactly one time, no more, no less | what people intuitively assume, but true exactly-once delivery across a network with independent failure points is either impossible or requires expensive coordination such as distributed transactions between the checkpoint store and the projection store | | **At-least-once** - what, in practice, essentially every event log or streaming system offers (Kafka, `EventStoreDB`, cloud pub/sub systems) | a consumer is guaranteed to eventually see every event, but may see some of them more than once, typically because the consumer crashed or was redeployed after processing an event but before durably recording that it had done so | ## Where the redelivery risk lives The projector's **checkpoint** (its saved read position) and its write to the projection store are two separate operations against two separate stores, and there's no way to make both happen atomically together without significant extra machinery, so a gap between them is where redelivery risk lives: if the crash happens after the store write but before the checkpoint advances, restart replays that same event again. ## Why idempotency is structural This is why idempotency isn't optional polish for projection handlers - it's a structural requirement of building on at-least-once infrastructure. An **idempotent handler** produces the same resulting state whether it processes a given event once or multiple times. - The most robust technique is to make the handler's operation an **upsert keyed on a stable, event-derived identity**: for example, keying a row update on aggregate id and event position so that reapplying the same event just overwrites the same row with the same values, rather than inserting a duplicate row or double-incrementing a counter. - Where the event itself carries the resulting value rather than a delta, a straightforward **'set' operation** is naturally idempotent. - Where events are inherently delta-based, handlers can instead track the last-applied position per record and explicitly **compare-and-skip**: before applying an incoming event, check whether its position is already less than or equal to the stored last-applied position for that record, and no-op if so. A related but distinct concern is **ordering**: within a single stream (a single aggregate's events), most event stores guarantee strict ordering, but across streams a subscriber consuming multiple streams concurrently may see events from different aggregates interleaved in ways that aren't globally ordered, and if a projection's logic accidentally depends on cross-aggregate ordering it hasn't been guaranteed, race-condition bugs can appear under concurrent load even without any redelivery involved. ## The trade-off The trade-off with strict idempotency is mostly code complexity and, sometimes, storage overhead: tracking per-record last-applied positions means extra columns or metadata beyond what the happy-path logic would otherwise need, and designing every handler around upsert-by-stable-key discipline requires more care than naive increment-in-place code. Some teams instead invest in reducing the probability of redelivery, for instance by checkpointing more frequently, but this only narrows the window, it never eliminates it, so idempotency remains the only fully reliable guarantee; treating reduced-frequency redelivery as good enough is a common source of rare, hard-to-reproduce production bugs. ## Failure modes when it is skipped Failure modes when this is skipped are recognizable and often subtle because they're intermittent, tied to exactly when a crash or redeploy happens to land relative to a checkpoint. 1. A **counter-style column**, such as order counts or inventory levels, implemented as increment-on-event rather than idempotent upsert will **silently drift upward** over time, with each production incident or deploy-triggered restart nudging it a little further from the truth - the kind of bug that's hard to catch in testing because it requires an actual crash-and-redeliver sequence to manifest, and shows up months later as 'why doesn't this total match the sum of the underlying orders.' 2. A second failure mode is **duplicate rows from naive inserts instead of upserts**: a handler that does a plain insert on every `OrderPlaced` event, rather than an upsert keyed on order id, produces duplicate rows on redelivery, which then double-count in aggregate queries built on top of that table. 3. A third, subtler failure is an idempotency check that uses **the wrong key** - for instance, deduplicating on event type alone instead of a truly unique identifier - which silently drops legitimate subsequent events of the same type rather than correctly allowing them through and only blocking true duplicates. ## A concrete example at the infrastructure level Kafka Streams is a widely used concrete example of a system that takes this seriously at the infrastructure level: it offers an exactly-once-semantics processing mode that internally still relies on at-least-once delivery from the broker, but achieves effectively-exactly-once processing by making the consumer offset commit and the output write part of the same transaction, which is precisely the kind of extra machinery referenced above as the alternative to writing idempotent handlers by hand.

  • Why can't event systems just guarantee true exactly-once delivery and remove this problem entirely?
    Guaranteeing exactly-once across independent failure points (the event store, the network, the consumer process, and its projection store) requires atomically coordinating a checkpoint update with a separate store's write, which either isn't possible without distributed transactions or is expensive enough that most systems don't offer it by default. It's cheaper and simpler system-wide to guarantee at-least-once and push idempotency responsibility onto the consumer, which only needs to solve it locally.
  • How would you detect that a projection has silently drifted due to a non-idempotent counter-increment handler?
    Periodically reconcile the derived counter against an independent recomputation from the source events - for example, comparing a lifetime-spend column against a fresh sum over that customer's PaymentCaptured events - and alert on mismatches. Because the drift is caused by rare crash timing, it won't show up in normal functional tests, so an ongoing reconciliation job is usually the only reliable detection method.
  • Does making a handler idempotent also make it safe against out-of-order delivery?
    No, those are different problems: idempotency handles seeing the same event more than once, while ordering handles seeing events in an unexpected sequence. A handler can be perfectly idempotent and still be wrong if it silently assumes cross-stream ordering that its event source doesn't actually guarantee, so both properties need to be reasoned about separately.

Like a mail carrier who, if unsure whether you already got a package, delivers it again rather than risk skipping it - fine if you just place the new package where the old one already sits (idempotent), but a disaster if 'delivering' means adding another $50 to a cash envelope every time (non-idempotent).

saying these in an interview costs you the question

  • assumes the event store guarantees exactly-once delivery by default
  • writes counter logic as blind increment-on-event with no dedup key
  • conflates idempotency with ordering guarantees
  • thinks more frequent checkpointing alone eliminates the need for idempotent handlers
  • uses plain insert instead of upsert-by-key for handlers that may be redelivered

context