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.
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.
- 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.
- 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.
- 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.
- 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.
- 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.
- 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.
- 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.
- 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
Further Reading
- Apache Kafka: Basic Operations — consumer groups and offset reset operations.
- Event Sourcing — when an event log is the source of truth for state.
- API Idempotency, Deduplication, and Safe Replays — designing operations that tolerate duplicate attempts.
- Dead-Letter Queues — isolating records that repeatedly fail processing.
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 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.
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.