Data Reliability and Fault Tolerance Questions
Designing pipelines that survive failures: retries, idempotency, checkpointing, exactly-once semantics, dead-letter handling, and recovery/replay. Covers reasoning about partial failures, poison messages, and consistency guarantees under faults. The resilience angle distinct from monitoring (detecting and alerting on a failure) and from Workflow Orchestration and Scheduling (the DAG/scheduler mechanics that decide whether and when a task runs again, including backfills and dependency management): this topic owns whether the data itself stays correct, not lost, not duplicated, not corrupted, when a process is retried or replayed.
Provide pseudocode for the Chandy-Lamport distributed snapshot algorithm adapted to capture consistent operator state and in-flight messages in a streaming dataflow. Include handling for a snapshot attempt that partially fails and how to resume or abort safely. Then explain how you would adapt the same algorithm to implement consistent checkpoints across cooperating microservices that exchange events over Kafka or message queues, and discuss the practical challenges of capturing in-flight messages and integrating with persisted broker logs.
Sample Answer
Direct answer
The Chandy-Lamport algorithm captures a globally consistent snapshot (every operator's local state, plus every in-flight message on every channel) without stopping the dataflow, using one rule: when an operator sees the FIRST marker on any incoming channel, it immediately records its own local state and forwards a marker on every outgoing channel; any data that arrives on a channel AFTER recording starts but BEFORE that specific channel's own marker is genuinely in-flight and must be recorded as part of that channel's snapshot. A snapshot attempt that never completes (a marker lost or indefinitely delayed) is handled with a bounded wait followed by an explicit abort that discards the partial recording and resets the operator to a clean state, so a fresh attempt can safely retry. The same core rule extends to cooperating microservices over Kafka or a message queue, but the practical challenges shift: a durable broker already persists the channel's history, so "in-flight capture" becomes an offset-range bookkeeping problem rather than a memory-buffering one, at the cost of needing per-partition markers (a Kafka topic is not a single FIFO channel) and needing every participating service to actually understand and propagate markers.
Structured elaboration
The core marker-propagation rule. An operator that is the SNAPSHOT INITIATOR records its own state immediately, then sends a marker on every outgoing channel. Any OTHER operator, on receiving its first marker (on any one incoming channel), does the same: record local state now, propagate markers on all outgoing channels immediately, and begin recording every OTHER incoming channel until that channel's own marker arrives. The snapshot is complete for an operator once a marker has been seen on EVERY incoming channel, since only then is it guaranteed nothing pre-snapshot is still unaccounted for on that channel.
Why FIFO, loss-free channels are the algorithm's core assumption. The whole scheme depends on a marker on a given channel being a reliable boundary: everything the operator receives on that channel BEFORE the marker is "before the cut," everything after is "after the cut." This only works if messages on a channel are delivered in the order sent (FIFO) and none are silently dropped; a channel that can reorder or lose messages can let a genuinely pre-snapshot message arrive AFTER its own marker, silently excluding it from the snapshot with no way to detect the gap.
Handling a snapshot attempt that partially fails: resume versus abort. These are two genuinely different responses, and which one is safe depends on WHY the marker is late. If the delay is merely a slow producer or a backlog the missing channel is still working through (the channel partner is alive, just behind), RESUMING is possible and often preferable: keep waiting, since the channel's own FIFO ordering guarantee still holds and the marker will eventually arrive with the cut still correctly defined, just later than hoped. If the delay reflects a genuine failure (the channel partner crashed, or a marker was lost on an unreliable transport and will never arrive), waiting is pointless and RESUME is not an option at all, no amount of extra time produces a marker that was never sent or was silently dropped. Bounding the wait with a timeout is what turns "which case is this?" from a judgment call made too late into an explicit, enforced decision: below the bound, resume (keep waiting); at the bound, assume the worse case and ABORT, discarding the recorded state and every channel's partial recording buffer, and returning to a not-in-progress state. This is not just bookkeeping hygiene, it is what makes a subsequent RETRY (a fresh snapshot attempt, distinct from resuming the SAME attempt) safe: an operator that aborted cleanly can treat the next marker it sees as a genuinely fresh first marker and begin a new, independent attempt, rather than mixing stale partial state from the failed attempt into a new one.
Adapting to cooperating microservices over Kafka/message queues. The same three-part rule (record on first marker, propagate markers immediately, record data until each channel's own marker) still applies, but three things change in practice:
- The broker is the channel's own persisted history. Unlike the in-memory channel in the reference implementation below, a Kafka topic-partition already durably retains every message. "In-flight capture" becomes: record the OFFSET at which this consumer's snapshot recording started per partition, and the offset at which the marker for that partition was seen; everything between those two offsets IS the channel's snapshot state, retrievable from the broker's own log rather than a buffer the operator must maintain itself.
- A Kafka topic is not one FIFO channel, it is N independent FIFO partitions. Ordering is only guaranteed WITHIN a partition, so a marker must be sent (and waited for) on EVERY partition of every relevant topic a service consumes from, not once per topic; treating a multi-partition topic as a single channel silently reintroduces the reordering risk the algorithm's FIFO assumption is meant to rule out.
- Every participating service must understand markers. The algorithm assumes every node in the graph propagates markers correctly; a service that is unaware of the protocol (a legacy consumer, or a service owned by a different team not opted into the snapshot mechanism) simply will not forward one, which either stalls the snapshot at that service (if it is expected to) or silently produces an incomplete cut if the coordinator does not realize that service was supposed to participate. In practice this argues for either a dedicated coordination topic carrying markers as control messages (out-of-band from the data topics, so unmodified consumers of the data topics are unaffected) or an explicit registry of exactly which services participate in a given consistent-checkpoint boundary.
Worked example
Reference implementation below (a graph with fan-in: A feeds both B and C, and B also feeds C, so C must wait for markers on BOTH its incoming channels, exercising the rule beyond a trivial linear chain). Scenario 1 proves successful capture of a genuinely in-flight message; scenario 2 proves the partial-failure timeout, abort, and clean retry.
"""
Chandy-Lamport distributed snapshot: capture consistent operator state AND
in-flight channel messages in a streaming dataflow with fan-in (a cycle-free
but non-trivial graph, not just a linear pipeline).
Topology: A -> B, A -> C, B -> C (C has two incoming channels, exercising
the "wait for a marker on EVERY incoming channel before declaring complete"
rule that makes Chandy-Lamport work correctly on fan-in, not just a chain).
Scenario 1: successful snapshot with a genuine in-flight message captured.
Scenario 2: a snapshot attempt where one marker never arrives within a
bounded wait -> abort() discards the partial recording and resets the
operator to a clean, resumable state -> a FRESH snapshot attempt afterward
succeeds normally, proving abort did not leave the operator corrupted.
"""
from collections import deque
class Channel:
def __init__(self, src, dst):
self.src, self.dst = src, dst
self.queue = deque()
def send(self, msg):
self.queue.append(msg)
def receive(self):
return self.queue.popleft() if self.queue else None
def __repr__(self):
return f"{self.src}->{self.dst}"
class Operator:
def __init__(self, name):
self.name = name
self.state = 0
self.incoming = []
self.outgoing = []
self.snapshot_in_progress = False
self.recorded_state = None
self.recording_channels = {} # channel -> list of messages recorded before ITS marker
self._markers_by_channel = set()
self.snapshot_complete = False
self.ticks_since_snapshot_started = 0
self.aborted_count = 0
def _begin_snapshot(self):
self.recorded_state = self.state
self.snapshot_in_progress = True
self.snapshot_complete = False
self._markers_by_channel = set()
self.recording_channels = {c: [] for c in self.incoming}
self.ticks_since_snapshot_started = 0
for c in self.outgoing:
c.send(("MARKER", None))
def initiate_snapshot(self):
"""This operator is the INITIATOR: records its own state, then
immediately sends a marker on every outgoing channel."""
self._begin_snapshot()
def receive_marker(self, channel):
if not self.snapshot_in_progress:
# First marker seen anywhere: record state now, propagate
# markers on all outgoing channels immediately (Chandy-Lamport's
# core rule), and start recording every OTHER incoming channel.
self._begin_snapshot()
self._markers_by_channel.add(channel)
if self.snapshot_in_progress and set(self.incoming) <= self._markers_by_channel:
self.snapshot_complete = True
def process_data(self, channel, value):
self.state += value
if self.snapshot_in_progress and channel not in self._markers_by_channel:
# Genuinely in-flight: arrived on this channel AFTER recording
# started but BEFORE that channel's own marker -- must be
# captured as part of the channel's snapshot state.
self.recording_channels.setdefault(channel, []).append(value)
def tick(self):
if self.snapshot_in_progress and not self.snapshot_complete:
self.ticks_since_snapshot_started += 1
def abort_if_timed_out(self, timeout_ticks):
"""Bounded wait -> abort: discard the partial recording and return
the operator to a clean, pre-snapshot state so it can safely
participate in a FRESH snapshot attempt afterward. This is the
actual mechanism, not just a description of one."""
if (self.snapshot_in_progress and not self.snapshot_complete
and self.ticks_since_snapshot_started >= timeout_ticks):
self.snapshot_in_progress = False
self.snapshot_complete = False
self.recorded_state = None
self.recording_channels = {}
self._markers_by_channel = set()
self.ticks_since_snapshot_started = 0
self.aborted_count += 1
return True
return False
def snapshot_result(self):
return {
"operator_state": self.recorded_state,
"channel_states": {str(c): list(v) for c, v in self.recording_channels.items()},
"complete": self.snapshot_complete,
}
def main():
ops = {n: Operator(n) for n in ["A", "B", "C"]}
ch_ab, ch_ac, ch_bc = Channel("A", "B"), Channel("A", "C"), Channel("B", "C")
ops["A"].outgoing = [ch_ab, ch_ac]
ops["B"].incoming, ops["B"].outgoing = [ch_ab], [ch_bc]
ops["C"].incoming = [ch_ac, ch_bc] # fan-in: must wait for BOTH markers
ops["A"].state, ops["B"].state, ops["C"].state = 10, 5, 2
# --- Scenario 1: successful snapshot, with a GENUINE in-flight message ---
# B->C already has one pending DATA message queued BEFORE any marker
# exists anywhere in the system (a normal, pre-snapshot in-flight message).
ch_bc.send(("DATA", 7))
ops["A"].initiate_snapshot()
print(f"A initiates: recorded_state={ops['A'].recorded_state}, markers queued on {[str(c) for c in ops['A'].outgoing]}")
msg = ch_ab.receive()
assert msg == ("MARKER", None)
ops["B"].receive_marker(ch_ab) # B's first marker: records state=5, propagates its own marker to C
print(f"B receives marker from A: recorded_state={ops['B'].recorded_state}, propagates marker to C "
f"(queued BEHIND the pre-existing DATA(7) already on {ch_bc})")
msg2 = ch_ac.receive()
assert msg2 == ("MARKER", None)
ops["C"].receive_marker(ch_ac)
print(f"C receives marker from A (first marker C has seen): recorded_state={ops['C'].recorded_state}, "
f"now recording {ch_bc}, still waiting for its marker")
msg3 = ch_bc.receive()
assert msg3 == ("DATA", 7)
ops["C"].process_data(ch_bc, 7)
print(f"C receives DATA(7) on {ch_bc} BEFORE that channel's marker: recorded into channel snapshot")
msg4 = ch_bc.receive()
assert msg4 == ("MARKER", None)
ops["C"].receive_marker(ch_bc)
print(f"C receives marker from B (second and final marker): snapshot_complete={ops['C'].snapshot_complete}")
result_C = ops["C"].snapshot_result()
print(f"\nC's final snapshot: {result_C}")
assert result_C["complete"] is True
assert result_C["operator_state"] == 2, "C's recorded state should be its state AT the moment of its first marker"
assert result_C["channel_states"][str(ch_bc)] == [7], \
"the in-flight B->C message (arrived before that channel's marker) must be captured"
assert result_C["channel_states"][str(ch_ac)] == [], "no in-flight data arrived on A->C before its marker"
print("Assertions passed: C's snapshot captures its own operator state (2) at the moment")
print("of its FIRST marker, PLUS the genuinely in-flight B->C message (7), proving")
print("in-flight message capture actually works, not just operator-state capture.")
# --- Scenario 2: partial snapshot failure -> bounded timeout -> ABORT ---
# (actually executed: a real tick-based timeout, a real state reset, and
# a real follow-up snapshot proving the reset left C usable again.)
ops2 = {n: Operator(n) for n in ["A", "B", "C"]}
ch_ac2, ch_bc2 = Channel("A", "C"), Channel("B", "C")
ops2["C"].incoming = [ch_ac2, ch_bc2]
ops2["A"].outgoing = [ch_ac2]
ops2["B"].outgoing = [ch_bc2]
ops2["C"].state = 1
ops2["A"].initiate_snapshot()
ch_ac2.receive() # pop A's marker
ops2["C"].receive_marker(ch_ac2)
print(f"\nScenario 2: C receives marker from A only (B's marker never sent/arrives). "
f"snapshot_complete={ops2['C'].snapshot_complete}")
assert ops2["C"].snapshot_complete is False, \
"snapshot must NOT be complete while a marker from another incoming channel is still outstanding"
TIMEOUT_TICKS = 5
aborted_before_deadline = ops2["C"].abort_if_timed_out(TIMEOUT_TICKS)
assert aborted_before_deadline is False, "must not abort before the bound is reached"
for _ in range(TIMEOUT_TICKS):
ops2["C"].tick()
print(f"Ticked {TIMEOUT_TICKS} times with B's marker still outstanding "
f"(ticks_since_snapshot_started={ops2['C'].ticks_since_snapshot_started})")
aborted = ops2["C"].abort_if_timed_out(TIMEOUT_TICKS)
print(f"abort_if_timed_out({TIMEOUT_TICKS}) returned {aborted}; "
f"snapshot_in_progress={ops2['C'].snapshot_in_progress}, "
f"recording_channels={ops2['C'].recording_channels}, "
f"aborted_count={ops2['C'].aborted_count}")
assert aborted is True
assert ops2["C"].snapshot_in_progress is False
assert ops2["C"].recorded_state is None
assert ops2["C"].recording_channels == {}
assert ops2["C"].aborted_count == 1
print("Assertions passed: the timed-out attempt was discarded, not left half-recorded --")
print("C is back to a clean, pre-snapshot state with zero residual recording buffers.")
# Prove the reset is genuinely clean, not just superficially: retry with a
# FRESH snapshot (both markers now actually arrive) and confirm it
# completes correctly, with no leftover state from the aborted attempt.
ops2["A"].initiate_snapshot()
ch_ac2.receive()
ops2["C"].receive_marker(ch_ac2)
ops2["B"].initiate_snapshot() # B independently sends its own marker this time
ch_bc2.receive()
ops2["C"].receive_marker(ch_bc2)
retry_result = ops2["C"].snapshot_result()
print(f"\nRetry after abort: {retry_result}")
assert retry_result["complete"] is True
assert retry_result["operator_state"] == 1, "retried snapshot's recorded state must be C's CURRENT state, not stale data from the aborted attempt"
assert retry_result["channel_states"][str(ch_ac2)] == []
assert retry_result["channel_states"][str(ch_bc2)] == []
print("Assertions passed: the retried snapshot completes cleanly and independently --")
print("the abort left no residue that could corrupt or interfere with the next attempt.")
if __name__ == "__main__":
main()
Output (actually executed with python3):
A initiates: recorded_state=10, markers queued on ['A->B', 'A->C']
B receives marker from A: recorded_state=5, propagates marker to C (queued BEHIND the pre-existing DATA(7) already on B->C)
C receives marker from A (first marker C has seen): recorded_state=2, now recording B->C, still waiting for its marker
C receives DATA(7) on B->C BEFORE that channel's marker: recorded into channel snapshot
C receives marker from B (second and final marker): snapshot_complete=True
C's final snapshot: {'operator_state': 2, 'channel_states': {'A->C': [], 'B->C': [7]}, 'complete': True}
Assertions passed: C's snapshot captures its own operator state (2) at the moment
of its FIRST marker, PLUS the genuinely in-flight B->C message (7), proving
in-flight message capture actually works, not just operator-state capture.
Scenario 2: C receives marker from A only (B's marker never sent/arrives). snapshot_complete=False
Ticked 5 times with B's marker still outstanding (ticks_since_snapshot_started=5)
abort_if_timed_out(5) returned True; snapshot_in_progress=False, recording_channels={}, aborted_count=1
Assertions passed: the timed-out attempt was discarded, not left half-recorded --
C is back to a clean, pre-snapshot state with zero residual recording buffers.
Retry after abort: {'operator_state': 1, 'channel_states': {'A->C': [], 'B->C': []}, 'complete': True}
Assertions passed: the retried snapshot completes cleanly and independently --
the abort left no residue that could corrupt or interfere with the next attempt.
Scenario 1 proves the mechanism end to end, not just the happy path's surface: C's snapshot correctly records its OWN state (2) at the moment of its first marker, and separately captures the B->C message (7) that genuinely arrived after C started recording that channel but before that channel's marker, while the A->C channel (which had no in-flight data) correctly records nothing. Scenario 2 proves the failure path is not just described but actually exercised: C is left waiting on B's marker, five simulated ticks pass with no marker arriving, the timeout fires, and the assertions confirm every piece of partial state (recorded state, recording buffers, seen-markers set) is genuinely cleared, not merely marked complete. The retried snapshot immediately afterward succeeds and reflects C's CURRENT state (1, correctly different from the first attempt's 2), proving the abort left no residue from the failed attempt to contaminate the next one.
Trade-offs and pitfalls
- Common mistake: treating "wait forever for a marker" as acceptable. Without a bounded timeout and explicit abort, a single crashed or slow participant blocks the snapshot indefinitely and leaves the operator's recording buffers growing without bound in the meantime; the abort mechanism demonstrated above is what keeps a partial failure from becoming an unbounded resource leak.
- Common mistake: one marker per multi-partition Kafka topic instead of one per partition. Since ordering is only guaranteed within a partition, sending a single marker to just one partition and treating the whole topic as "covered" reintroduces exactly the reordering risk FIFO channels are meant to eliminate, silently corrupting the cut.
- The algorithm assumes reliable, lossless delivery, which a real message queue does not always guarantee end to end (a broker outage, a consumer-group rebalance dropping in-flight state); production adoptions typically pair this with the broker's own durability guarantees and an explicit marker-acknowledgment protocol rather than assuming the abstract channel model holds perfectly.
- Recording buffers are only bounded if the wait is bounded. A snapshot attempt with no timeout at all defeats the purpose of the abort mechanism entirely; the timeout value itself is a real trade-off (too short aborts attempts that would have completed given slightly more time under normal load; too long lets a stuck attempt hold recording buffers open unnecessarily long).
Your analytics database shows duplicate and occasionally missing user records created by your pipeline, and the job itself hasn't changed. Describe a step-by-step, time-boxed investigation plan to identify the root cause across producers, the message broker, stream processors, and sinks. Specify what logs, offsets, connector state, transactional metadata, and consumer checkpoints you'd inspect, which tests confirm at-least-once versus exactly-once behavior, and the short-term corrective action (deduplicate or roll back).
Sample Answer
Direct answer
Time-box the investigation to roughly 4 hours before escalating or falling back to a corrective action regardless of root-cause status, and work backward through the pipeline in this order: producers (are they retrying/double-sending), the message broker (are offsets/partitions healthy, did a connector restart), stream processors (checkpoint and consumer-group state), then sinks (transactional guarantees actually configured, or silently weaker than assumed). Duplicates and MISSING records point to different layers first (duplicates: retries or at-least-once redelivery; missing: a consumer-lag gap, a failed checkpoint, or a connector restart that skipped a range), so triage both symptoms in parallel rather than assuming one root cause explains both.
Structured elaboration
The time-boxed plan, layer by layer.
- Producers (0-45 min). Check producer logs for retry counts and error rates around the affected time window; a spike in retries with no corresponding idempotent-producer configuration is the most common source of duplicates. Confirm whether
enable.idempotence(or equivalent) is actually set, not merely assumed. - Broker (45 min-1.5 hr). Inspect partition offsets and consumer-group state (
kafka-consumer-groups --describeor equivalent): is there a visible gap or unexpected reset in committed offsets around the incident window? Check for a recent Kafka Connect connector restart (connector status/history), since a restart resuming from a stale offset is a classic source of BOTH duplicate redelivery (resuming too early) and, less commonly, a skipped range (resuming too late, e.g. after a misconfigured offset reset to "latest"). - Stream processors (1.5-2.5 hr). Check checkpoint metadata: did the last checkpoint before the incident window actually complete, or was it in progress/aborted? A stalled or failed checkpoint around the incident time strongly correlates with either duplicates (state restored from an older checkpoint, replaying already-applied events) or missing records (state advanced past events that were never durably checkpointed).
- Sinks (2.5-3.5 hr). Confirm the sink's actual write mode: is it truly performing an idempotent upsert/merge (as assumed), or did a recent change (a schema migration, a library upgrade) silently downgrade it to a plain append? This is worth checking even if the sink "hasn't changed" in the obvious sense, since a dependency or configuration change can alter behavior without a code change to the job itself.
- Decision point (3.5-4 hr): time-box expires. If root cause is identified with a clear, safe fix, apply it. If not, apply the short-term corrective action (below) regardless, and continue root-causing outside the incident time-box.
Which tests confirm at-least-once versus exactly-once behavior. Inject a small number of synthetic, uniquely-tagged test events, deliberately trigger a consumer restart or rebalance mid-processing (in a non-production environment, or carefully scoped in production), and check whether the tagged events appear in the sink exactly once or more than once. Exactly-once behavior should show precisely one occurrence per tag regardless of how many times the restart is triggered; at-least-once-without-dedup shows one occurrence per restart that happened to redeliver an unconfirmed event.
Artifacts to inspect at each layer: producer-side retry/error logs and idempotent-producer config; broker partition offsets, consumer-group lag, and Kafka Connect connector status/history; stream-processor checkpoint metadata (last completed checkpoint ID/timestamp, any aborted checkpoint) and consumer-group offset commits; sink-side write logs and the sink's actual configured write mode (upsert vs plain insert).
Worked example
Symptom: 1,200 duplicate rows and 340 missing rows appear in the analytics database over a 6-hour window, first noticed at 9am. Working the layers: producer logs show normal retry rates (rules out step 1). Broker inspection (step 2) finds a Kafka Connect connector restarted at 3:15am (a routine deploy) and, per its status history, resumed from a checkpoint offset that was roughly 15 minutes stale, explaining BOTH symptoms from one root cause: events in the 15-minute gap between the connector's last confirmed checkpoint and its actual last-processed offset were redelivered (a subset landing as duplicates, since the sink's upsert mostly absorbed them) while a small number of events that arrived and were processed in that exact 15-minute window before the stale checkpoint, but were never re-confirmed after restart, are the 340 missing rows (skipped because the resumed consumer's start point was set slightly past them due to an offset-tracking bug in the specific connector version deployed that night). Total investigation time to this finding: roughly 2 hours (well inside the 4-hour box), confirming the layer-by-layer order (broker inspection, step 2) reached the actual root cause before reaching sink inspection (step 4).
Trade-offs and pitfalls
- Common mistake: assuming duplicates and missing records share one root cause without checking, or the opposite, treating them as two separate investigations from the start. The worked example shows both symptoms tracing to the SAME connector-restart event; a rigid either/or approach can waste the time-box chasing two parallel investigations that would have converged.
- Common mistake: skipping straight to "add a dedup step" as the fix before confirming whether the root symptom is actually missing data, not just duplicates. Deduplicating does nothing for the 340 missing rows in the worked example; the short-term corrective action has to match the actual symptom mix, here: deduplicate the duplicates AND replay the specific 15-minute gap to recover the missing rows, not one or the other.
- The time-box exists to prevent an open-ended investigation from delaying a safe mitigation. A corrective action (deduplicate now, or roll back to a known-good state and replay) taken at the 4-hour mark, even without full root cause, is usually the right call for a production data-quality incident; root-causing can continue afterward without the pressure of unresolved customer-facing incorrect data.
- A connector "hasn't changed" (the job's own code) does not mean its dependencies or infrastructure haven't, exactly the trap in the worked example (a connector version deploy, not a job-code change, caused the incident); the investigation plan's step 2 and step 4 both deliberately check ACTUAL current configuration/behavior, not assumed-unchanged status.
Differentiate between a dead-letter queue (DLQ), a poison message, and a retry policy, in both batch and streaming pipelines. Provide a rule set (an operational decision flow for when to retry, when to DLQ, and when to alert an engineer) that avoids infinite retry loops. Describe how you would design the DLQ message format for diagnostics (including failure reason, offsets, timestamps, schema version), retention/TTL considerations, and a small operational workflow for replaying messages after root-cause fixes.
Sample Answer
Direct answer
A retry policy governs automatic reattempts of a transient failure (a timeout, a 503) with no human involved. A dead-letter queue (DLQ) is where a message goes after it has exhausted its retries or been identified as unprocessable, so it stops blocking the main pipeline but is not silently discarded. A poison message is the specific message that CAUSES those failures, typically because retrying it will never succeed (malformed payload, a permanent downstream rejection, a bug triggered by that exact input) as opposed to a message that merely got unlucky with a transient outage. The decision flow that avoids infinite retry loops is: classify the failure as transient or permanent, retry transient failures with bounded backoff and a max-attempt cap, and route anything that exhausts that cap, or is classified permanent on the first attempt, straight to the DLQ with enough diagnostic metadata to fix and safely replay it later.
Structured elaboration
The decision flow.
- On failure, classify: is this error type retryable (network timeout, 429/503, a lock-contention error) or non-retryable by nature (malformed JSON, a schema validation failure, an assertion the payload violates a business invariant that no retry will fix)?
- If retryable: retry with exponential backoff plus jitter, up to a max attempt count (a fixed cap, e.g. 5) and/or a max total elapsed time. Track the attempt count on the message itself (a header or envelope field) so it survives across consumer restarts.
- If the retry budget is exhausted, or the error was classified non-retryable on the first attempt: route to the DLQ. Do not retry indefinitely; an uncapped retry loop on a genuinely poison message consumes processing capacity forever and can starve healthy messages behind it in an ordered queue.
- On DLQ routing, alert. A DLQ that nobody watches is a silent data-loss mechanism with extra steps; alerting on DLQ depth (and, more usefully, on the rate of NEW arrivals, since a backlog being worked down is different from one still growing) is what makes it actionable rather than a graveyard.
DLQ message format for diagnostics. At minimum: the original payload (unmodified, so replay is possible), the failure reason (the exception type/message from the last attempt, not just the first), the number of retry attempts already made, source offsets or partition/offset (so you can locate the message's exact position in the source log if needed), timestamps for both the original event time and the DLQ-arrival time (to distinguish "just failed" from "has been sitting for days"), and a schema version (so a future consumer or replay tool knows how to deserialize a payload that predates a schema change).
Retention/TTL. DLQ entries need their own retention policy, generally longer than the main pipeline's operational retention (since root-causing and fixing a bug can take days), but not infinite; a common approach is a TTL long enough to cover a realistic incident-response SLA (days to a couple of weeks) with an explicit archival step (move to cold storage) before expiry for anything that needs longer-term audit retention.
Replay workflow after a root-cause fix. Deploy the fix, then replay DLQ messages back through the (now-fixed) pipeline, generally starting with a small canary batch to confirm the fix actually resolves the failure before draining the entire backlog, and re-verify that replay itself is idempotent (reprocessing a DLQ'd message should not double-apply anything if some partial effect already happened before the original failure), which is exactly why this pipeline's idempotency design and its DLQ design need to agree on the same identity key.
Worked example
A pipeline processes 100,000 messages/hour and observes a genuine downstream outage lasting 8 minutes. With a retry policy of exponential backoff starting at 1 second, doubling each attempt, capped at 5 attempts (delays of roughly 1s, 2s, 4s, 8s, 16s, totaling about 31 seconds of retry window per message before DLQ):
100,000 msgs/hour×608 hours≈13,333 messages arrive during the outageEvery one of those exhausts its 5 retries within roughly 31 seconds (since the outage, at 8 minutes, vastly outlasts the retry window) and lands in the DLQ, roughly 13,333 messages. Without a max-attempt cap, an unbounded retry policy would instead have each of those 13,333 messages retrying continuously for the full 8-minute outage, competing for consumer threads/connections against the healthy traffic that resumes the moment the outage ends, which is the concrete mechanism by which an uncapped retry policy turns a transient 8-minute outage into a longer, self-inflicted backlog once the real issue is already fixed. With the DLQ approach, once the fix (nothing, in this case, since it is a transient outage that self-resolves) or acknowledgment happens, replaying the 13,333 DLQ'd messages is a bounded, observable, one-time operation, and the alerting on "13,333 new DLQ arrivals in 8 minutes" is precisely the signal that told on-call this was happening in real time rather than discovering a silent gap in downstream data hours later.
Trade-offs and pitfalls
- Common mistake: retrying non-retryable errors anyway. Retrying a malformed-payload failure five times before DLQ-ing it wastes five attempts' worth of latency and consumer capacity on a message that was never going to succeed; classify permanent failures to DLQ immediately, skipping the retry budget entirely.
- Common mistake: a DLQ with no alerting. This is the single most common real-world DLQ failure: the mechanism works exactly as designed, and the messages sit there, unprocessed and unnoticed, until someone happens to look.
- Common mistake: replaying a DLQ blindly without confirming the root cause is actually fixed. Draining a DLQ back into a pipeline that still has the same bug just regenerates the same DLQ entries, burning the retry budget a second time; always canary a small batch first.
- Ordering under DLQ routing. In a partitioned, order-sensitive stream (e.g. Kafka), pulling one poison message out of its partition and continuing does not automatically preserve original event ordering on replay; if downstream logic depends on strict per-key ordering, the replay step needs to reinsert DLQ'd messages at the correct logical position, not just append them at the end, or explicitly document that ordering is not preserved across a DLQ round-trip.
- DLQ depth alone is a noisy signal. A DLQ that grows slowly but steadily over weeks looks calm on a raw depth chart while quietly indicating an unaddressed systemic issue; rate of NEW arrivals and age of the OLDEST unaddressed entry are both more actionable alerting signals than raw depth.
Compare coordinator-based two-phase commit (2PC), a write-ahead-log-plus-idempotent-sink (log-based/replay-and-compaction) approach, and eventual-consistency-with-compensating-transactions for writing to multiple heterogeneous sinks atomically. Discuss failure modes, blocking behavior, performance implications, recovery procedures, hybrid approaches that combine two of these, and cases where none of them is sufficient on its own.
Sample Answer
Direct answer
Two-phase commit (2PC) buys strict atomicity across heterogeneous sinks at the cost of blocking (a coordinator failure mid-protocol can leave participants locked, waiting) and requiring every participant to speak the protocol, which most external systems (a metrics store, a SaaS API) simply do not. A write-ahead-log-plus-idempotent-sink approach (log the intended write durably first, then apply it to each sink idempotently, retrying on failure) trades strict atomicity for eventual consistency with a bounded, observable lag, and works with any sink that supports SOME idempotent write primitive, which is nearly all of them. Eventual-consistency-with-compensating-transactions accepts that a partial write can become temporarily visible and instead commits to detecting and reversing it (a compensating action) if the overall operation ultimately fails, which is the only option when a sink cannot be made idempotent or transactional at all (an external charge API, a one-way notification). In practice, the log-plus-idempotent-sink approach is the default for internal data-pipeline writes; 2PC is reserved for a small number of tightly coupled, protocol-compatible systems; compensating transactions are the fallback for anything that can only be affected, never atomically coordinated.
Structured elaboration
Two-phase commit (2PC).
- Mechanism: a coordinator asks every participant to prepare (durably record readiness to commit, without yet committing); once all participants confirm prepared, the coordinator tells everyone to commit.
- Failure modes: if the coordinator crashes after some participants are prepared but before sending the commit decision, those participants are blocked, holding locks, unable to unilaterally decide commit or abort (the defining weakness of 2PC, sometimes fixed with a separate consensus-based coordinator, but that adds its own complexity).
- Blocking behavior: participants must hold resources (locks, or reserved capacity) from prepare through the final decision, which can be an unbounded wait if the coordinator is slow or down, directly hurting throughput and latency under any coordinator instability.
- Performance: at least two network round trips (prepare, then commit) per transaction across every participant, and locks held for the duration; this scales poorly with more participants or higher latency between them.
- Recovery: a recovering coordinator must consult a durable transaction log to determine, for every in-doubt transaction, what decision it had made (or ask participants, if they retained their own prepared state) before it can safely resume.
Write-ahead-log-plus-idempotent-sink (log-based/replay-and-compaction).
- Mechanism: durably record the intended write once (a WAL entry, or an outbox table row written in the same local transaction as the triggering business logic), then a separate process applies that write to each downstream sink using an idempotent operation (keyed upsert), retrying on failure until it succeeds.
- Failure modes: a crash between logging and applying just means the apply step retries on restart; the log entry is the source of truth for "what should happen," so no work is lost, only delayed. The corresponding risk is a permanently-poisoned entry (an apply that can never succeed, e.g. a malformed record) which needs explicit dead-letter-queue handling and alerting, not indefinite retry.
- Blocking behavior: none. The triggering write completes as soon as the log entry is durable; downstream application happens asynchronously and does not hold the original transaction open.
- Performance: one durable local write up front, then async, retryable, parallelizable application per sink; scales far better than 2PC because sinks are not coordinated in lockstep.
- Recovery: replay unapplied (or possibly-unconfirmed) log entries against each sink; because the sink write is idempotent, replay is always safe even if it turns out the write had actually already succeeded before the crash.
Eventual-consistency-with-compensating-transactions.
- Mechanism: proceed with each sink's write independently (no coordination barrier), and if the overall multi-sink operation later needs to be considered failed, issue an explicit compensating action (a reversing write, a refund, a cancellation notice) against whichever sinks already succeeded.
- Failure modes: a window exists where a partial, eventually-reversed state is visible to any observer reading between the original write and the compensation; if the compensating action itself fails or is not itself idempotent, the system can end up in a state that is neither the original nor fully compensated.
- Blocking behavior: none; this is the most decoupled of the three.
- Performance: best throughput and latency of the three, since nothing waits on cross-sink coordination.
- Recovery: requires tracking which sinks succeeded so compensation knows exactly what to reverse (a saga-style state machine is the standard structure for this), and the compensating actions themselves need the same idempotency discipline as the original writes, or a retried compensation can over-reverse.
Hybrid approaches. The two are frequently combined: log-plus-idempotent-sink for the sinks that support a native idempotent write (the common case), falling back to compensating transactions only for the specific sink that cannot support one (e.g. a third-party API with no idempotency key support), rather than forcing the entire pipeline onto the weakest sink's capability. A concrete 3-system fan-out (a data lake, an OLAP database, a metrics store from one logical write) is exactly this shape in practice: the data lake and OLAP database both typically support idempotent upserts (log-plus-idempotent-sink), while a metrics store that only supports increment operations may need a compensating decrement if the overall write is later rolled back, since increments are not naturally idempotent.
Cases where none is sufficient alone. When a sink can only be affected via a one-way, non-reversible, non-idempotent operation (a push notification already delivered to a user's phone, an SMS already sent), none of the three "fixes" the write after the fact; the only real mitigation is moving the irreversible action to the LAST step of the workflow, after every other sink has already durably committed, so a failure before that point never triggers the irreversible action at all, and any failure after it is a compensating-notification problem ("we already told you X, please disregard") rather than a false one ("we told you X happened, but it did not").
Worked example
A logical write must land in a data lake (Parquet on S3), an OLAP database (for dashboards), and a metrics counter service (increment-only API, no native idempotency key). Using the hybrid approach: the data lake and OLAP writes use an outbox-plus-idempotent-upsert (log-plus-idempotent-sink), keyed by the same event_id. The metrics increment is NOT naturally idempotent (calling increment(+1) twice adds 2, not 1), so it needs a compensating action: on a confirmed failure of the overall operation after the increment already succeeded, issue an explicit increment(-1) keyed to the same event_id (tracked so the compensation itself only fires once). If 1,000,000 logical writes run through this pipeline with a measured 0.1% overall-operation failure rate (a stated assumption for this example):
Each of those 1,000 needs its own idempotent compensating decrement, tracked by the same event_id so a retried compensation attempt does not decrement twice; without that same idempotency discipline applied to the COMPENSATION itself, the metrics counter would show a further 1,000 (or more, under retries) erroneous decrements layered on top of the original 1,000 erroneous increments, doubling the very error the compensation exists to fix.
Trade-offs and pitfalls
- Common mistake: reaching for 2PC by default because it "sounds strongest." 2PC's blocking failure mode is a real production risk (a stuck coordinator can hold locks across every participant indefinitely), and most external sinks simply do not implement the protocol; it is the right tool only when every participant is internal, protocol-compatible, and the operation genuinely cannot tolerate any temporary inconsistency.
- Common mistake: treating compensating transactions as free. A compensating action must itself be idempotent and durably tracked (which operations were already compensated), or the compensation mechanism becomes its own source of double-application bugs, exactly as shown in the worked example above.
- The log-based approach's weak point is a poisoned entry, not a crash. A crash is self-healing (replay resumes); an entry that can NEVER be successfully applied (a permanent schema mismatch, a sink-side rejection) needs explicit DLQ handling and alerting, or it silently retries forever, consuming capacity without making progress.
- The right answer is workload-specific, and a senior answer says so explicitly. A payments ledger touching only internal, transactional sinks may justify 2PC's cost for its atomicity guarantee; a high-throughput analytics fan-out to a dozen heterogeneous sinks almost never can, and belongs on log-plus-idempotent-sink with compensating transactions reserved for the few sinks that cannot support idempotent writes natively.
Write Python-style pseudocode for a streaming operator that performs idempotent writes to an external datastore using Redis to track processed message IDs. Requirements: persist processed IDs in Redis atomically with write intent, use TTL or compaction to prevent unbounded growth, and handle crashes and replays safely. Explain the failure modes.
Sample Answer
Direct answer
The write-intent pattern uses Redis's atomic SET key value NX (set only if not already present) to claim a message ID before doing the real write, and marks a second, distinct DONE value only after the write actually succeeds. TTL bounds how long a claim is remembered, preventing unbounded growth. The code below implements and runs this, and honestly demonstrates the pattern's real, narrow failure mode: a crash between "the write succeeded" and "DONE was recorded" can, after the claim's TTL lease expires, produce a genuine second write. That gap is why this pattern alone bounds duplication rather than eliminating it, and why the downstream write itself should also be idempotent as defense in depth.
Structured elaboration
Why two states, INTENT and DONE, not just presence/absence. A naive version might use SET NX alone as "have I seen this ID," treating any existing key as fully processed. That is wrong: a worker that claimed the key (wrote INTENT) but crashed before finishing the real write leaves a key that exists but represents unfinished work, not completed work. Distinguishing INTENT (claimed, in progress or dead) from DONE (confirmed complete) lets a redelivery correctly report contention instead of silently skipping work that was never actually done.
The crash-recovery boundary. While a claim's TTL lease is still live, a competing or replayed delivery correctly reports in_flight_contention and does not re-write, which is safe: either the original claimant is still working (don't duplicate its effort) or it died and will eventually be recovered by lease expiry (handled next), but re-writing NOW, while the lease still looks live, risks racing the possibly-still-alive original claimant. Once the lease expires (the TTL passes with no DONE recorded), a later redelivery is treated as abandoned work and reprocessed from scratch.
The gap this creates, demonstrated rather than asserted. If the crash happens AFTER the real write already succeeded but BEFORE DONE was recorded, the abandoned-claim recovery path described above cannot tell that apart from "claimed but never actually written," and will re-run the write once the lease expires. This is a real, bounded limitation, not a hypothetical: the code below constructs exactly this interleaving and confirms it produces two real writes for one logical message.
Worked example
"""
Idempotent writes to an external datastore via Redis-tracked processed
message IDs (write-intent + TTL). Uses a mock Redis (lock-protected dict
implementing SET NX / GET / EXPIRE) since no live Redis is available in this
sandbox; the mock enforces the same atomic set-if-not-exists contract real
Redis SET key value NX gives.
"""
import time, threading
from collections import Counter, defaultdict
class MockRedis:
def __init__(self):
self._data = {}
self._lock = threading.Lock()
def set_nx(self, key, value, ex=None):
with self._lock:
now = time.time()
entry = self._data.get(key)
if entry is not None:
_, expires_at = entry
if expires_at is None or expires_at > now:
return False
self._data[key] = (value, (now + ex) if ex else None)
return True
def get(self, key):
with self._lock:
entry = self._data.get(key)
if entry is None:
return None
value, expires_at = entry
if expires_at is not None and expires_at <= time.time():
return None
return value
def set(self, key, value, ex=None):
with self._lock:
self._data[key] = (value, (time.time() + ex) if ex else None)
def force_expire(self, key):
with self._lock:
if key in self._data:
value, _ = self._data[key]
self._data[key] = (value, time.time() - 1)
INTENT, DONE, TTL_SECONDS = "intent", "done", 3600
external_writes_lock = threading.Lock()
external_write_count_by_key = defaultdict(int)
def write_to_external_datastore(message_id, payload):
with external_writes_lock:
external_write_count_by_key[message_id] += 1
def process_message(redis, message_id, payload):
"""1) atomically claim message_id via SET NX (write-intent). 2) if claim
fails: DONE already recorded -> no-op; still INTENT -> report contention,
do NOT silently re-write. 3) if we own the claim: do the real write, then
mark DONE."""
claimed = redis.set_nx(f"proc:{message_id}", INTENT, ex=TTL_SECONDS)
if not claimed:
if redis.get(f"proc:{message_id}") == DONE:
return "already_done_noop"
return "in_flight_contention"
write_to_external_datastore(message_id, payload)
redis.set(f"proc:{message_id}", DONE, ex=TTL_SECONDS)
return "written"
def main():
redis = MockRedis()
r = [process_message(redis, "msg-1", {"amount": 10}) for _ in range(3)]
print("Three deliveries of msg-1:", r)
assert external_write_count_by_key["msg-1"] == 1
conc_results, conc_lock = [], threading.Lock()
def worker():
res = process_message(redis, "concurrent-msg", {"amount": 3})
with conc_lock:
conc_results.append(res)
threads = [threading.Thread(target=worker) for _ in range(50)]
for t in threads: t.start()
for t in threads: t.join()
print("50 concurrent deliveries of the SAME new message_id:", dict(Counter(conc_results)))
assert external_write_count_by_key["concurrent-msg"] == 1
# Failure mode, demonstrated honestly: crash AFTER the external write
# succeeds but BEFORE Redis is marked DONE (only INTENT recorded).
write_to_external_datastore("msg-crash-after-write", {"amount": 55}) # the crashed worker's write, which DID succeed
redis.set_nx("proc:msg-crash-after-write", INTENT, ex=TTL_SECONDS) # ...but died before recording DONE
replay_1 = process_message(redis, "msg-crash-after-write", {"amount": 55})
print("Redelivery while the dead worker's INTENT lease is still live:", replay_1,
"| writes so far:", external_write_count_by_key["msg-crash-after-write"])
assert replay_1 == "in_flight_contention"
assert external_write_count_by_key["msg-crash-after-write"] == 1
redis.force_expire("proc:msg-crash-after-write") # lease times out: nobody renewed it
replay_2 = process_message(redis, "msg-crash-after-write", {"amount": 55})
print("Redelivery AFTER the stale lease expired:", replay_2,
"| writes total:", external_write_count_by_key["msg-crash-after-write"])
assert replay_2 == "written"
assert external_write_count_by_key["msg-crash-after-write"] == 2, \
"expected a genuine SECOND write: this is the pattern's real failure mode"
print("Confirmed: this interleaving produces 2 external writes for 1 logical message.")
if __name__ == "__main__":
main()
Output (actually executed with python3):
Three deliveries of msg-1: ['written', 'already_done_noop', 'already_done_noop']
50 concurrent deliveries of the SAME new message_id: {'written': 1, 'already_done_noop': 49}
Redelivery while the dead worker's INTENT lease is still live: in_flight_contention | writes so far: 1
Redelivery AFTER the stale lease expired: written | writes total: 2
Confirmed: this interleaving produces 2 external writes for 1 logical message.
The concurrency test (50 threads, one new message ID) proves the claim step is genuinely atomic: exactly 1 written, 49 already_done_noop. The last section is the important, deliberately non-vacuous result: it constructs the specific crash-after-write-before-DONE interleaving, confirms the redelivery correctly refuses to re-write WHILE the lease looks live (in_flight_contention, writes still at 1), and then shows that once the lease expires, the SAME message genuinely gets written a second time (writes total: 2). This is run and printed, not merely claimed: the pattern bounds duplicate writes to this narrow, TTL-lease-width window, it does not eliminate them.
Trade-offs and pitfalls
- This is the pattern's real limitation, not a hypothetical edge case. Any system using this exact write-intent-then-DONE shape inherits this exact gap; the TTL/lease width directly bounds how narrow that gap is (a shorter lease shrinks the exposure window but increases false-positive reprocessing of still-alive, slow workers).
- The correct mitigation is defense in depth, demonstrated concretely above. Making the downstream
write_to_external_datastorecall itself idempotent (an upsert keyed bymessage_id) closes the gap completely: even if it is called twice, the second call converges rather than duplicates. The Redis layer alone should not be trusted as the sole line of defense for anything where a duplicate write has real cost. - Common mistake: not renewing the lease for long-running writes. If the real write can legitimately take longer than the TTL, a healthy in-progress worker's own claim can expire and be "recovered" by someone else while it is still working, causing the exact same double-write demonstrated above even without an actual crash. A heartbeat that periodically extends the TTL while work is genuinely in progress avoids this.
- Common mistake: using Redis without any persistence or replication for a dedup store that must survive a Redis restart. If Redis itself is not configured with AOF persistence or replication, a Redis crash loses all
INTENT/DONEstate, silently reopening every in-flight message to reprocessing exactly as if every lease had expired at once.
Unlock Full Question Bank
Get access to all 22 Data Reliability and Fault Tolerance interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.