Replay and Reprocessing in Event-Driven Systems

Learn how to replay event streams safely, choose offsets and cursors, isolate side effects, and recover projections without confusing replay with retry.

published: reading time: 15 min read author: GeekWorkBench
Quick Summary

Replay rebuilds or repairs a consumer’s output by reading a chosen range of retained events, while retry repeats a narrow failed operation and backfill loads history from another source. The guide covers offsets, isolated destinations, idempotent writes, disabled side effects, and bounded runbooks. Operators can use its checks to compare rebuilt state with business invariants before cutover.

Replay and Reprocessing in Event-Driven Systems

Introduction

A projection can need rebuilding after a consumer bug while live events keep arriving. Replay lets a separate consumer read retained history into an isolated destination, then compare the result before readers switch over.

That job is different from retrying one failed operation: replay covers a chosen historical range, so it needs explicit boundaries, safe repeated writes, and external side effects disabled by default. This guide walks through those controls and the checks to make before cutover.

Replay, Retry, and Backfill

These terms describe different recovery jobs, even if all of them read records more than once.

Operation What happens Typical reason Main risk
Retry The same failed operation is attempted again, usually soon after failure Temporary dependency failure A prior attempt may have succeeded despite a lost acknowledgement
Replay A consumer reads retained events again over a chosen range Rebuild a projection, repair a consumer bug, or validate a new projection Repeating external effects or mixing old and live writes
Backfill Historical data is processed to populate missing or newly introduced state New index, field, report, or downstream consumer Large load on the source and destination; incomplete historical coverage

A retry is usually narrow and operationally urgent. A replay is a planned pass over prior records. A backfill can use an event log, a database export, or another historical source; it does not necessarily mean the system is event-sourced. With event sourcing, the event log is authoritative and rebuilding state from it is a normal capability. In a broker-based event-driven system, the broker may retain messages only for a bounded period.

For the event-sourced case, see event sourcing. For deduplicating repeated commands and events, see API idempotency and safe replays.

Offsets, Cursors, and Retention

An offset or cursor is a consumer’s position in an ordered stream. It answers “where should this consumer resume?” It does not record whether every business side effect caused by earlier records happened exactly once.

In Kafka, offsets are tracked per consumer group and partition. A committed offset marks the position from which that group should resume; applications commonly commit the next offset to read. Since partitions advance independently, a replay can begin at different positions in different partitions. A timestamp-based start is also an estimate of a range, not a guarantee that every event with a matching business timestamp is present.

Before choosing a start position, check:

  • Retention: Are the needed records still available? Kafka retention may remove records by age or size. Compaction can retain the latest value for a key while removing intermediate history.
  • Scope: Which topics, partitions, tenants, and time ranges are in scope? Write them down before execution.
  • Ordering: Does the consumer depend on order within a key or partition? Keep related events on the same ordered path where the design requires it.
  • Schema compatibility: Can the current consumer decode old event versions? Use versioned schemas or upcasters where needed.
  • Position ownership: Is the consumer group shared with live processing? A replay must not unexpectedly move the live group’s committed position.

Kafka’s official documentation describes changing group offsets and the operational behavior of offset reset: Kafka basic operations: consumer group offsets. Treat a reset as a production change: verify the group, topic, partitions, and intended start positions before applying it.

Immutable Events and Corrections

If an event is wrong, editing the historical record in place destroys the audit trail and can leave consumers with different views of history. Prefer an explicit correction event that states what changed, references the original event or entity, and carries the reason and effective time when those concepts matter.

Replay then has a clear meaning: consumers process the same immutable history, including its correction. A corrected event is a new fact in the log, not a silent replacement. Some systems do transform old schema versions while reading; that upcasting should be deterministic and documented because it changes how history is interpreted.

Do not assume every broker preserves an eternal event log. If replay beyond broker retention is a requirement, keep an appropriately governed archive or event store and test restoration from it.

When to Use and When Not to Use

Use replay when the event history is available and you need to recompute a consumer’s derived state. For example, replay a fixed projection after correcting its reducer, or let a new analytics consumer read retained events to build its first dataset. In both cases, keep the destination and progress independent from live processing until you have checked the result.

Choose a retry when one recent operation failed because a dependency timed out or returned a temporary error. Retrying that operation is narrower and avoids reading a large historical range. Make the retry idempotent because the first attempt may have completed even if its response was lost.

Use a backfill when the required history lives in a database snapshot, warehouse, or archive rather than the event stream. Replay cannot reconstruct records that retention removed. A backfill can load that source, with an explicit boundary to prevent overlap or gaps when it joins live events.

Do not replay merely to correct a past business fact or repeat an external action. Append a correction event when history needs a new fact; use a separately reviewed, deduplicated job when a payment, notification, or webhook must happen. If the current projection is already correct and no consumer needs historical input, replay adds load without fixing a problem.

A Safe Replay Flow

The following flow puts isolation and verification around the actual read. The live consumer keeps its own position while a replay builds a separate result.

