Event-Driven Architecture and Asynchronous Messaging Questions
Designing systems around events and message passing: publish/subscribe, message queues, event streaming, choreography versus orchestration, and decoupling producers from consumers. Covers delivery semantics (at-least-once, at-most-once), ordering, backpressure, dead-letter handling, and the operational tradeoffs of asynchronous flows. Includes async processing patterns for offloading slow work.
Design a Dead Letter Queue (DLQ) processing workflow. Requirements: safe reprocessing of failed messages, visibility into failure reasons, quarantine for poison messages, and automation to replay or archive. Explain checks to run before re-enqueueing (idempotency, schema compatibility), and how to monitor DLQ health.
Sample Answer
Direct answer
Treat the dead-letter queue (DLQ) as a state machine, not just a holding queue: every quarantined message moves through explicit states, and nothing gets replayed until it passes two specific automated checks, has this exact message already been applied downstream (idempotency), and does its payload still match what the current consumer code expects (schema compatibility). Skipping either check on replay is the single most common way a "fixed" DLQ incident turns into a second incident.
Structured elaboration
Idempotency check before replay. A message can land in the DLQ after a later step in its processing failed, not the first one, for example a payment charge that succeeded but a subsequent confirmation-email step that did not. Replaying that message from the top without first checking whether it already partially or fully succeeded would double-apply the parts that already worked. Before replay, re-run the exact same dedup or idempotency lookup the normal delivery path would use, so a message that actually did succeed is recognized and skipped rather than blindly reprocessed.
Schema compatibility check before replay. A message can sit in the DLQ for days or weeks while the consumer code keeps evolving. Replaying it against the current code without checking whether its payload still matches the current expected schema risks a deserialization failure (landing right back in the DLQ) or, worse, being silently misinterpreted by code that has since changed its assumptions about a field's meaning. Validate the payload against the current schema before replay; a compatible message replays normally, an incompatible one gets transformed if a safe mapping exists, or archived with a note for manual handling rather than replayed blind.
Visibility into failure reasons. Every quarantined message should carry its classification (why it failed, and which state it is in) so a reviewer, or the automation itself, can act without re-diagnosing from scratch.
Automated replay or archive. Auto-classify at arrival (the same failure-reason tagging as the DLQ architecture question), auto-archive messages that pass a retention cutoff without being addressed, and auto-replay only for a narrowly pre-approved class of situations (for example, "downstream service X was down for known maintenance window Y, safe to replay everything quarantined during that window") with the idempotency and schema checks still applied even to auto-replays, not just manual ones. Anything outside that narrow, pre-approved class gates on human approval.
Monitoring DLQ health. Depth alone is not enough. Track the age of the oldest quarantined message (a slow leak looks fine on depth if replay keeps pace, age catches it), the replay success rate (a replay that lands right back in the DLQ signals the underlying fix did not actually work), and the arrival rate of new quarantines relative to a normal baseline.
stateDiagram-v2
[*] --> Quarantined: retries exhausted
Quarantined --> UnderReview: on-call or automation inspects
UnderReview --> SchemaCheck: candidate for replay
SchemaCheck --> IdempotencyCheck: schema compatible
SchemaCheck --> Archived: schema incompatible
IdempotencyCheck --> Replayed: safe to reprocess
IdempotencyCheck --> Archived: unsafe or already applied
Replayed --> [*]
Archived --> [*]
Worked example
An illustrative incident, numbers chosen for the walkthrough, not measured: a downstream service outage sends 200 messages to the DLQ. On recovery, the on-call engineer runs the automated checks before batch-replaying. Of the 200: 12 had actually already succeeded on a delayed retry that landed just before the DLQ escalation was processed, the idempotency check filters these out and marks them replayed-as-already-done rather than reprocessing them; 3 have a payload shape that predates a schema migration made two weeks earlier, these route to archive with a note for manual follow-up; the remaining 185 pass both checks and replay cleanly. As a sanity check on the walkthrough's own numbers: 12 + 3 + 185 = 200, accounting for the full batch.
Trade-offs and pitfalls
- Skipping the idempotency re-check on replay is the most common way a DLQ "fix" causes a new incident, double-processing something that had already partially succeeded.
- An auto-replay policy that is too broad (for example, "replay everything older than a fixed age" with no per-message check) reintroduces exactly the failure mode this workflow exists to prevent.
- Monitoring only depth misses a slow, steady quarantine growth that never actually gets worked down, age and replay-success-rate catch what depth alone cannot.
- Not capturing why a message was archived (versus replayed) leaves nothing for someone auditing the incident later to reconstruct the decision.
Define idempotency in the context of event-driven architectures. As a Solutions Architect, design a simple idempotency/deduplication strategy for an email-sending consumer that reads messages with payload {email_id, recipient, template}. Describe storage choices for dedup keys, TTL policies, memory vs disk trade-offs, and how to handle retries and long outage recovery.
Sample Answer
Direct answer
Idempotency in an event-driven system means processing the same message any number of times produces exactly the same effect as processing it once, which matters because message brokers commonly redeliver: a crash between processing and acknowledging, a network blip, or a producer retry can all cause the same logical message to arrive twice. For the email-sending consumer, that means checking a durable, atomic "have I already sent this email_id" record before every send, not trusting that the broker will only ever deliver a message once.
Structured elaboration
Dedup key
- Use
email_idas the deduplication key if oneemail_idalways corresponds to exactly one intended send; if a singleemail_idcould legitimately target multiple recipients, use(email_id, recipient)as a composite key instead. - Store a small record per key: status (in progress or sent) and when it was written.
Storage choices for dedup keys
- A fast in-memory store (a Redis-style key-value cache) for the hot path: low latency, high throughput, supports an atomic check-and-set operation needed to avoid a race between two consumer instances both trying to claim the same message.
- A durable, disk-backed store (a relational or key-value database) as the system of record: survives a full cache restart or a long outage, at higher latency than the in-memory tier.
- Hybrid: check the fast in-memory tier first; on a miss, fall back to the durable store before concluding the message is genuinely new, so a cache restart cannot cause a previously-sent email to be resent.
Time-to-live (TTL) policies
- The in-memory tier's entries expire after a TTL sized to the realistic redelivery window (long enough that any plausible retry or redelivery is still caught, short enough not to grow unbounded).
- The durable tier does not need a TTL for correctness (it is the long-term source of truth for "was this ever sent"), though it may have its own separate retention policy for storage cost or compliance reasons, independent of the dedup TTL.
Memory vs disk trade-offs
- Memory: fast, ideal for the common case (a redelivery arriving within seconds to minutes of the original), but limited capacity and lost on a restart unless backed by persistence.
- Disk: slower per lookup, but durable and cheap to retain for long windows, which is exactly what is needed for the "long outage recovery" requirement below.
- The two-tier design exists specifically to get the low latency of memory for the common case without losing correctness when memory alone is not enough.
Handling retries and long outage recovery
- Use an atomic check-and-set (an operation like Redis's
SETNX, or an equivalent conditional write) to claim a key before sending, so two concurrent consumer instances processing a redelivered pair of the same message cannot both decide they are the one to send it. - After a long outage (the in-memory cache is empty on restart, or entries aged out past their TTL during the outage), a redelivered message must not be trusted as new just because the fast cache has no record of it; falling back to the durable store catches exactly this case.
Worked example
The following script implements the two-tier dedup store (an in-memory "hot" tier with a TTL, backed by a durable tier with no TTL) and an idempotent consume function for the email-sending consumer, and was executed exactly as shown. A logical tick counter stands in for wall-clock time so TTL expiry is deterministic and reproducible, rather than asserting any real elapsed-time claim.
"""
Pinned, deterministic demo of a consumer-side dedup store for an email-sending
consumer, A logical tick counter
stands in for wall-clock time so TTL expiry is reproducible (no sleep()).
Run with: python3 idempotent_email_consumer.py
"""
from dataclasses import dataclass
from typing import Dict, Optional
@dataclass
class DedupRecord:
status: str # "in_progress" | "sent"
written_at_tick: int
class DedupStore:
"""Emulates a Redis-style store with SETNX (atomic check-and-set) and a
TTL measured in logical ticks, backed by a durable table for anything
that survives past the TTL window (the 'long outage recovery' path)."""
def __init__(self, ttl_ticks: int):
self.ttl_ticks = ttl_ticks
self._hot: Dict[str, DedupRecord] = {} # Redis-like cache
self._durable: Dict[str, DedupRecord] = {} # DB-backed, no TTL
def _expire_if_needed(self, key: str, now_tick: int):
rec = self._hot.get(key)
if rec and now_tick - rec.written_at_tick > self.ttl_ticks:
del self._hot[key]
def try_claim(self, key: str, now_tick: int) -> bool:
"""Atomic check-and-set: returns True if this call claimed the key
(i.e., no prior attempt is in flight or completed), False if a
duplicate delivery should be skipped."""
self._expire_if_needed(key, now_tick)
if key in self._hot:
return False # in_progress or sent, already claimed
if key in self._durable:
# survived a hot-cache eviction or full outage; still dedup
self._hot[key] = self._durable[key]
return False
self._hot[key] = DedupRecord(status="in_progress", written_at_tick=now_tick)
return True
def mark_sent(self, key: str, now_tick: int):
rec = DedupRecord(status="sent", written_at_tick=now_tick)
self._hot[key] = rec
self._durable[key] = rec # durable write-through, no TTL
def status(self, key: str) -> Optional[str]:
if key in self._hot:
return self._hot[key].status
if key in self._durable:
return self._durable[key].status
return None
class FakeEmailProvider:
def __init__(self):
self.sent_log = []
def send(self, email_id: str, recipient: str, template: str):
self.sent_log.append((email_id, recipient, template))
def consume(store: DedupStore, provider: FakeEmailProvider, message: dict, now_tick: int):
key = message["email_id"]
claimed = store.try_claim(key, now_tick)
if not claimed:
return "skipped_duplicate"
provider.send(message["email_id"], message["recipient"], message["template"])
store.mark_sent(key, now_tick)
return "sent"
if __name__ == "__main__":
TTL = 5 # ticks
store = DedupStore(ttl_ticks=TTL)
provider = FakeEmailProvider()
msg = {"email_id": "welcome-e100", "recipient": "a@example.com", "template": "welcome_v2"}
print("=== First delivery ===")
r1 = consume(store, provider, msg, now_tick=0)
print(f"result={r1}, provider.sent_log={provider.sent_log}")
assert r1 == "sent"
assert provider.sent_log == [("welcome-e100", "a@example.com", "welcome_v2")]
print("\n=== At-least-once redelivery of the SAME message, 1 tick later (well within TTL) ===")
r2 = consume(store, provider, msg, now_tick=1)
print(f"result={r2}, provider.sent_log={provider.sent_log}")
assert r2 == "skipped_duplicate"
assert len(provider.sent_log) == 1, "must not send twice"
print("\n=== Two consumers race on the SAME message at the SAME tick (concurrent redelivery) ===")
store2 = DedupStore(ttl_ticks=TTL)
provider2 = FakeEmailProvider()
claim_a = store2.try_claim("race-1", now_tick=0)
claim_b = store2.try_claim("race-1", now_tick=0)
print(f"consumer A claimed={claim_a}, consumer B claimed={claim_b}")
assert claim_a is True and claim_b is False, "only one consumer may claim the key"
print("\n=== Hot cache entry expires after TTL, but durable store still dedups (outage-recovery path) ===")
# Simulate a long outage: the in-memory/Redis-style cache is empty on
# restart (process restarted, or the key aged out past its TTL), but the
# durable table still has the record from before the outage.
after_outage_tick = 0 + TTL + 1 # past the TTL window
store._expire_if_needed("welcome-e100", after_outage_tick)
print(f"hot cache has key? {'welcome-e100' in store._hot}")
assert "welcome-e100" not in store._hot, "TTL should have evicted the hot entry"
r3 = consume(store, provider, msg, now_tick=after_outage_tick)
print(f"redelivery after TTL expiry + simulated outage -> result={r3}, provider.sent_log={provider.sent_log}")
assert r3 == "skipped_duplicate", "durable store must still catch the duplicate"
assert len(provider.sent_log) == 1, "still must not have sent twice"
print("\n=== A genuinely NEW message with a different email_id sends normally ===")
msg2 = {"email_id": "welcome-e101", "recipient": "b@example.com", "template": "welcome_v2"}
r4 = consume(store, provider, msg2, now_tick=after_outage_tick)
print(f"result={r4}, provider.sent_log={provider.sent_log}")
assert r4 == "sent"
assert len(provider.sent_log) == 2
print("\nALL ASSERTIONS PASSED")
Actual output from running python3 idempotent_email_consumer.py:
=== First delivery ===
result=sent, provider.sent_log=[('welcome-e100', 'a@example.com', 'welcome_v2')]
=== At-least-once redelivery of the SAME message, 1 tick later (well within TTL) ===
result=skipped_duplicate, provider.sent_log=[('welcome-e100', 'a@example.com', 'welcome_v2')]
=== Two consumers race on the SAME message at the SAME tick (concurrent redelivery) ===
consumer A claimed=True, consumer B claimed=False
=== Hot cache entry expires after TTL, but durable store still dedups (outage-recovery path) ===
hot cache has key? False
redelivery after TTL expiry + simulated outage -> result=skipped_duplicate, provider.sent_log=[('welcome-e100', 'a@example.com', 'welcome_v2')]
=== A genuinely NEW message with a different email_id sends normally ===
result=sent, provider.sent_log=[('welcome-e100', 'a@example.com', 'welcome_v2'), ('welcome-e101', 'b@example.com', 'welcome_v2')]
ALL ASSERTIONS PASSED
Walking the scenarios: the first delivery of welcome-e100 sends normally. A redelivery of the identical message one tick later (well within the 5-tick TTL) is skipped, and the send log still shows only one entry, proving no duplicate email went out. Two consumer instances racing on the same key at the same tick show only one successfully claims it (consumer A claimed=True, consumer B claimed=False), which is what the atomic check-and-set is for. After simulating the hot cache's TTL expiring (past tick 5) to represent a long outage, a redelivery of welcome-e100 is still correctly skipped, because the durable tier (no TTL) still has the record even though the fast tier does not, exactly the outage-recovery path the design requires. A genuinely new message (welcome-e101) sends normally and independently, confirming the dedup logic does not over-match on unrelated messages.
Trade-offs and pitfalls
- Common wrong turn: relying on the in-memory cache alone with no durable backing. It is fast for the common case, but a cache restart or an outage that outlasts the TTL causes exactly the duplicate-send failure idempotency was supposed to prevent.
- Common wrong turn: using a non-atomic "check, then set" (two separate operations) instead of a single atomic check-and-set. Between the check and the set, a second consumer instance can slip through and both end up believing they are the one to send.
- Common wrong turn: generating a fresh key on every delivery attempt (for example, from the broker's own delivery identifier) instead of using the message's own stable
email_id. A key tied to delivery mechanics changes on redelivery and stops deduplicating the thing that actually matters. - Senior signal: treating the TTL'd fast tier and the untimed durable tier as two different concerns (latency versus correctness-across-outages) rather than a single cache with an arbitrary expiry. The same two-tier idempotent-consumer pattern applies unchanged to other at-least-once consumers with a different natural key, for example a payment-notification consumer keyed by
payment_idinstead ofemail_id.
A customer needs an immutable audit trail and the ability to rebuild multiple read models quickly. Compare event sourcing + CQRS against a traditional relational database augmented with Change Data Capture (CDC). Discuss complexity, operational cost, replayability, schema evolution, developer ergonomics, and scenarios where event sourcing is or is not justified.
Sample Answer
Direct answer
For an immutable audit trail with fast rebuild of multiple read models, both event sourcing plus Command Query Responsibility Segregation (CQRS) and a relational database augmented with change-data-capture (CDC) can work, but they solve it differently: event sourcing stores domain intent as the source of truth and replays it deterministically, while CDC turns an existing relational database's row-level changes into a changelog that downstream consumers materialize into views. Pick event sourcing when the business needs the "why," not just the "what," and needs a canonical replayable log; pick relational-plus-CDC when the relational model already fits the domain and the team wants a lighter operational footprint.
Structured elaboration
Complexity
- Event sourcing + CQRS: higher architectural complexity. Domain events are the source of truth; the team must implement an append-only event store, event versioning, snapshotting, and one or more projections, and must handle at-least-once delivery into those projections.
- Relational + CDC: lower incremental complexity if the team already runs a relational database. CDC (typically reading the database's write-ahead log, for example via Debezium) streams committed transactional changes to downstream systems without changing the core domain model.
Operational cost
- Event sourcing needs operational maturity for the event store itself: scaling, backups, compaction, and snapshotting are all new operational surfaces, plus projection workers and a message bus.
- CDC leans on existing database tooling and log-shipping infrastructure; fewer bespoke components, generally lower operational overhead, but the pipeline is now coupled to the source database's internal log format and retention.
Replayability (folding the append-only-log / materialized-view nuance)
- Event sourcing's replay is native and deterministic: any read model or a brand-new projection can be rebuilt by replaying the event log from the start (or from a snapshot). The event log is intentionally an append-only log of domain intent, and every read model is a materialized view derived from it.
- CDC is structurally similar (the database's write-ahead log is also an append-only log, and CDC consumers building denormalized tables are also building materialized views), but the two differ in what the log records: CDC's changelog captures row-level state deltas ("column X became Y"), not business intent ("customer upgraded their plan"). Reconstructing intent from a stream of row deltas is lossy and sometimes ambiguous, and a full CDC-based rebuild requires retaining the source database's change history for as long as you might need to replay it, which is a weaker retention guarantee than an event store's log is designed to give.
Immutable audit trail and intent
- Event sourcing stores intent explicitly as first-class domain events; the audit trail directly answers "what business action happened and why."
- CDC provides a factual log of persisted state changes, useful for audit of what changed, but not necessarily why, unless the application already wrote that intent into the row (e.g., an explicit
reasoncolumn).
Schema evolution
- Event sourcing requires an explicit event-versioning strategy (adding fields with defaults, or introducing a new event type for a breaking change, plus upcasters that translate old event versions when replayed).
- CDC schema evolution tracks the source table's schema; adding a nullable column is generally safe, but column renames, type changes, or table restructuring can break the CDC pipeline's mapping and every downstream consumer of it.
Developer ergonomics
- Event sourcing has a steeper learning curve: developers must think in events, eventual consistency, and projection design, but gain very clear auditability and strong support for temporal questions ("what did this account look like on March 3rd?").
- CDC lets developers keep writing familiar create/read/update/delete (CRUD) code against the relational schema; projections consume the CDC stream with comparatively little domain-model rework.
When event sourcing + CQRS is justified
- A complex domain with rich auditability or regulatory replay requirements, where the business needs to reconstruct exact past system state.
- Many read models that change frequently and need fast, correct rebuilds.
- The business explicitly wants to capture intent, not just state, for analytics or downstream machine learning use.
When it is not justified
- A straightforward create/read/update/delete domain where the relational schema already models the business well and the database's own transaction log already satisfies audit and retention requirements.
- Limited engineering bandwidth or a need for fast delivery where the extra event-sourcing machinery would slow the team down for no corresponding benefit.
Worked example
A payments team is choosing between the two approaches for a ledger requiring 7 years of audit retention and 3 read models (customer statement, fraud-review queue, regulatory export).
- Event sourcing: the event store holds roughly 40 million domain events over 7 years (about 15,600 events/day on average for a mid-size ledger). A new read model, say a fourth "tax reporting" view added in year 5, is built by replaying those 40 million events once, deterministically, against the new projection logic; correctness is verifiable because the same event log produces the same output every time it is replayed.
- Relational + CDC: the same 7-year retention means either keeping 7 years of database write-ahead log history available to the CDC pipeline (expensive and often beyond what most databases retain by default) or accepting that a full historical rebuild of a new view is not actually possible from CDC alone, only from that point forward. This is the concrete cost of CDC's weaker replay guarantee versus a purpose-built event store: it shows up exactly when the business asks for a new view of old data.
Trade-offs and pitfalls
- Do not choose event sourcing for its audit-trail marketing value alone; a relational database with proper CDC and immutable audit columns can satisfy many audit requirements at a fraction of the operational cost.
- Do not underestimate CDC's replay ceiling: if "rebuild any read model from any point in history" is a hard requirement, verify the source database's log retention actually supports it before committing to CDC as the long-term answer.
- Senior signal: naming the retention and rebuild requirement in concrete terms (how far back, how many read models, how often they change) before picking a side, rather than treating this as a purely stylistic architecture preference.
How would you choose a Kafka partitioning key for a user-events topic such that ordering is preserved per user but partitions are balanced across the cluster? Discuss hashing strategies, handling hot users, and approaches for multi-tenant fairness when some tenants generate far more events than others.
Sample Answer
Direct answer
Use user_id itself as the partition key so Kafka's partitioner (a mod-based hash of the key) always routes every event for a given user to the same partition, which is what gives per-user ordering, it falls directly out of using the same key for the same entity, not something built separately. Balance across the cluster then comes from having enough distinct keys and a reasonably uniform hash, with a deliberate, separate remediation for the minority of disproportionately hot keys.
Structured elaboration
Hashing strategy. Kafka's default partitioner computes hash(key) % num_partitions (using a murmur2-based hash internally); the exact hash function matters less than the property that it is deterministic, the same key always maps to the same partition number for a fixed partition count. Because the partition assignment is a pure function of the key, ordering per key is automatic once you commit to using that key consistently, not an extra mechanism you build on top.
Why "balanced" is a statistical property, not a guarantee. With many distinct user ids and a well-distributed hash, the law of large numbers makes total load spread roughly evenly across partitions, if every key generates traffic at roughly the same rate. The question specifically names the case where that assumption breaks: hot users.
Handling hot users. Two remediations, with different costs:
- Key salting (sub-partitioning): append a bounded suffix to a hot user's key only once their rate crosses a threshold (
user_id#0throughuser_id#K-1), spreading that one user's traffic across up to K partitions. The explicit cost: per-user ordering across the salted shards is no longer guaranteed, only ordering within each shard is. This is a real trade-off to name out loud, not a free win, if the consuming system genuinely needs strict per-user order even for hot users, salting defeats that and a downstream re-sequencing step would be needed, which mostly defeats the purpose of salting in the first place. - Isolating known hot keys: if hot users are identifiable ahead of time (service accounts, bots, high-volume enterprise integrations), route them to a separate topic or partition set with its own capacity plan, keeping the main topic's "roughly uniform" assumption valid for the long tail of ordinary users.
Multi-tenant fairness. Partitioning purely by user_id does nothing to stop one tenant (an account with many users) from occupying a disproportionate share of partitions and consumer capacity relative to a smaller tenant, since partition assignment has no concept of tenant at all. If fairness across tenants is a hard requirement: give large or noisy tenants a dedicated topic or partition range so they cannot starve others' consumer lag budget, and layer per-tenant rate limiting or quota enforcement on the consumer or producer side (the same token-bucket idea used for rate limiting generally), since a static partition assignment cannot adapt to traffic that changes over time the way an explicit quota can.
Worked example
A small, executed illustration of the mechanism (using Python's built-in zlib.crc32 purely as a stand-in deterministic hash function to demonstrate the mod-based mechanic; Kafka's actual default partitioner uses a different, murmur2-based hash internally, the mechanism, not the specific hash, is what this illustrates).
import zlib
num_partitions = 6
for uid in ["user-1042", "user-8831", "user-2207", "user-77519"]:
h = zlib.crc32(uid.encode())
print(uid, "partition=", h % num_partitions)
user-1042 partition= 3
user-8831 partition= 4
user-2207 partition= 4
user-77519 partition= 0
Every event for user-1042 always lands on partition 3, giving that user strict per-user order for free. Now simulate salting the hot user user-8831 into 12 sub-keys (user-8831#0 through user-8831#11):
counts = {}
for shard in range(12):
key = f"user-8831#{shard}"
p = zlib.crc32(key.encode()) % num_partitions
counts[p] = counts.get(p, 0) + 1
print(counts)
{4: 5, 0: 1, 5: 5, 1: 1}
The 12 salted sub-keys land on 4 of the 6 partitions instead of all landing on partition 4 as the unsalted key would. The spread is not perfectly even at this small sample size (partitions 4 and 5 got 5 each, partitions 0 and 1 got 1 each), which is expected statistical noise at n = 12, not a flaw in the technique; at real production volumes for a genuinely hot key, the same law of large numbers that balances ordinary keys also evens out a hot key's salted shards.
Trade-offs and pitfalls
- Salting trades away strict per-user order for the salted user specifically, in exchange for spreading their load; state this cost explicitly when proposing it, do not present it as a free fix.
- Changing the partition count later reshuffles the mod-based assignment for every existing key, not just new ones, since
hash(key) % num_partitionschanges for essentially every key whennum_partitionschanges. Under-provisioning partition count up front makes a later increase a disruptive, coordinated migration, not an incremental capacity add. - "Balanced across the cluster" and "fair across tenants" are different guarantees: the first is about hash uniformity over many keys, the second is about the business meaning of who owns which keys, no hashing scheme fixes tenant fairness on its own.
- Picking the wrong entity granularity for the key, for example partitioning by
session_idwhen the actual ordering requirement is per-user across sessions, silently breaks the guarantee the system was supposed to provide.
Design an idempotent consumer for processing payment events from a message queue. Message format: {payment_id, user_id, amount, currency, timestamp}. Requirements: prevent double-charges on redelivery, support at-least-once delivery semantics from the broker, allow retries, and maintain low latency. Describe the deduplication store, choices between in-memory, Redis, or relational DB, TTL strategy, and cleanup considerations.
Sample Answer
Direct answer
Derive the deduplication (dedup) key from the payment's own business identity, payment_id, not from anything the broker generates, then make claiming that key and processing the payment a single atomic step: check-and-claim, process, mark done. Redelivery of the same message must always produce the same payment_id, so a redelivered message collapses onto the same key and gets skipped instead of charged again.
Structured elaboration
Why payment_id, not the other fields. The message is {payment_id, user_id, amount, currency, timestamp}. user_id alone is wrong (one user can legitimately make two separate payments). (user_id, amount, timestamp) is also wrong: two real, distinct charges can share the same amount, and timestamp is set by the producer, so a retried publish can carry a different timestamp for the same logical payment while a coincidentally identical tuple could describe two different real charges. payment_id is the one field that is supposed to be stable across every redelivery of the same logical event, which is exactly the property a dedup key needs.
The atomicity requirement. "At-least-once delivery from the broker" means the same message will arrive more than once, by design, not as an edge case. If the consumer does exists = store.get(payment_id) followed by a separate store.set(payment_id, ...), two redeliveries handled concurrently (two consumer threads, or two instances in the same consumer group during a rebalance) can both read "not found" before either writes, and both charge. The claim has to be a single atomic operation: "insert this key only if it does not already exist," succeed-or-fail in one round trip.
Store choice: in-memory vs. Redis vs. relational database.
| Store | Atomic claim primitive | Round-trip cost | Durability across restarts | Shared across consumer instances | Best fit here |
|---|---|---|---|---|---|
| In-memory (process-local dict or least-recently-used, LRU, cache) | Check-then-insert under a single process lock | Lowest | None, lost on restart | No | A fast pre-check layer in front of a durable store, never the sole source of truth |
| Redis | SET key value NX PX ttl_ms (atomic in one command) | Low, one round trip | Best-effort unless persistence is explicitly enabled | Yes | Primary dedup store when "maintain low latency" is a hard requirement |
| Relational database | INSERT ... ON CONFLICT DO NOTHING or a unique constraint on payment_id | Higher, a transactional write | Full, same guarantees as the ledger | Yes | When the claim and the actual charge write must commit as one atomic transaction |
For a payment consumer, Redis is usually the right default for the hot path because "maintain low latency" is an explicit requirement and the claim is a single command. The relational option earns its extra latency specifically when you want the dedup row and the ledger write to be atomically consistent with each other, for example inserting the dedup row in the same transaction as the row that records the charge, so a crash between "claim" and "charge" cannot leave the dedup store and the ledger disagreeing.
Probabilistic pre-filter. At high volume, a Bloom or cuckoo filter can sit in front of the authoritative store as a cheap "have I possibly seen this before" check: a negative result is certain (definitely new, skip the store lookup and go straight to processing), a positive result is only "maybe," and must still fall through to the authoritative store, because treating a false positive as certain would silently drop a legitimate first-time charge. A cuckoo filter additionally supports deletion, which matters if you want the filter itself to track a rolling window rather than growing forever. Either filter is an optimization that reduces load on the durable store for the (usually large) fraction of genuinely-new events; it never replaces the store as the source of truth.
TTL strategy. The dedup key's time-to-live (TTL) must be at least as long as the broker's own maximum redelivery window plus a safety margin. If the broker can, in the worst case, redeliver a message up to 12 hours after first delivery (for example, a long consumer outage followed by catch-up), a TTL shorter than that lets the key expire and a legitimate-looking "new" charge slip through on a very late, otherwise-normal redelivery. TTL should not be sized off how long you want to retain history; it should be sized off how long the broker can still legally hand you a duplicate.
Cleanup considerations. Redis expires keys natively via the TTL, no separate job needed. A relational dedup table has no native per-row expiry in most databases, so it needs an explicit periodic sweep (a scheduled DELETE WHERE created_at < now() - retention_window, i.e. a dedup-table-plus-TTL pattern implemented as a cron job rather than a database feature), or the table grows without bound and both storage cost and index lookup latency degrade over time.
flowchart TD
A[Payment event arrives] --> B{Dedup key exists in store?}
B -- Yes --> C[Skip: return prior result, no charge]
B -- No --> D[Atomically claim key: SETNX with TTL]
D --> E[Process payment]
E --> F[Mark key as done, keep TTL]
F --> G[Ack message to broker]
C --> G
Worked example
Assume 500 payment events per second sustained, and a 24 hour TTL matching the broker's stated maximum redelivery window.
live keys=rate×TTL=500×86,400=43,200,000At roughly 150 bytes per key plus value plus Redis's own per-key overhead:
bytes=43,200,000×150B≈6.48GBThat is a meaningful working set for a single Redis instance. It is exactly the kind of number that motivates either sharding the dedup store across a Redis cluster keyed by payment_id, shortening the TTL if the broker's actual redelivery window is smaller than assumed, or moving completed, already-settled payments' dedup records into the cheaper relational table once they age past the point where redelivery is still plausible.
Trade-offs and pitfalls
- Using a coarse or derived key (amount, timestamp, or a hash of the whole payload) instead of
payment_idis the single most common mistake: it either misses real duplicates (payload varies slightly across redeliveries, for example a re-serialized timestamp) or falsely merges two distinct legitimate charges. - A GET-then-SET pattern instead of a single atomic command reintroduces the exact race the dedup layer exists to prevent; this only shows up under real concurrent redelivery, so it can pass casual testing and still double-charge in production.
- Redis becoming unavailable forces an explicit choice: fail closed (stop processing new payments until Redis is back, safe but reduces availability) or fail open (keep processing without dedup protection, risks double-charges). Silently defaulting to one without deciding on purpose is itself the pitfall.
- Sizing the relational table's cleanup job off total historical volume instead of the broker's actual redelivery window wildly over-retains data and slows the unique-constraint index down for no correctness benefit.
Unlock Full Question Bank
Get access to all 25 Event-Driven Architecture and Asynchronous Messaging interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.