Recoverable pipelines

Buffering, replay, and duplicate handling in an event pipeline.

Illustrative reference architecture 2 min read

Figure 01 / Event processing path

  1. Event sourceStable event ID; retries unacknowledged writes
    acknowledge after durable append
  2. Durable logFinite retention and capacity
    append
    read from recorded progress
  3. ProcessorCan restart and replay
    run
    write output, then record progress
  4. Output storeIdempotent write by event ID
    write

The processor reads events and writes the result. It records progress only after the output is durable.

Illustrative example. No external side effects beyond the output write.

The assumptions

This example follows an event from a source, through a durable log and a processor, into an output store. The source keeps pending events until they are acknowledged and retries with the same event ID. It cannot replay acknowledged history on its own.

The log has finite retention. The output store can atomically record an event ID and apply its effect. Ordering is required within an entity, not across the entire system. There are no external side effects beyond that output write.

The normal path

The source sends an event with a stable identifier. The log acknowledges it after a durable append. If an acknowledgement is lost, the source retries the same event. A repeated ID with different content is a conflict, not a new event.

The processor reads the event and opens an output transaction. It records the identifier and applies the change together. A unique constraint prevents a second copy from applying the same effect. Only after that transaction commits does the processor record its progress in the log.

The gap that matters

A crash before the output commits leaves work to retry. A crash after commit but before the checkpoint causes the same event to be read again. The output transaction must recognize that duplicate and return success without repeating the effect.

Repeated delivery is expected here. The design controls repeated effects at the destination.

With parallel processing, progress advances only beyond a contiguous sequence of completed records. A later success must not let the checkpoint skip an earlier failure. Updates that depend on entity order also need ordering or version checks.

Recovery has a limit

Track the age of the oldest unprocessed event, output failures, retries, and remaining retention headroom. Recovery needs enough processing capacity to handle new arrivals and reduce the backlog.

If the outage exceeds retention, the missing history must come from an archive or reconstruction at the source. Deduplication records also need to cover the supported retry and replay period. Rebuilding a projection with different logic needs a separate version or destination; otherwise old deduplication records may suppress the new work.

Invalid records need a durable quarantine path and an explicit rule for continuing past them. Silently skipping them changes the recovery contract.

What this buys, and what it costs

The log separates arrival from processing and gives the system time to recover. It also adds storage, retention management, and another service to operate. For a small workload that can safely retry directly against a transactional destination, that extra layer may be unnecessary.

This example does not provide a general exactly-once guarantee for messages, payments, or other external actions. Each additional effect needs its own retry and reconciliation design.

References

Search the site

Search experience, studies, articles, projects, and contributions.

Try a topic