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.
A client asks you to recommend between Kafka, Amazon SQS, and RabbitMQ given requirements: durable event storage, long retention, consumer replay, low-latency processing, and moderate ops complexity. Compare the three in terms of ordering guarantees, retention/replay, scaling model, operational cost, and typical use-cases (audit log, task queue, pub/sub notification).
Sample Answer
Direct answer
Given durable event storage, long retention, and consumer replay as hard requirements, Kafka (an open-source distributed event-streaming platform built around a durable, partitioned, replayable log) is the strongest fit of the three, because Amazon SQS (Simple Queue Service, a fully managed point-to-point queue) and RabbitMQ do not treat "replay the same history from an arbitrary point" as a first-class capability. The trade-off to name explicitly is that Kafka's strength here comes with more operational complexity than SQS and a different latency/throughput profile than RabbitMQ, so the recommendation should acknowledge that moderate ops complexity is genuinely in tension with the retention/replay requirement, not free.
Structured elaboration
Compare all three against the five named axes:
| Axis | Kafka | Amazon SQS | RabbitMQ |
|---|---|---|---|
| Ordering guarantees | Strict order within a partition (messages with the same key land on the same partition and are read in order); no global order across partitions. | Standard queues: no ordering guarantee. FIFO (first-in-first-out) queues: strict order per message group, at reduced throughput. | Per-queue order is preserved for a single consumer; with multiple competing consumers, per-message order across consumers is not preserved unless deliberately partitioned by key. |
| Retention/replay | Configurable retention (commonly days to indefinite), consumers track their own offset and can rewind to reprocess history; this is the platform's defining feature. | Messages are deleted once acknowledged (or after the retention window if unconsumed, max 14 days); no concept of replaying already-consumed messages. | Messages are removed once acknowledged; no built-in replay of consumed messages (a separate durable log would be needed alongside it). |
| Scaling model | Scales by adding partitions and brokers; throughput scales with partition count, at the cost of needing to plan partition count and key distribution up front. | Scales automatically and transparently; AWS (Amazon Web Services) operates the scaling, no partition planning required. | Scales by adding queues/nodes and clustering; requires more manual cluster management to scale horizontally than SQS. |
| Operational cost | Highest if self-managed (broker fleet, cluster coordination/metadata layer, monitoring); substantially lower if using a managed offering (e.g., Amazon MSK or Confluent Cloud), but still generally pricier and more operationally involved than SQS. | Lowest: fully managed, pay-per-request, effectively zero operational burden. | Moderate: typically self-hosted or managed via a vendor; less operational surface than a self-run Kafka cluster but more than SQS. |
| Typical use case | Audit log (durable, replayable history of everything that happened); analytics event backbone. | Task queue (discrete units of work, e.g. a thumbnail-generation job) where you want zero-ops simplicity. | Pub/sub notification and flexible routing (exchange-based routing to multiple queues) where you want more routing flexibility than SQS offers without running Kafka. |
Worked example
Given the stated requirements (durable event storage, long retention, consumer replay, low-latency processing, moderate ops complexity), map each requirement to what it rules out: "long retention + consumer replay" rules out SQS and RabbitMQ as the primary event store, since neither is designed to let a new consumer rewind and reprocess history that has already been acknowledged by a different consumer. That leaves Kafka, but "moderate ops complexity" is now in direct tension with a self-managed Kafka cluster, so the concrete recommendation is a managed Kafka offering (e.g., Amazon MSK or Confluent Cloud) rather than self-hosted Kafka: it satisfies durability, retention, and replay natively, and shifts broker operations (patching, scaling the cluster, replication management) to the vendor, bringing operational cost back down toward "moderate." Low-latency processing is achievable on Kafka (consumers read from the log with millisecond-scale poll latency once caught up), but the design should note that "low latency" here refers to steady-state consumer read latency, not a guarantee about time-to-first-byte during a cold-start replay of a large retained history, which is inherently slower by the size of the backlog being replayed.
Trade-offs and pitfalls
The main pitfall is treating "moderate ops complexity" and "long retention with replay" as independently satisfiable without acknowledging the tension between them; a client that hears "yes to all five requirements" without the managed-Kafka caveat will be surprised by the operational reality of running Kafka themselves. A second pitfall is choosing SQS FIFO queues as a workaround for the ordering/replay requirements, because FIFO addresses ordering within a message group but still does not provide replay of consumed messages, so it does not actually satisfy the stated requirement even though "FIFO" sounds like the right keyword. A senior answer states the recommendation as a specific product choice (managed Kafka) rather than just "Kafka" in the abstract, since the operational-cost axis changes materially between self-hosted and managed.
Implement (in Python) an HTTP consumer handler that receives JSON POST events with a unique 'event_id'. Use Redis to ensure idempotent processing by atomically marking an event as processed and only invoking the handler once per event_id. Show the atomic check-and-set logic and TTL handling to prevent unbounded state.
Sample Answer
Direct answer
The HTTP handler needs one atomic Redis command that both checks whether event_id has been processed and claims it in the same round trip: SET key value NX PX ttl_ms. NX means "only set if the key does not already exist," which is the compare-and-set; PX attaches an expiry in milliseconds so a claim from a crashed request eventually self-heals instead of blocking that event_id forever.
Structured elaboration
Approach. On each POST, build the Redis key as processed:{event_id}. Attempt SET key "processing" NX PX ttl_ms. If it returns success, this call owns the event: run the handler, then overwrite the value to "done" (refreshing the same TTL). If it returns failure, the key already exists, meaning either another request already handled this event_id or one is in flight right now, so this call returns immediately without invoking the handler.
import fakeredis # in-memory, wire-compatible with redis-py; swap for a real
# Redis client in production, the API is identical
from flask import Flask, request, jsonify
r = fakeredis.FakeStrictRedis(decode_responses=True)
app = Flask(__name__)
def business_handler(payload: dict):
"""Stand-in for the real per-event side effect (e.g. update a record,
send a notification). Idempotency is enforced below it, not inside it."""
...
def claim_and_handle(r, event_id: str, payload: dict, handler, ttl_seconds: int = 86400) -> bool:
"""
Claims event_id using a single atomic SET key value NX PX ttl.
NX = only set if the key does not already exist (the compare-and-set).
PX = expiry in milliseconds, so a stuck/never-completed claim self-heals.
Returns True if THIS call ran the handler, False if it was a duplicate.
"""
key = f"processed:{event_id}"
claimed = r.set(key, "processing", nx=True, px=ttl_seconds * 1000)
if not claimed:
return False
handler(payload)
r.set(key, "done", px=ttl_seconds * 1000) # refresh marker after success
return True
@app.post("/events")
def receive_event():
"""The actual HTTP consumer handler: parses the JSON POST body, pulls
out event_id, and delegates to the same atomic claim-and-handle logic."""
body = request.get_json(force=True)
event_id = body["event_id"]
ran = claim_and_handle(r, event_id, body, business_handler, ttl_seconds=60)
status = "processed" if ran else "duplicate_skipped"
return jsonify({"status": status, "event_id": event_id}), 200
Key points.
NXis what makes this a true compare-and-set: a plainGETfollowed by a plainSETwould let two concurrent requests both observe "not found" before either writes, and both would run the handler.- The TTL is reapplied after success so the marker's lifetime is measured from the last state change, not just the original claim.
- The handler is invoked exactly once inside the branch that won the claim; every other call for the same
event_id, whether truly concurrent or a much later HTTP retry, returnsFalsewithout side effects.
Worked example
Executed against fakeredis (an in-memory server that speaks the same protocol as redis-py, used here purely so the demo needs no running Redis server; the client-facing API is identical to a real deployment), driven two ways: first hitting the HTTP layer directly through Flask's test client with real JSON POST requests, then a lower-level run isolating the claim function's TTL behavior.
Through the actual HTTP endpoint, posting the same JSON body three times (simulating a client retry plus a broker redelivery):
resp1: 200 {'event_id': 'evt-http-1', 'status': 'processed'}
resp2: 200 {'event_id': 'evt-http-1', 'status': 'duplicate_skipped'}
resp3: 200 {'event_id': 'evt-http-1', 'status': 'duplicate_skipped'}
business handler invocation count: 1
Isolating the claim function's own TTL behavior directly (bypassing HTTP to control timing precisely):
first call ran handler : True
second call ran handler : False
third call ran handler : False
handler invocation count : 1
stored marker value : done
ttl seconds remaining : 60
key exists before expiry : True
key exists after 1.2s : False
handler ran after expiry : True
The three POSTs to /events with the identical event_id return 200 every time (an HTTP consumer should not error on a duplicate, it should acknowledge and no-op), but the business handler runs exactly once, confirmed by the invocation counter. The second run isolates the claim function alone: the first call claims the key and runs the handler once; the second and third calls (simulated redeliveries of the same event_id) both see the key already present and skip the handler, so handler invocation count stays at 1. Its last four lines use a separate 1 second TTL: the key exists immediately after the claim, is gone 1.2 seconds later, and a fourth call after expiry runs the handler again, proving the TTL is what prevents the dedup store from growing without bound rather than leaking claimed keys forever.
Trade-offs and pitfalls
Complexity. One Redis round trip per event, O(1) time. Memory is O(number of live keys), which is bounded by arrival_rate x TTL rather than by total historical event count, since expired keys are reclaimed automatically.
Edge cases.
- Slow handler vs. TTL: if the handler can run longer than the TTL, a second delivery arriving mid-processing will see the key expired and re-run the handler concurrently with the first, unclaimed run. Either extend the TTL past the handler's worst-case duration, or refresh the TTL periodically while the handler is still running.
- Crash after claim, before completion: if the process crashes after
SET ... NXsucceeds but before the handler finishes, the key stays"processing"until it expires, at which point a redelivery will legitimately re-run the handler. This means the handler itself should be safe to run from the start on a redelivery (its own side effects should be idempotent or wrapped in their own transaction), a pure lock cannot fully substitute for that. - Two events with the same
event_idbut different payloads: the design here trustsevent_idas the sole identity; if the same id can legitimately carry different content it is not a valid dedup key, that would need a content hash bound into the key or the schema fixed upstream. - Redis unavailability: the handler has no fallback path here; a production version needs an explicit decision (reject the request, or accept the risk of duplicate processing) when the
SETcall itself fails, not just when it returns a normal "already claimed" result.
Implement an exponential backoff retry strategy as a middleware for a Node.js message consumer. The middleware should support max retries, jitter, and configurable base/backoff multipliers. Provide code or clear pseudocode showing retry logic, how failures are bubbled to DLQ after max retries, and where to insert idempotency checks.
Sample Answer
Direct answer
Structure the middleware as a loop around the handler: an idempotency check runs first (before any retry bookkeeping, so a redelivery of an already-processed message never consumes a retry slot), then on each failure compute a capped exponential delay and sleep a random amount between zero and that cap (full jitter), and once the configured maximum retries is exceeded, push the message to a dead-letter queue (DLQ) and stop.
Structured elaboration
Where the idempotency check belongs. It runs before the retry loop, not inside it. A message redelivered after it already succeeded should short-circuit immediately, it should not be treated as attempt 1 of a fresh retry sequence, and it should not re-run the handler's side effects.
The retry loop. On each failure: if the attempt count exceeds maxRetries, escalate to the DLQ and stop; otherwise compute delay = min(maxDelayMs, baseMs * multiplier ** (attempt - 1)) and sleep a uniformly random value between 0 and that delay (full jitter) before the next attempt.
Configurable parameters. baseMs and multiplier control how fast the delay grows per attempt; maxDelayMs caps it so a message near its retry limit does not wait an unreasonably long time; maxRetries bounds total attempts before DLQ escalation.
function makeBackoffMiddleware({ baseMs, multiplier, maxRetries, maxDelayMs, rng, sendToDLQ, idempotencyCheck }) {
return async function process(message, handler) {
// Idempotency check runs BEFORE any retry bookkeeping: a redelivered
// message already committed downstream short-circuits here instead of
// consuming a retry slot or re-running side effects.
if (await idempotencyCheck(message)) {
return { status: 'skipped-duplicate', attempts: 0, delaysMs: [] };
}
let attempt = 0;
const delaysMs = [];
for (;;) {
attempt += 1;
try {
await handler(message);
return { status: 'processed', attempts: attempt, delaysMs };
} catch (err) {
if (attempt > maxRetries) {
await sendToDLQ(message, err, attempt - 1);
return { status: 'dlq', attempts: attempt - 1, delaysMs, error: err.message };
}
// Full jitter: sleep a UNIFORM random value between 0 and the capped
// exponential delay, not the exponential value itself, this is what
// prevents every instance that failed at the same moment from
// retrying in lockstep (the thundering-herd failure mode).
const capped = Math.min(maxDelayMs, baseMs * multiplier ** (attempt - 1));
const jitterMs = Math.round(capped * rng());
delaysMs.push(jitterMs);
// production code awaits a real timer here: await sleep(jitterMs)
}
}
};
}
Worked example
Executed with Node.js, baseMs = 100, multiplier = 2, maxRetries = 4, maxDelayMs = 5000, and a seeded pseudo-random number generator (mulberry32, seeded with 42) standing in for Math.random() purely so the printed jitter values are reproducible for this demo; production code would use Math.random() or crypto.randomInt() instead.
case1 (recovers within budget): {"status":"processed","attempts":4,"delaysMs":[60,90,341]}
case2 (exhausts retries -> DLQ): {"status":"dlq","attempts":4,"delaysMs":[67,35,211,219],"error":"downstream-500"}
dlq contents: [{"id":"evt-2","attempts":4,"error":"downstream-500"}]
case3 (duplicate redelivery skipped): {"status":"skipped-duplicate","attempts":0,"delaysMs":[]}
Case 1: a handler that fails 3 times then succeeds recovers on attempt 4, within the 4-retry budget, with 3 recorded jitter delays (one per failure before the successful attempt). Case 2: a handler that always fails exhausts all 4 retries and escalates to the DLQ, with the DLQ entry recording the event id, attempt count, and error. Case 3: redelivering the same event id from case 1 (already committed) is caught by the idempotency check and returns immediately with zero attempts, confirming the check runs ahead of the retry loop rather than inside it.
Trade-offs and pitfalls
Complexity. O(1) work per attempt (one comparison, one exponentiation, one random draw); total attempts bounded by maxRetries + 1, so worst-case work per message is O(maxRetries).
Edge cases.
-
A message redelivered mid-retry (a second copy of the same event arrives while the first copy's retry loop is still running) needs the idempotency check to be safe under concurrency too, a check-then-act pattern here has the same race risk as any other idempotency check, it needs an atomic claim, not just a lookup.
-
If
handleritself is not safe to partially re-run (it has non-idempotent side effects), retrying at all is unsafe regardless of how well the backoff is tuned, the retry middleware assumes the handler's work is safe to repeat. -
A cap (
maxDelayMs) that is too low relative to how long the downstream outage actually lasts causes retries to keep hammering a still-down dependency at the capped rate instead of backing off further. -
Sleeping the exponential value itself instead of a random value up to it (no jitter at all) reintroduces the thundering-herd risk this design exists to avoid.
-
Forgetting to insert the idempotency check ahead of the retry loop (instead of, say, only checking once right before the DLQ push) lets an already-succeeded message consume retry attempts and, worse, potentially re-run a non-idempotent handler on every redelivery.
You are on call for an asynchronous data-ingestion pipeline with an SLA to persist every event within 60 seconds. Design the observability and alerting strategy: what would you track, what would actually be worth paging someone for, and what would your on-call runbook tell them to check first?
Sample Answer
Direct answer
For a 60-second persist service-level agreement (SLA), the metric that actually matters is the projected time-to-persist for the oldest unprocessed event. Everything else, producer rate, consumer processing rate, dead-letter-queue rate, error rate, raw throughput and latency, exists to explain why that projection is moving, and paging should be tied to SLA risk crossing a threshold, not to any single metric wobbling on its own.
Structured elaboration
What to track.
- Producer rate: events published per second, the input side of the balance.
- Consumer processing rate (throughput): events actually persisted per second, the output side.
- Consumer lag: how far behind the consumer is, in event count or in the age of the oldest unprocessed event, the single number that most directly answers whether the pipeline is on track to breach the 60-second SLA.
- Dead-letter queue (DLQ) rate and depth: events that failed processing and were routed out of the normal flow. Growth here means a subset of events are silently not going to be persisted on time, or at all, until someone intervenes.
- Error rate: the fraction of processing attempts failing, whether or not they eventually land in the DLQ, often a leading indicator that rises before lag does.
- End-to-end latency: measured from event creation to successful persist, the metric the SLA is actually stated in terms of, and it should be tracked directly wherever the event carries a creation timestamp, not just inferred from lag and rate.
What's worth paging for. Not every metric wobble. Page on SLA risk specifically: when the projected additional latency for the oldest currently unprocessed event crosses a threshold meaningfully inside the 60-second budget, enough time left for someone to actually act, not an alert that fires the same moment the SLA is already breached, or when DLQ depth grows continuously rather than staying flat. A flat, near-zero DLQ is healthy; a climbing one means a systemic issue is actively producing failures faster than anyone is triaging them. A momentary lag blip that self-resolves within the next reporting interval, or a brief error-rate spike that recovers, is exactly the kind of thing that belongs on a dashboard, not something that should wake someone up.
Runbook: what to check first. With lag rising or an SLA-risk page firing, the ordered checks are: first, is the producer rate elevated, an unusual traffic spike that may just need consumer autoscaling to catch up, or is the consumer processing rate depressed, a downstream dependency slowdown, a bad deploy, resource exhaustion on the workers. This distinguishes "we need more capacity" from "something is broken." Second, check the DLQ: is it growing, and if so, sample a few entries to see whether they share a common cause, one bad event shape from one producer versus a systemic downstream failure. Third, check recent deploys and recent consumer restarts or scaling events; a consumer-group churn event, workers restarting or the assignment of work across them changing, is a common, usually self-resolving cause of a short lag spike, worth ruling out before assuming a deeper problem. Fourth, check the downstream dependency the persist step actually writes to, since a slow or unavailable datastore on the write side shows up first as processing-rate degradation, not as anything on the ingestion side.
Worked example
Suppose consumer lag is currently 12,000 unprocessed events and the consumer is processing at a steady 500 events/sec:
additional latency=50012000=24 sThe oldest unprocessed event would take about 24 seconds to be reached and persisted at the current rate, comfortably inside the 60-second SLA. If lag instead grows to 27,000 events at the same 500 events/sec processing rate, the projected additional latency is 27000/500 = 54 seconds, past a page threshold set at 75% of the SLA budget:
page threshold=0.75×60=45 sand closing in on the 60-second breach itself. That is the SLA-risk crossing that should page, not the raw lag number in isolation, the same 12,000-event lag would be a non-event if the processing rate were higher, or would already be an active breach if the rate had dropped low enough.
Trade-offs and pitfalls
- Paging on raw lag count without dividing by current processing rate produces false urgency during high-throughput periods and false calm during low-throughput periods. The SLA is stated in time; the alert should be too.
- Alerting on every DLQ message individually instead of on DLQ growth rate and depth creates alert fatigue fast. A slow, steady trickle into the DLQ from known bad producer data is a triage-queue item, not a page; a sudden acceleration is.
- Treating a transient consumer-group reassignment as an incident before checking whether lag self-resolved wastes on-call attention on a self-healing event; conversely, dismissing every lag spike as probably just a reassignment, without checking, risks missing a real regression. The runbook should say check it, not assume it.
- Tracking throughput and latency as separate, unrelated numbers instead of connecting them to the SLA budget makes the dashboard descriptive but not actionable. The projection, lag divided by rate, compared to the SLA, is what turns raw metrics into a decision.
Explain the difference between message queues and event streams using: (1) simple high-level definitions, (2) step-by-step architecture differences (point-to-point vs publish-subscribe, retention and consumer offsets), (3) concrete use cases (task queue for background jobs vs event log for analytics), (4) discuss delivery semantics (at-least-once, at-most-once, exactly-once) and operational trade-offs.
Sample Answer
Direct answer
A message queue is a transient work-distribution mechanism: a message is delivered to one consumer and then typically removed. An event stream is a durable, ordered, append-only log that many consumers can read and re-read independently, each tracking their own position. Use a queue for background jobs that need to happen once; use an event stream when multiple consumers need to replay or independently process the same history of events, such as for analytics.
Structured elaboration
(1) Simple high-level definitions. A message queue holds discrete units of work; once a consumer successfully processes (and acknowledges) a message, it is generally gone from the queue. An event stream holds an ordered, append-only sequence of events that persists for a configured retention period regardless of whether any consumer has read them yet, and multiple independent consumers can read the same stream at their own pace.
(2) Step-by-step architecture differences.
- Delivery model: a queue is point-to-point: competing consumers pull from the same backlog, and each message goes to exactly one of them. An event stream is publish-subscribe at the storage layer: every consumer group gets its own full view of the log, so 5 independent consumer groups can each read every event, not compete for it.
- Retention: a queue typically discards a message once acknowledged (or after a short retention window for safety); an event stream retains events for a configured duration (hours to indefinitely) independent of consumption, which is what makes replay possible.
- Consumer offsets: a queue broker tracks per-message state (delivered, in-flight, acknowledged) on the broker's behalf; a consumer of an event stream tracks its own offset (a pointer into the log marking "I have processed up through here"), which the consumer can reset backward to reprocess history or forward to skip it, something a queue's consume-once model does not support.
(3) Concrete use cases. A task queue for background jobs (resize an uploaded image, send a single email, generate a PDF) is the queue use case: each job runs once, and once done there is no reason to revisit it. An event log for analytics (every purchase event, retained for 30+ days, read independently by a real-time dashboard consumer, a daily batch-aggregation consumer, and a fraud-detection consumer, each at their own offset) is the event-stream use case: the same history needs to serve multiple, independently-paced readers, and a new consumer added next month should be able to read events from before it existed.
(4) Delivery semantics and operational trade-offs. Both queues and streams commonly offer at-least-once delivery (a message may be delivered more than once, so consumers must be idempotent) as the practical default, since it is the easiest guarantee to build reliably. At-most-once (a message may be silently lost but never duplicated) is rarely chosen deliberately; it falls out of not retrying failed deliveries and is generally the wrong default for anything that matters. Exactly-once, from the consumer's experience, is achievable at the application layer by combining at-least-once delivery with a consumer-side idempotency check (a deduplication store keyed by message ID with a time-to-live, TTL), rather than relying on a broker guarantee alone; broker-level exactly-once semantics (EOS) protocols exist on some platforms but their internals are a streaming-platform implementation detail, not something the consuming application needs to reason about. Operationally, queues are simpler to run (bounded backlog, straightforward monitoring of queue depth) but lose replay ability; event streams give you replay and multiple independent consumer groups at the cost of needing retention-size planning, offset-lag monitoring per consumer group, and generally a heavier broker to operate.
Worked example
An e-commerce platform's "PurchaseCompleted" event needs to reach a real-time inventory-decrement consumer (must process within seconds), a nightly analytics-aggregation job (reads the whole day's events once, in batch), and a newly-added fraud-review consumer (needs to backfill and reprocess the last 7 days of purchases to retrain a rule set). This is only possible on an event stream: the inventory consumer reads from the current offset forward, the analytics job resets its offset to the start of each day, and the fraud-review consumer resets its offset back 7 days on rollout, all three reading the same retained log independently. If "PurchaseCompleted" had instead been placed on a point-to-point queue, only one of the three would ever have received a given event, and none of them could replay history that had already been consumed.
Trade-offs and pitfalls
The most common pitfall is reaching for an event stream by default because it seems more capable, then paying its operational cost (retention planning, per-consumer-group offset-lag monitoring, a heavier broker) for a workload that was really a simple once-only task queue. The opposite pitfall is placing something that needed independent replay, like the fraud-review backfill above, onto a queue and discovering months later that the history needed to replay it was already gone. A senior answer treats "does more than one independent party need to read this, possibly at different times or replaying history" as the deciding question, not "which technology is newer or more scalable."
Unlock Full Question Bank
Get access to all 47 Event-Driven Architecture and Asynchronous Messaging interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.