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 producer team wants to remove a field that several downstream consumer teams currently read from an event. Two consumer teams say removing it will break their service; a third says they don't use it. Walk through how you would let the producer team make this change (and future ones like it) without breaking consumers who depend on the old shape, and what you would need in place beforehand for that to be possible.
Sample Answer
Direct answer
Do not trust the third team's "we don't use it" self-report as the basis for removing the field: verify usage empirically (via consumer-driven contract tests or a lineage/usage audit against the registry) before touching anything, then make the change additive and reversible rather than destructive: deprecate the field with a stated window while the two dependent teams migrate to whatever replaces it, and only physically remove it after the registry and the audit both agree nobody depends on it. What has to be in place beforehand for this (and every future change like it) to be low-risk is a schema registry with enforced compatibility, consumer-driven contracts wired into the producer's continuous integration (CI) pipeline, and field-level usage visibility, not policy written down but unchecked by tooling.
Structured elaboration
Step 1: verify the claims, don't just count votes
- Two teams say removal breaks them; one says they don't use it. Before acting on any of these statements, check field-level usage against real traffic (a lineage/audit view keyed to the registry, or the third team's own consumer-driven contract, if one exists) rather than relying on memory or a quick grep of code that might miss dynamic field access. A team that "doesn't use" a field today may still have a downstream job or dashboard reading it that the team itself is not aware of.
- If no usage-tracking tooling exists yet, this is itself the finding: the producer team cannot safely make this change (or any future one) without it, and building that visibility becomes the first deliverable, not an afterthought.
Step 2: never remove directly, deprecate first
- Mark the field deprecated in the schema registry with a stated time-to-live (a concrete window, for example 60 to 90 days depending on how disruptive the change is), rather than removing it in the next release. Deprecation is reversible; removal is not.
- If the field is being replaced rather than dropped outright (a common case: a flatter
carrier_codestring replaced by a richercarrierobject), ship the replacement as an additive, backward-compatible field alongside the old one first, so the two dependent teams can migrate to the new field on their own schedule while the old one still works.
Step 3: give the two dependent teams a real migration path, not a deadline
- Communicate the deprecation and its window through the same channel every consumer team already watches (the registry's deprecation surface, not a one-off message that is easy to miss).
- Track each dependent team's migration status against the deprecation window on a shared dashboard, so "day 89 and two teams still on the old field" is visible before it becomes an incident, not discovered at the deadline.
- If a team cannot realistically migrate within the standard window, that is a scheduling negotiation with a visible tracked extension, not silent, indefinite delay and not a forced break either.
Step 4: remove only once verified safe
- Remove the field only after both signals agree: the deprecation window has closed, and the usage audit independently confirms zero real traffic reads it, including the third team's prior "we don't use it" claim, now verified rather than assumed.
What must be in place beforehand for this, and future changes like it, to be possible
- A schema registry that enforces compatibility on every change automatically, so an accidental breaking change cannot ship silently regardless of what any individual engineer intends.
- Consumer-driven contracts (or equivalent field-level usage tracking / lineage) wired into the producer's continuous integration (CI) pipeline, so "who actually depends on this field" is a query, not a Slack thread across three teams.
- A documented, tooled deprecation workflow (a way to mark a field deprecated with a TTL, and tooling that surfaces that to consumers and eventually blocks removal until the window and the usage audit both clear), so this becomes a repeatable, low-drama process rather than a one-off negotiation every time a producer wants to evolve its data.
- Without these three things in place, the producer team's only honest options for a genuinely breaking change are: negotiate with every consumer team by hand each time (slow, does not scale past a handful of consumers), or ship the breaking change and hope (which is what created this exact situation).
Worked example
Say the field in question is legacy_shipping_zone on a shipment-events topic with 8 consumer teams subscribed. Team A and team B say they read it (matches the scenario's "will break" teams); team C says they don't.
- A usage audit against 30 days of production traffic (using consumer group lag/read metrics keyed to which fields a consumer's deserializer actually accesses, or the registry's usage tracking if wired up) confirms: team A genuinely reads it in a nightly job, team B reads it in a real-time dashboard, and team C's claim holds, zero reads from team C's consumer group in the 30-day window.
- The field is marked deprecated with a 60-day TTL. A replacement
shipping_zone_v2field ships additively alongside it in the same event, so team A and team B can migrate independently: team A's nightly job is updated in week 2 (low urgency, batch job, easy to redeploy); team B's real-time dashboard, wired into a customer-facing screen, is updated in week 7 after more careful testing. - At day 60, the audit is re-run: both team A and team B now read
shipping_zone_v2exclusively; nobody readslegacy_shipping_zone. The field is removed. Team C, whose original claim was correct, was never blocked or delayed by any of this.
Trade-offs and pitfalls
- Common wrong turn: trusting a consumer team's self-reported "we don't use it" without verification. Self-reports are honest but incomplete; a dynamic field access, a downstream job the reporting team forgot about, or stale documentation can all make a well-intentioned "we don't use it" wrong.
- Common wrong turn: removing the field once the two dependent teams say they've migrated, without an independent audit confirming zero remaining traffic. "We think we're done migrating" and "the traffic confirms nobody reads the old field" are different claims, and only the second one is safe to act on.
- Common wrong turn: treating this as a one-time negotiation to get through, rather than as evidence that the team needs standing tooling (registry, contracts, usage visibility) so the next field removal does not require the same manual, three-team back-and-forth.
- Senior signal: naming the prerequisite tooling explicitly as the actual answer to "what would you need in place beforehand," rather than only describing the sequence of steps for this one field.
Event-sourcing stores all state changes as events. Discuss the trade-offs between storing only events versus introducing periodic snapshots. As a data engineer, explain snapshotting frequency, snapshot storage, snapshot validation, rehydration cost, and strategies for compaction or archival to control event-store growth.
Sample Answer
Direct answer
Storing only events gives perfect auditability and deterministic rebuilds, but rehydration cost (the work to replay events into current state) grows with the event count, so as a stream ages, reads and rebuilds get slower and more expensive. Periodic snapshots cap that cost by giving rehydration a recent starting point instead of the beginning of time, at the price of extra storage and a validation problem: a snapshot must be provably consistent with the events it claims to summarize. As a data engineer, pick a snapshot cadence and compaction policy driven by measured rehydration cost, not by a fixed rule of thumb.
Structured elaboration
Snapshotting frequency
- Event-count-based: snapshot every N events (commonly in the low thousands, tuned to the aggregate). Predictable worst-case replay cost regardless of how much wall-clock time has passed.
- Time-based: snapshot daily or hourly, better for aggregates with low or bursty event velocity where event count alone is not a reliable trigger.
- Hybrid: aggressive event-count triggers for hot, frequently-updated aggregates; time-based triggers for long-lived but rarely-updated ones.
- Whichever trigger is chosen, tune it against measured rehydration cost (replay time and CPU per rehydration), not a value picked without data.
Snapshot storage
- Store snapshots as immutable objects in cost-efficient storage, tagged with the aggregate identifier, the schema version, and the sequence number of the last event folded into the snapshot.
- Keep hot, frequently-accessed snapshots in fast storage or a cache; move cold snapshots to cheaper storage tiers with lifecycle rules.
Snapshot validation
- A snapshot must carry: the last-applied event sequence number, a schema version, and a checksum of its contents.
- On load, verify the checksum and confirm sequence continuity (no gap between the snapshot's last-applied sequence and the next event to replay). On mismatch, fall back to a full replay from the last known-good snapshot or from the beginning.
- Run periodic background verification jobs that independently rebuild a sample of aggregates from events alone and diff against the stored snapshot, to catch silent snapshot corruption before it is relied upon.
Rehydration cost
- Rehydration cost is: load the snapshot, then replay only the events recorded after that snapshot's sequence number. Cost scales with events-since-snapshot, not with the aggregate's total lifetime event count.
- The snapshot cadence directly bounds the worst case: with a snapshot taken every N events, no rehydration ever replays more than N-1 events.
Compaction and archival to control event-store growth (folding the replay/backfill and derived-dataset-reprocessing nuance)
- Compact by taking a full snapshot and, only where retention rules permit it, archiving or truncating events older than that snapshot's sequence to cold, cheaper storage rather than deleting them outright; legal or audit requirements often forbid true deletion.
- The archival tier still has to remain replayable: any derived dataset (an analytics table, a machine learning feature store, a rebuilt read model) that was originally built by consuming the full event history needs that same history available if it must ever be reprocessed, for example after a bug fix in the transformation logic. A retention policy that only optimizes for "rehydrate a live aggregate quickly" and quietly discards old events breaks backfill for any derived dataset that depended on that full history, even though the live aggregates themselves are unaffected. Decide retention and compaction against both use cases explicitly, not just the aggregate-rehydration one.
- Use a tiered retention: hot recent events in the primary event store, older events moved to compressed, partitioned cold storage with an index that supports selective replay by aggregate and time range.
- Mark truncation points explicitly (a compaction marker event, or metadata) so any consumer replaying the stream can detect where the live tier's history starts and knows to fetch older ranges from the archive if it needs them.
Worked example
An account aggregate receives 50 events per day on average.
- Without snapshots, rehydrating that account after 5 years of history means replaying 5 * 365 * 50 = 91,250 events.
- With a snapshot taken every 1,000 events, the worst-case replay after any snapshot is 999 events (999 / 91,250 = 1.09%, so rehydration cost is bounded to roughly 1% of the no-snapshot case at year 5, not strictly under 1%: it is a hair over).
- At 50 events/day, a snapshot every 1,000 events fires roughly every 20 days (1,000 / 50 = 20), so the account accumulates at most 20 days' worth of unsnapshotted events at any time.
- If the business later needs to reprocess the full 5-year history to backfill a new derived dataset (say, a new fraud-scoring feature that needs every historical event, not just the latest snapshot), the archived event range covering all 91,250 events must still be retrievable even though live rehydration never touches most of them.
Trade-offs and pitfalls
- Common wrong turn: choosing a snapshot cadence without measuring actual rehydration cost, then discovering it under- or over-snapshots (too frequent wastes storage and write bandwidth on snapshotting; too infrequent leaves rehydration slow).
- Common wrong turn: treating snapshot validation as optional. An unvalidated, silently corrupt snapshot is worse than no snapshot: it produces confidently wrong current state instead of forcing a (correct, if slow) full replay.
- Common wrong turn: designing retention purely around live-aggregate rehydration speed and discovering, only when a derived dataset needs to be rebuilt, that the events required for that rebuild were already archived out of reach or deleted.
- Senior signal: naming the specific numeric relationship between snapshot cadence and worst-case replay cost, and treating archival policy as a decision that serves more than one consumer (live rehydration and derived-dataset reprocessing), not just the first one that comes to mind.
Explain event-driven architecture and contrast it with synchronous request-response architectures. As a data engineer, identify the core components (producers, brokers, topics/queues, consumers), typical data-pipeline use cases (CDC, audit trails, streaming enrichment), and the trade-offs (coupling, latency, fault isolation, operational complexity) when you choose event-driven designs for data workloads.
Sample Answer
Direct answer
Event-driven architecture (EDA) structures a system around producers emitting events (facts that something happened) to a broker, which independently delivers them to one or more consumers, rather than a caller directly invoking another service and waiting for a response. The core trade-off versus synchronous request-response is that EDA buys loose coupling, independent scaling, and fault isolation at the cost of latency (results aren't immediate) and operational complexity (more moving pieces, harder end-to-end debugging).
Structured elaboration
Contrast with synchronous request-response. In a synchronous model, service A calls service B directly (commonly over HTTP or a remote procedure call) and blocks waiting for B's response; A knows about B specifically, and if B is slow or down, A is directly affected. In an event-driven model, A publishes an event describing what happened and moves on immediately; A does not know or care who, if anyone, consumes that event, and a consumer processes it whenever it's able to, independent of A's timing.
Core components.
- Producers: services or components that emit events when something of interest happens (an order was placed, a user updated their profile).
- Brokers: the intermediary system that receives events from producers and delivers them to consumers (examples include a managed pub/sub service or a distributed log/queueing platform); the broker decouples producers from consumers so neither needs a direct network reference to the other.
- Topics/queues: the named channels within the broker that events are published to and consumed from; a topic is typically fan-out (many consumers can each receive a copy of every event), while a queue is typically point-to-point (one message is delivered to exactly one consumer among a competing group).
- Consumers: services that subscribe to a topic or read from a queue and process the events they receive, often producing further events of their own as a result.
Typical data-pipeline use cases. Beyond application messaging, event-driven patterns are a common backbone for data pipelines specifically:
- Change-data-capture (CDC): capturing every row-level insert/update/delete in a source database as a stream of events, so downstream systems (a search index, a cache, an analytics store) stay in sync without querying the source database directly or on a slow batch schedule.
- Audit trails: because events are an immutable record of "what happened, in order," they naturally serve as an audit log, useful for compliance and for reconstructing how a piece of data reached its current state.
- Streaming enrichment: consuming a raw event stream and augmenting it with additional context (looking up a customer's segment, geocoding a location) before republishing an enriched event for downstream consumers, so that enrichment logic lives in one place rather than being duplicated by every consumer that needs it.
Trade-offs when choosing event-driven for data workloads.
- Coupling: EDA reduces coupling significantly, a producer doesn't need to know which or how many consumers exist, so new consumers can be added without changing the producer at all; synchronous designs couple the caller directly to the callee's availability and interface.
- Latency: synchronous calls give an immediate result; event-driven processing is asynchronous by nature, so there's an inherent delay (typically milliseconds to low seconds under healthy conditions, but with no hard upper bound unless explicitly engineered and monitored) between an event being produced and a consumer acting on it.
- Fault isolation: if a consumer is down or slow in an event-driven design, events queue up and are processed once it recovers, the producer and other consumers are unaffected; in a synchronous chain, a single slow or failing downstream service can cascade failure back to the caller.
- Operational complexity: event-driven systems introduce more infrastructure to run and reason about (the broker itself, delivery guarantees, ordering, retry and dead-letter handling, distributed tracing to debug a chain of asynchronous hops), which is real added complexity a purely synchronous system doesn't have to manage.
Worked example
A CDC pipeline: a relational database's write-ahead log is tailed by a CDC connector, which emits one event per row change, for example {"table": "orders", "op": "UPDATE", "id": 4821, "after": {"status": "shipped"}}, to a topic named db.orders.changes. A streaming-enrichment consumer reads that topic, looks up the customer's loyalty tier for order 4821, and republishes an enriched event to orders.enriched containing both the original change and the loyalty tier. A separate audit consumer independently reads the same db.orders.changes topic and appends every event, unmodified, to a durable audit log. Contrast this with the synchronous alternative: the order-status-update code path would need to directly call a search-index-update function, a loyalty-lookup function, and an audit-log-write function in line, meaning a bug or slowdown in any one of those three calls could block or fail the original status update itself; in the event-driven version, those three concerns are fully independent consumers of the same event, and a failure in one does not affect the others or the original write.
Trade-offs and pitfalls
The most common mistake in evaluating this trade-off is treating "asynchronous" as strictly better across the board; a workflow where the caller genuinely needs an immediate answer (checking whether an item is in stock before showing "add to cart") is usually a poor fit for a fully asynchronous redesign, since the user experience needs a synchronous response even if some other part of the system reacts to the resulting order asynchronously. A second common pitfall is underestimating the debugging cost: tracing a single business action through multiple independent consumers requires deliberate observability investment (correlation ids, distributed tracing) that a synchronous call stack gives you for free via a single request's own logs and stack trace. Event-driven data pipelines specifically also need to account for events arriving out of order or being redelivered, which a naive consumer written like a simple database trigger will not handle correctly by default.
flowchart LR
P1[Producer] -->|publishes event| B[(Broker: topic or queue)]
B -->|delivers event| C1[Consumer 1]
B -->|delivers event| C2[Consumer 2]
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.
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.
Unlock Full Question Bank
Get access to all 28 Event-Driven Architecture and Asynchronous Messaging interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.