pattern

Projection Checkpoint

also called Projection Offset, Consumed Position Marker

The durably recorded stream position a read model has applied up to - written atomically with the model itself so a crash can neither skip events nor re-apply them undetectably.

event-sourcingcqrsprojectionsidempotencyrecovery

A projection applies a batch of 1,200 events to a balance view, commits the view, and the pod is killed before it records how far it got. On restart it re-reads those 1,200 events and applies them again. Every balance built by addition is now wrong, by exactly one batch, and no error was raised anywhere.

Reverse the order and the failure reverses with it: record the position first, crash before writing the view, and those 1,200 events are permanently skipped. The projection is confident it is up to date, the numbers are quietly short, and the only way to find out is to rebuild and compare.

The checkpoint is not bookkeeping. It is the thing that decides whether at-least-once delivery becomes effectively-once state, and the ordering of two writes is the whole mechanism.

Why it matters

Replay is the capability an event log is bought for, and replay is only trustworthy if the projection knows precisely what it has consumed. The scale makes the point: a stream of 1.2 billion events replayed at 50,000 events a second takes about 6.7 hours. A projection with no durable checkpoint pays that on every restart, which in production turns a routine deployment into an outage and makes an autoscaler actively dangerous.

The second reason is diagnosis. Lag — events behind head — is computed from the checkpoint, so a projection without one has no lag metric, and a stalled read model is indistinguishable from a quiet one.

Implementation patterns

  • Write the checkpoint in the same transaction as the view update. One row in the same database as the projection. This is the only arrangement that gives effectively-once without idempotent handlers.
  • If the view store has no transactions, make handlers idempotent instead, keyed on aggregate id plus version with a conditional write, and treat the checkpoint purely as a restart optimisation.
  • Checkpoint per batch, and let the batch size set the worst-case rework. A 1,000-event batch at 0.3 ms per event is about 300 ms of repeated work after a crash, which is a sensible trade against one checkpoint write per batch.
  • Record the handler's code version alongside the position. When the fold logic changes, continuing from the middle produces a view that is part old-logic and part new; the version makes that a deliberate rebuild rather than a silent hybrid.
  • Never keep the checkpoint on local disk or in memory. A rescheduled pod then starts from zero, which is the 6.7-hour rebuild arriving unannounced.

Industry example

Broker-managed consumer offsets, as Kafka provides, are checkpoints you do not control: the broker advances them on a schedule that knows nothing about whether your view transaction committed, which is why auto-commit plus a non-idempotent projection is the most reproducible way to corrupt a read model in production. Teams running projections at volume move the position into their own store precisely so that it commits with the data.

Failure scenarios

  • Position ahead of data — a permanent hole, invisible until a rebuild disagrees with the live view.
  • Position behind data — double application, visible as inflated counters and balances that drift upward.
  • Two instances sharing one checkpoint, each advancing it, so the view is an interleaving of two partial passes.
  • Checkpoint lost on reschedule, causing an unplanned full replay at the moment of a deployment.
  • A logic change applied from the middle of the stream, producing a view no replay can reproduce — the defect that makes a projection un-rebuildable.

Trade-offs

A transactional checkpoint ties the projection to a store that supports transactions, which rules out some fast read stores and adds a write to every batch. Idempotent handlers free you from that at the cost of a version column, a conditional write per event and handler logic that must be provably commutative. The cheap third option, auto-commit with no idempotency, is not a trade-off; it is an undetected correctness bug whose cost arrives as a reconciliation.

When not to use it

When a projection is small enough to rebuild from scratch on every start — a few million events, seconds of work. Then the simplest correct design is to have no checkpoint at all, rebuild into a fresh table on boot and swap. It flips when rebuild time exceeds the restart budget, which for most teams is the point where a deployment can no longer afford a cold start.

Interview question

Q: Your read model's numbers are slightly too high and the event log is intact. Before you look at any code, what are the two orderings you suspect, and how would you prove which one happened?

What a strong answer covers: view-then-checkpoint causing re-application versus checkpoint-then-view causing skips · the direction of the error distinguishing them, since double application inflates and skipping deflates · proving it by rebuilding into a shadow table and diffing · and the fix being either one transaction over both writes or idempotent handlers keyed on version, not a larger batch or more retries.

Quick check

Quiz: A projection commits its view and then its stream position, and the process dies in between. What is the consequence? Those events are re-applied after restart, so any additive field is overstated by one batch, with no error raised.

Flashcard: What has to be true of a projection's checkpoint for at-least-once delivery to produce correct state? — Either the checkpoint commits in the same transaction as the view update, or every handler is idempotent keyed on aggregate version; the checkpoint must also be durable off-pod and stamped with the handler's code version.