flowchart LR
    Log[(Retained event log)] -->|range and partitions| Replay[Replay consumer]
    Replay --> Validate[Decode and validate]
    Validate --> Project[Isolated projection or staging table]
    Project --> Compare[Compare counts and invariants]
    Compare -->|approved| Switch[Switch read traffic]
    Replay -.->|live side effects disabled| SideEffects[Email, payment, external APIs]
    Live[Live consumer group] -->|independent offsets| Log

“Exactly once” should not be inferred from replaying a log. A consumer may receive a record again after a crash, and a database update or HTTP call may already have succeeded before the consumer’s offset was committed. Idempotent writes, deduplication keys, transactions supported by the relevant components, and explicit effect controls address different parts of that problem. None makes unrelated systems atomically commit together by default.

Kafka-Style Replay Runbook

This procedure avoids relying on a particular command-line tool version. The same safeguards apply whether an operator uses an admin client, a deployment job, or a one-off consumer.

  1. State the repair. Name the projection or consumer bug, affected range, expected output, owner, and rollback path. Establish whether the event source is complete for that range.
  2. Choose the destination. Prefer a new projection table, index, or consumer group. Keep live offsets untouched. If the application supports a shadow projection, write there first.
  3. Set replay boundaries. Record the topic, partition set, start and end offsets (or equivalent cursor), schema version, and any tenant filters. Confirm retention covers the entire range.
  4. Make processing safe. Disable email, payment, notification, and other external calls. Make writes idempotent by stable event ID or business key. If an external effect is part of the repair, give it a separate reviewed job with its own deduplication strategy.
  5. Run a small sample. Process a limited range and check parse failures, duplicate handling, row counts, ordering assumptions, and resource use. Compare known entities against the source of truth.
  6. Process in bounded batches. Apply rate limits and backpressure. Track the last completed position per partition. Avoid committing beyond work that the destination has durably accepted.
  7. Validate before cutover. Compare counts and business invariants, inspect lag and errors, then direct reads to the rebuilt result. Keep the prior projection available for rollback until the new one is accepted.
  8. Close out. Record final positions, code version, operators, duration, exceptions, and cleanup actions. Remove temporary resources only after the rollback window ends.

For a generic consumer, the core loop is:

for each partition in selected_partitions:
    position = configured_start[partition]
    while position < configured_end[partition]:
        records = read_batch(partition, position, configured_end[partition])
        for record in records:
            event = decode_and_validate(record)
            upsert_projection(event, idempotency_key=record.event_id)
        wait_until_destination_is_durable()
        position = next_position(records)
        record_replay_checkpoint(partition, position)

The checkpoint belongs to the replay job and destination. In a real implementation, ensure checkpoint advancement cannot get ahead of durable writes. If the destination and checkpoint cannot share a transaction, replaying a batch after a crash should be safe.

Isolation and Idempotency

Use a separate consumer group or explicit range reader for replay. Reusing a live group can cause a reset to rewind or advance live consumption. A distinct group also makes replay lag and progress visible, though it does not isolate writes by itself; the output destination must be separated too.

For projections, a common pattern is an upsert keyed by entity ID plus a processed-event ledger keyed by stable event ID. Whether the ledger is useful depends on the projection’s update semantics and storage cost. If each event deterministically replaces a projection with a versioned state, a conditional upsert may be enough. If it increments a counter, duplicate delivery can inflate the result unless the update is deduplicated or made idempotent.

External effects need their own boundary. A replay should default to suppressing them. If the business task requires sending a corrected notification, decide which recipients qualify and use an effect ledger or provider-supported idempotency key. Do not route replay output back into the same topic without a deliberate design; that can create event loops or duplicate downstream work.

Production Failure Scenarios

Failure What you may observe Mitigation
Live consumer group is reset by mistake Live traffic rereads old events or skips expected work Use a dedicated replay group; review the group and partition plan before changing offsets; alert on unexpected offset movement
Replay writes twice after a crash Counts or balances drift after restarting from an earlier checkpoint Make writes idempotent; couple durable output and checkpoint where possible; otherwise accept safe batch repetition
Replay overloads a database or search index Query latency, write timeouts, or elevated live lag Throttle replay, use a separate destination, reserve capacity, and pause automatically when live SLOs degrade
Old events no longer decode Poison records or a large failure spike near a schema change Test representative old versions, add deterministic upcasters, and quarantine malformed records with a repair process
History is incomplete The rebuilt projection disagrees with source totals Check retention, compaction, archival gaps, and start offsets before cutover; restore from the authoritative archive if available
External side effect is repeated Duplicate email, webhook, or charge Disable side effects during projection rebuilds; use stable idempotency keys and an independently controlled effect job

Some records will still fail. A dead-letter queue can preserve failures for inspection, but it does not fix bad data automatically. Record the original topic, partition, offset, event ID, error, and replay job so a repaired event can be traced to its source.

Trade-off Analysis

