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.
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.
As a Solutions Architect, detail the decision criteria you use to choose synchronous (HTTP/REST, gRPC) versus asynchronous (message queues, event streams) service-to-service communication. Discuss trade-offs around latency, reliability, coupling, operational complexity, developer productivity, and how each choice affects deployment independence.
Sample Answer
Direct answer
Choose synchronous communication (HTTP/REST, or gRPC, a high-performance remote-call framework built on protocol buffers) when the caller needs an immediate answer to proceed and strong end-to-end reliability guarantees matter more than decoupling; choose asynchronous communication (message queues, event streams) when the work can be deferred, the caller does not need the result to continue, or you need producers and consumers to fail, scale, and deploy independently. The decision criteria are latency requirements, reliability/failure-isolation needs, coupling tolerance, operational complexity budget, developer productivity, and how much deployment independence the teams involved actually need.
Structured elaboration
Walk each axis explicitly rather than picking a style by habit:
- Latency. Synchronous calls give the caller a result (or an error) within the request's timeout window, which is required whenever a human or an immediately-dependent step is waiting (e.g., "is this card valid"). Asynchronous messaging trades immediate response for throughput and resilience: the producer gets an acknowledgment that the message was accepted, not that the work finished.
- Reliability and failure isolation. A synchronous call couples the caller's availability to the callee's: if the downstream service is slow or down, the caller blocks or fails too, and a chain of synchronous calls compounds failure probability (if three downstream services each have 99.9% availability, the chain that calls all three synchronously is bounded by roughly 0.9993≈0.997, i.e. about 3x the failure rate of any single hop). Asynchronous messaging inserts a durable buffer (the queue or log) between producer and consumer, so a consumer outage delays processing but does not fail the producer's request.
- Coupling. Synchronous calls create temporal coupling (both sides must be up at the same instant) and often contract coupling (the caller depends on the callee's exact response shape and latency profile). Asynchronous messaging only requires agreement on the event/message schema; producer and consumer do not need to be online simultaneously.
- Operational complexity. Synchronous systems are simpler to trace (a single call stack, straightforward distributed tracing) and debug. Asynchronous systems add a broker to run and monitor, delivery-semantics decisions (at-least-once handling, ordering), dead-letter queues (DLQs, queues that hold messages a consumer could not process after retries) for poison messages, and eventual-consistency reasoning that developers have to learn.
- Developer productivity. Synchronous request/response is the default mental model most engineers already have; it is faster to build and test for simple CRUD-style interactions. Asynchronous flows require additional skills (idempotent consumers, correlation IDs for tracing a request across services, compensating logic) and slower local iteration (you cannot just curl an endpoint and see the final state).
- Deployment independence. This is the axis most often underweighted. A synchronous caller that depends on a callee's API contract must coordinate deploys carefully around breaking changes (versioned endpoints, backward-compatible fields). An asynchronous consumer reading from a durable topic can be redeployed, scaled, or even paused independently of the producer, because the event log absorbs the gap; this is what lets teams ship on independent cadences, which is usually the real reason "event-driven" gets proposed for a set of otherwise unrelated services.
A simple framework: ask "does the caller need the result to proceed, right now, in this request?" If yes, go synchronous. Then ask "if the callee is degraded, should the caller degrade too, or should the caller succeed and the work catch up later?" If the caller should still succeed, go asynchronous even if latency were not a constraint, because the failure-isolation property is what you are actually buying.
Worked example
A checkout service calling three downstream services synchronously (inventory check, tax calculation, fraud score) with each at 99.9% availability has a combined dependency availability of 0.9993≈0.997, meaning roughly 3 requests in 1,000 fail purely from the chaining, even though each service individually is healthy 999 times in 1,000. If two of those three calls (tax calculation and fraud score) can tolerate a few hundred milliseconds of extra latency and their results are not needed to authorize the immediate step, moving them to asynchronous event handlers that publish a "checkout.enriched" event removes two of the three synchronous dependencies from the critical path, leaving only inventory check (which genuinely must block, since you cannot confirm an order for stock you do not have) synchronous. The chain's availability floor improves to roughly 0.9991=0.999, and the tax/fraud services can now be redeployed or scaled without coordinating a maintenance window with checkout.
Trade-offs and pitfalls
A common wrong turn is defaulting to asynchronous "for scalability" on a step the caller genuinely needs the result of right now; that just relocates the wait (the caller polls or blocks on a callback) while adding a broker, a correlation mechanism, and eventual-consistency bugs, with no user-facing benefit. The opposite pitfall is defaulting to synchronous everywhere because it is simpler to write, then discovering that one flaky downstream service now takes the whole call chain down with it. A senior answer treats each interaction independently on these six axes rather than applying one style architecture-wide, and explicitly names deployment independence as a first-class criterion, not an afterthought, since it is usually the criterion that determines whether decoupling was worth the added operational complexity.
You must choose between Apache Kafka and RabbitMQ as the backbone for an internal eventing platform. Constraints: peak 50k events/sec, need durable replay for analytics, many consumers require at-least-once delivery, and some consumers require ordered delivery per key. List evaluation criteria, recommend one option with justification, and describe mitigations for its weaknesses.
Sample Answer
Direct answer
For an internal eventing platform at 50,000 events/sec with durable replay for analytics, at-least-once delivery for many consumers, and per-key ordering for some consumers, Kafka is the better fit: its partitioned, replayable log natively supports per-key ordering (all events for a given key land on the same partition) and long retention, which RabbitMQ does not provide as first-class capabilities. RabbitMQ remains the better choice when the workload is closer to task distribution with complex routing needs and no replay requirement, which this workload is not.
Structured elaboration
Evaluation criteria, applied to the four stated constraints:
- Sustained throughput at 50k events/sec: both platforms can reach this rate, but the design work to get there differs. Kafka scales throughput by increasing partition count (more parallelism), which is a capacity-planning exercise done up front. RabbitMQ scales by adding queues and consumers, and can also sustain very high throughput, but does not offer partition-based ordered parallelism as a native primitive.
- Durable replay for analytics: this is the decisive criterion. Kafka retains events for a configured window regardless of consumption and lets any consumer group replay from an arbitrary offset. RabbitMQ removes messages once acknowledged; replaying "already-processed" history is not a built-in capability and would require bolting on a separate durable store.
- At-least-once delivery for many consumers: both support at-least-once delivery natively. Kafka's consumer-group model lets many independent consumer groups each read the full stream at-least-once without competing with each other; RabbitMQ supports at-least-once via manual acknowledgment, but achieving "many independent full-stream readers" specifically (rather than competing consumers on one queue) requires fanning the exchange out to a separate queue per consumer group, which is more configuration to maintain as consumers are added.
- Ordered delivery per key for some consumers: Kafka provides this natively via partition-key hashing, exactly the "some consumers require ordered delivery per key" requirement. RabbitMQ can approximate per-key ordering with a single active consumer per queue and a routing key per entity, but that limits parallelism for that queue to one consumer, which is a real constraint at 50k events/sec.
Recommendation: Kafka, specifically because the combination of durable replay and per-key ordering are both native to its architecture, whereas RabbitMQ would need workarounds for both.
Mitigations for Kafka's weaknesses. Kafka's two well-known weak points relative to RabbitMQ are operational complexity (a broker cluster, plus metadata quorum management, to run and monitor) and the lack of complex, content-based routing that RabbitMQ's exchange types provide natively. Mitigate the operational-complexity weakness by using a managed offering (e.g., Amazon MSK or Confluent Cloud) rather than self-hosting the cluster. Mitigate the routing-flexibility weakness by keeping routing logic in the consuming application (a consumer subscribes to specific topics and filters or re-publishes as needed) rather than expecting the broker to route by message content, which is a reasonable trade given this workload's requirements are about ordering and replay, not complex routing.
Worked example
Sizing partitions for 50,000 events/sec with per-key ordering: assume, as a planning input for this example (not a general Kafka benchmark), that a single partition on the chosen broker instance size can sustainably sustain roughly 5,000 events/sec for this message payload size. The minimum partition count is:
Npartitions≥⌈5,000 events/s per partition50,000 events/s⌉=10
Ten partitions is the bare minimum to sustain the target throughput with zero headroom; in practice you would provision with headroom for growth and for uneven key distribution (a small number of very active keys, or "hot keys," can overload their assigned partition even when the average across all partitions is fine), so a reasonable starting point is roughly double the bare minimum, i.e. 20-24 partitions, revisited once real key-distribution data from production traffic is available. This calculation is illustrative: the 5,000 events/sec per-partition figure is a stated planning assumption for this example, to be replaced with a measured number from a load test against your actual message size and hardware before committing to a partition count in production.
Trade-offs and pitfalls
The main pitfall is under-provisioning partitions based only on average throughput and ignoring hot-key skew; a partition count that comfortably covers the 50k events/sec average can still fall over if one key (e.g., one very active customer or entity) accounts for a disproportionate share of traffic, since that key's events are all pinned to a single partition by design. A second pitfall is choosing partition count too high without a plan to ever reduce it: Kafka does not support shrinking partition count on an existing topic without breaking the per-key-to-partition mapping (and therefore ordering) for every key, so partition sizing decisions are effectively one-directional and should be made deliberately, not just rounded up generously. A third pitfall specific to this scenario is assuming "at-least-once for many consumers" and "ordered per key for some consumers" are the same requirement: ordering is a per-partition property that some consumers need to respect (by not parallelizing within a key), while at-least-once is a delivery guarantee every consumer needs to handle via idempotency regardless of whether it cares about ordering.
Explain backpressure and flow-control techniques across networked services and message brokers: TCP flow control, HTTP/2 flow control, reactive streams (request-n), credit-based broker flow control, and producer-side throttling. For each technique explain when it's most appropriate and how you'd design end-to-end controls to prevent OOMs and cascading failures.
Sample Answer
Direct answer
These five techniques operate at different layers, from the raw byte stream up to an individual message, and they compose: none of them alone protects an entire pipeline from an out-of-memory (OOM) crash or a cascading failure, the end-to-end design comes from layering them so each protects the specific thing the layer below it cannot.
Structured elaboration
| Technique | Mechanism | Layer | Most appropriate when |
|---|---|---|---|
| Transmission Control Protocol (TCP) flow control | Receiver advertises a byte window in each acknowledgment; sender never has more unacknowledged bytes outstanding than that window | Transport, automatic for any TCP connection | Baseline protection for any point-to-point byte stream, especially when you do not control the application protocol running on top of it |
| HTTP/2 flow control | A credit-like byte window, but per multiplexed stream as well as per connection | Application-transport boundary | Many logical exchanges (for example gRPC calls) multiplexed over one TCP connection, where one slow stream should not starve the others sharing that connection |
Reactive streams (request(n)) | The subscriber explicitly tells the publisher how many items it is ready to receive next; the publisher cannot push more than requested | Application, in-process or library-level | An async pipeline or stream-processing library where backpressure needs to be an explicit, item-oriented part of the consumer's own control flow, not inferred from bytes |
| Credit-based broker flow control | A consumer grants the broker a bounded number of unacknowledged in-flight messages (a prefetch or credit limit); each acknowledgment replenishes one unit | Message-broker to consumer relationship | Bounding one consumer's own in-flight work without needing the producer to coordinate with it directly, the broker mediates |
| Producer-side throttling | The producer itself is rate-limited before sending, based on a fixed budget or on live downstream health signals (queue depth, consumer lag, error rate) | Outermost, protects the whole system including the broker | Protecting the system as a whole, including the broker's own storage, not just one consumer |
Why layering matters, not picking one. TCP and HTTP/2 flow control protect the wire, an application can still accumulate an unbounded in-memory backlog even while the underlying connection is perfectly healthy, because those layers only bound bytes in flight on the network, not messages queued in application memory waiting to be processed. Credit-based broker flow control is what actually bounds an individual consumer's own memory footprint, by capping how many unacknowledged messages the broker will hand it. Producer-side throttling is the only one of the five that protects the broker's own memory or disk, without it, a healthy consumer relationship does nothing to stop the broker itself from accumulating an unbounded queue if producers keep publishing faster than the system can drain.
Preventing cascading failure specifically. Without an upstream signal, a slow consumer causes the broker's queue depth to grow unbounded, risking the broker's own out-of-memory (OOM) crash or disk exhaustion. If the broker then starts shedding load (rejecting connections or writes) to protect itself, producers that do not back off in response to that rejection can retry-storm the broker, making the incident worse instead of better. This is why producer-side throttling needs to be dynamic, responsive to live broker or consumer health signals, rather than a fixed rate calibrated only for normal conditions; a fixed rate does nothing extra during the exact moment it is needed most.
Worked example
A consumer configured with a prefetch (credit) limit of 50 unacknowledged messages, average payload size 2 kilobytes (KB).
With credit-based flow control: the broker never has more than 50 unacknowledged messages outstanding to this consumer, so its in-flight memory footprint for this consumer's queue is bounded at roughly 50 x 2 KB = 100 KB, regardless of how large the total backlog waiting behind those 50 becomes.
Without it (an unbounded prefetch): the broker keeps pushing every available message, so the consumer's in-flight count, and its memory footprint, grows with the size of the backlog itself, with no structural bound, until the backlog stops growing or the consumer runs out of memory.
The structural difference (bounded versus unbounded growth) is what the mechanism buys you; the 100 KB figure is a direct, pinned-input calculation (50 messages times 2 KB each), not a measurement.
Trade-offs and pitfalls
- Setting broker credit or prefetch too low sacrifices throughput, many small round trips to replenish a small credit budget.
- Setting it too high defeats the purpose, it looks bounded on paper but the bound is too large to meaningfully protect memory.
- Assuming TCP or HTTP/2 flow control alone is "enough" backpressure is the single most common pitfall named implicitly by this question: those protect the wire, not an application's own queues or memory.
- Static, fixed-rate producer throttling calibrated for normal conditions does not actually prevent cascading failure during a real incident, since the fixed rate was never designed to respond to the incident happening; dynamic throttling that reacts to consumer lag or broker health is what closes that gap.
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.
Unlock Full Question Bank
Get access to all 26 Event-Driven Architecture and Asynchronous Messaging interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.