How do you replay quarantined records after a fix without double-loading rows that already landed?
answer
- run it through the real pipeline
- key the write, never plain insert
- the quarantined copy may be the older one
- replay only what the fix addressed
- the marts do not know anything changed
basics
~20 sReplay through the normal pipeline with an idempotent keyed merge, select only the rejects whose reason matches the fix, guard the write with the record's own version or event time so stale payloads cannot overwrite newer state, and mark each record replayed under a conditional status update.
solid answer
~60 sFour rules make replay safe. **Same code path** — feed the quarantined payloads back through the ordinary ingestion job with a replay run id, never through a one-off script, so the fix is actually the thing being exercised. **Idempotent write** — merge on the business key rather than inserting, so a record already loaded by a later run is updated, not duplicated. **Staleness guard** — a record quarantined three days ago may be older than what the target now holds, so make the update conditional on the record's own event timestamp or version being newer; otherwise a well-meaning replay regresses live data. **Scoped selection with a lifecycle** — replay only the rows whose reason code matches the fix, within the affected time window, and move them from `pending_replay` to `replayed` with a conditional update so two concurrent replays cannot both claim the same rows. Records that fail again are re-quarantined with an incremented attempt count. Finally, tell downstream: refresh the aggregates for the affected partitions, or the marts stay wrong even though the base table is now right.
code
sql · 6 lines-- claim rows for this replay so two concurrent replays cannot both take them
UPDATE orders_quarantine
SET status = 'replaying', replay_run_id = :run_id
WHERE status = 'pending_replay'
AND reject_reason = 'date_unparseable'
AND rejected_at >= :window_start;go deeper
Know that a replay must not create duplicates, and that keying the write on a business identifier is how that is avoided. Recognising the quarantined record as the input to the normal job is the right instinct here.
Explain the merge, the reason-scoped selection and the status transitions. Be ready to say why a bespoke replay script is worse than reusing the pipeline that was just fixed.
Bring up the staleness guard unprompted — an old rejected payload overwriting newer state is the incident that replay causes — plus concurrency-safe claiming, attempt caps and post-replay reconciliation.
Own replay as a standard, audited operation across pipelines: who authorises it, how it is scoped and reported, and how downstream refresh is guaranteed so a corrected base table actually reaches consumers.
## Replay is an ingestion run, not a rescue script The first decision is architectural. A replay that goes through a hand-written script written during the incident tests the script, not the fix, and it usually bypasses the very validations the pipeline exists to apply. The right shape is to make the ingestion job accept an alternative input source — the quarantined payloads — and otherwise behave exactly as it does for fresh data: same parsing, same validation, same write, same metrics, tagged with a `replay_run_id` so its output is attributable. If a replayed record fails again, it goes back to quarantine through the normal path with an incremented attempt count. ## Idempotency: merge on the business key An `INSERT` of replayed rows duplicates anything that arrived by another route in the meantime — a later batch, a manual correction, a full refresh. The write must therefore be keyed: a `MERGE`/upsert on the natural or business key, so replaying the same record twice converges on one row. Where the target is append-only by design (an events table, a raw landing zone), idempotency has to come from a deduplication key stored with the row — the source locator or a content hash — so a re-run can be filtered rather than deduplicated after the fact. ## The staleness guard, which is the part people forget Quarantined records are, by definition, old. Suppose an order row was rejected on Monday for a bad currency code, and on Wednesday the source sent a corrected version that loaded cleanly. Replaying Monday's record on Thursday with a plain upsert overwrites Wednesday's good data with Monday's bad shape. The replay has now caused a worse incident than the original rejection. The guard is to make the update conditional on the incoming record being newer than the stored one, using the record's own **source event timestamp** or **version**, not the ingestion time: `WHEN MATCHED AND src.source_updated_at > tgt.source_updated_at THEN UPDATE`. For change-capture-style sources, the natural version is the source's own sequence or commit position. Where no version exists at all, the honest alternative is not to replay the payload but to **re-extract the current state of that key from the source**, which sidesteps the ordering question entirely — and is often the better answer for small reject sets. ## Selecting the right rows A blind "replay everything in quarantine" re-runs failures that nothing has fixed, burning attempts and producing noise that hides the records that would now succeed. Select deliberately: - **By reason code** — replay only `date_unparseable` because that is what the deploy fixed. - **By time window or run id** — only the runs affected by the defect. - **By status** — only records a human moved to `pending_replay`, never `discarded` ones. This is why the reason code needs to be a stable identifier rather than an exception string: the replay predicate is a `WHERE` clause on it. ## Lifecycle transitions that survive concurrency Two engineers triaging the same incident can both launch a replay. Claim the rows with a conditional update — set `status = 'replaying'` and the `replay_run_id` only `WHERE status = 'pending_replay'` — and process only the rows the update actually claimed. On success move them to `replayed`; on repeated failure back to `new` with `attempt_count + 1`, and stop retrying a record past a cap so it is triaged by a human instead of cycling. Keeping the record rather than deleting it preserves the audit trail: you can later prove which rows entered the target by replay and when. ## Verify, then tell downstream After a replay, reconcile: the count of records marked `replayed` should equal the change in target row count plus the number that matched and updated existing rows. If the numbers do not add up, something was filtered by the staleness guard — that is fine, but it should be *reported*, not discovered later. Then handle the derived layer. Aggregates, incremental models and extracts built from the target on Tuesday do not know that Tuesday's data changed on Thursday. A replay that lands rows into base tables without triggering a refresh of the affected partitions leaves every downstream number stale, which is a subtler version of the original problem. If downstream models are incremental with a watermark on load time, a replay of old business dates may be invisible to them entirely — that is the case where the model needs a full refresh of the affected range, not another incremental run. ## What good looks like Replay is a first-class, repeatable operation: scoped by reason, run through the real pipeline, written idempotently with a version guard, tracked by status transitions that tolerate concurrency, reconciled afterwards, and followed by a targeted downstream refresh.
- Why does replaying an old quarantined record risk making things worse?Because the source may have re-sent a corrected version that already loaded. A plain upsert then overwrites good current data with the stale rejected payload. Guard the write on the record's own event timestamp or version so only genuinely newer payloads win, or re-extract that key's current state from the source instead of replaying.
- What if the record has no version or event timestamp to compare?Do not replay the payload. Use the quarantined row only as a list of affected keys and re-extract their current state from the source, which makes ordering irrelevant. For small reject sets this is usually simpler and safer than inventing a version; for large ones it argues for getting a version into the contract.
- Why can downstream models stay wrong after a successful replay?Incremental models typically select by a load or business-date watermark that has already passed, so rows landing late for old dates are never picked up. After a replay, trigger a full refresh of the affected partition range rather than another incremental run, and treat that as part of the replay procedure.
saying these in an interview costs you the question
- Inserting replayed rows instead of merging on the business key
- Replaying everything in quarantine regardless of reason code
- Letting an old rejected payload overwrite newer loaded state
- Deleting quarantine rows on replay and losing the audit trail
- Declaring the replay done without refreshing downstream aggregates