Choice Advantages Costs and risks Good fit
Rebuild into a new destination Live reads and offsets remain stable; easy side-by-side comparison Needs temporary storage and a cutover plan Large projection repairs or risky code changes
Replay into the current destination Less infrastructure and no read switch Partial results can be visible; rollback is harder Small, idempotent corrections with a clear repair boundary
Reset the existing consumer group Uses the normal consumer deployment path Can rewind live processing or alter lag unexpectedly Controlled environments with coordinated downtime and verified offsets
Separate replay group Independent progress and observability Duplicates broker reads; downstream writes still need isolation Most production replays where live processing continues

The “new destination” option often costs more up front but gives operators a clean comparison point. Replaying into the current table is reasonable when the operation is genuinely idempotent and partial state is acceptable; otherwise, the shortcut tends to hide whether the repair is complete.

Observability Checklist

Before and during a replay, make the following visible on a dashboard or run record:

  • Replay job ID, owner, code version, source range, and destination name
  • Start and end offsets per partition, plus the current checkpoint
  • Records read, accepted, rejected, deduplicated, and written
  • Processing rate, consumer lag, estimated time remaining, and retry count
  • Decode and validation errors grouped by event type and schema version
  • Destination latency, capacity, and impact on live traffic
  • External effect suppression status and any separately authorized effects
  • Business-level checks, such as entity counts, balances, or aggregate totals
  • Cutover decision, rollback state, and final checkpoint

Consumer lag alone is not a completion signal. A job can reach the end offset while dropping malformed records, writing to the wrong destination, or violating a business invariant.

Security and Data Retention

Replay access can expose old personal or confidential data that a live consumer no longer routinely reads. Apply the same authorization and audit controls as production processing, limit the operator’s access to the requested topics and range, and avoid copying raw payloads into broadly accessible logs.

Retention policy should state how long event data remains in the broker and whether a governed archive exists. A replay plan must honor deletion, legal hold, and data minimization requirements. If data must be erased, rebuilding a projection from an archive that still contains it may reintroduce the deleted data. Design deletion markers or filtering rules into the replay path and verify the resulting projection before cutover.

Common Pitfalls / Anti-Patterns

  • Treating a replay as a safe substitute for retry, even when a remote call may already have succeeded.
  • Resetting the live consumer’s offsets without recording the original positions.
  • Starting from a timestamp and assuming it maps exactly to the desired business interval.
  • Rebuilding a projection against compacted data as if every historical update were still retained.
  • Updating events in place to “correct” history instead of appending an explicit correction.
  • Committing offsets before the destination write is durable.
  • Turning on external side effects because the replay consumer uses the normal handler code.
  • Declaring success from record counts alone without comparing business invariants.
  • Keeping replay permissions or temporary archives after the recovery window.

Quick Recap Checklist

  • Name the repair and the intended output before reading old records.
  • Confirm retention, partitions, schema compatibility, and exact replay boundaries.
  • Keep the live consumer position independent from replay progress.
  • Write to an isolated destination when partial results would be risky.
  • Make repeated writes safe; suppress external effects by default.
  • Track durable checkpoints, errors, lag, and business invariants.
  • Compare results before cutover and document rollback and cleanup.

Interview Questions

1. How does replay differ from retry?
A retry repeats a failed operation, usually for a narrow set of work. A replay reads a chosen range of historical events again, often to rebuild or repair a consumer’s output. Both can deliver the same event more than once, so replay does not remove the need for idempotency.
2. Why use a separate consumer group for a replay?
A separate group has independent offsets, so replay progress does not rewind or advance the live group. It also gives the job its own lag and checkpoint view. It does not isolate the destination or external effects; those need separate controls.
3. When should corrected data be emitted as a new event?
When the original event was an accepted historical fact that later needs correction, append a correction or compensating event with a reference to the original. Avoid editing immutable history in place because consumers and audit tools may have already observed it.
4. Does committing a Kafka offset prove that processing happened exactly once?
No. A consumer can write to a database and crash before committing its offset, causing the record to be read again. Or it can commit before a non-transactional side effect completes. Use idempotent effects and coordinate output with checkpoints where the system supports it; do not assume an offset alone makes separate systems atomic.

Further Reading

Conclusion

Replay is a recovery tool, not a guarantee about delivery or side effects. Identify the output to rebuild, choose an explicit event range, keep replay progress separate from live consumption, and make repeated writes safe. Then validate the rebuilt state against business invariants before switching readers. Those steps turn “read the topic again” into a controlled repair.

Category

Related Posts

Event Envelopes and Metadata for Reliable EDA

Learn how event envelopes separate transport context from payload, apply CloudEvents attributes, propagate correlation IDs, and validate messages safely.

#event-driven-architecture #cloudevents #messaging

Event Security and Sensitive Data in EDA

Secure event-driven systems with least-privilege identities, encrypted transport and payloads, careful data minimization, and a deliberate retention plan.

#event-driven-architecture #security #data-privacy

Idempotency, Deduplication, and Safe Replays

Design idempotent API operations and deduplication records so clients can retry after timeouts without creating duplicate payments, jobs, or updates.

#api-design #idempotency #retries