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.
Explain Command Query Responsibility Segregation (CQRS). As a data engineer, when is CQRS valuable for analytics or operational workloads? Discuss trade-offs including complexity, eventual consistency of read-models, and strategies to make reads 'fresh' when required.
Sample Answer
Direct answer
Command Query Responsibility Segregation (CQRS) is the pattern of using a different model for writes (commands that change state) than for reads (queries that return state), instead of forcing one schema to serve both well. As a data engineer, it earns its added machinery when the write side's transactional shape and the read side's analytical or lookup shape diverge enough that one schema serves neither well, for example narrow row-level operational writes versus wide, pre-aggregated reporting reads. Treat it as a deliberate trade: schema simplicity for the ability to scale, model, and store reads and writes independently.
Structured elaboration
What CQRS actually separates
- Command side: an authoritative model (a normalized transactional store, or an event log) that enforces write-time invariants and produces state changes.
- Query side: one or more purpose-built read models (denormalized tables, search indexes, in-memory caches), each shaped for a specific access pattern rather than for correctness enforcement.
- A projector connects the two asynchronously: it consumes the write side's changes (domain events, or change-data-capture (CDC) records) and updates the read model(s).
When CQRS is valuable for analytics workloads
- Reporting or business intelligence (BI) queries need aggregation shapes (rollups by category, hour, region) that would otherwise require expensive joins or full scans against the transactional schema.
- Several independent consumers need different projections of the same data (a finance rollup, a fraud-detection view, a customer dashboard); three denormalized read models are cheaper to operate than three sets of ad-hoc joins against the online transaction processing (OLTP) store.
- Analytical queries would otherwise contend for locks and I/O with operational writes on the same tables.
When CQRS is valuable for operational workloads
- Write throughput and read throughput need to scale independently and at different rates (write-heavy ingestion feeding a low-cardinality operational dashboard).
- Write-side invariants are complex enough (state machines, multi-step validation) that mixing them with read-optimization concerns would make the write model harder to reason about.
- Not valuable: a small application with one read pattern that already matches the write schema. There CQRS adds a projector, extra storage, and extra failure modes with no offsetting benefit.
Trade-off: complexity
You now operate an additional pipeline (the projector), additional storage (one or more read stores), and additional failure modes: projector lag, projector crashes mid-batch, and schema drift between the write shape and the read shape.
Trade-off: eventual consistency of read models
Because the read model updates asynchronously, a query issued immediately after a write can observe stale data. The size of that staleness window is a direct function of projector throughput and batching, not something a team can design around by ignoring it.
Strategies to make reads "fresh" when required
- Read-your-writes for the writer: return enough state in the command response (or a version/sequence number) that the client who just wrote never needs to trust the read model for its own write.
- Tighten the pipeline: smaller batches and event-driven push instead of periodic batch pull shrinks the staleness window, at the cost of more frequent projector invocations.
- Expose staleness explicitly: attach a last-updated version or timestamp to read-model responses so callers can judge whether the data is fresh enough, instead of the system silently presenting stale data as current.
- Selective synchronous update: for a small, well-identified set of critical fields, update the read model synchronously in the write path (accepting some coupling) while everything else stays asynchronous.
Worked example
An order system accepts writes at 500 orders per minute (about 8 to 9 orders per second) into a transactional order table. A "revenue by category, per hour" read model is built by a projector that drains the order-events stream every 60 seconds and applies that batch of updates.
- Worst-case staleness for that read model equals the batch interval: 60 seconds. An order committed just after a batch run will not appear until the next run.
- Average staleness is roughly half the batch interval, about 30 seconds, if orders arrive close to uniformly across the minute.
If the product requirement is "the dashboard must reflect a new order within 10 seconds," this projector cadence fails outright: 60 seconds worst case exceeds the 10-second bound. The fix is either to drop the batch interval below 10 seconds, or to read the specific "orders placed today" counter synchronously from the write side while the rest of the dashboard stays on the 60-second cadence.
Trade-offs and pitfalls
- Common wrong turn: adopting CQRS because it sounds architecturally sophisticated for a workload that has a single read pattern already matching the write schema. That is pure overhead with no payoff.
- Common wrong turn: treating "eventually consistent" as a detail to sort out later. Staleness needs an explicit, stated bound (or an explicit "no bound" with a user-facing affordance for it) decided at design time, not discovered in production when a user cannot see the order they just placed.
- Senior signal: naming a concrete staleness budget and matching the pipeline's cadence to it, rather than discussing CQRS only in the abstract.
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.
How would you choose a Kafka partitioning key for a user-events topic such that ordering is preserved per user but partitions are balanced across the cluster? Discuss hashing strategies, handling hot users, and approaches for multi-tenant fairness when some tenants generate far more events than others.
Sample Answer
Direct answer
Use user_id itself as the partition key so Kafka's partitioner (a mod-based hash of the key) always routes every event for a given user to the same partition, which is what gives per-user ordering, it falls directly out of using the same key for the same entity, not something built separately. Balance across the cluster then comes from having enough distinct keys and a reasonably uniform hash, with a deliberate, separate remediation for the minority of disproportionately hot keys.
Structured elaboration
Hashing strategy. Kafka's default partitioner computes hash(key) % num_partitions (using a murmur2-based hash internally); the exact hash function matters less than the property that it is deterministic, the same key always maps to the same partition number for a fixed partition count. Because the partition assignment is a pure function of the key, ordering per key is automatic once you commit to using that key consistently, not an extra mechanism you build on top.
Why "balanced" is a statistical property, not a guarantee. With many distinct user ids and a well-distributed hash, the law of large numbers makes total load spread roughly evenly across partitions, if every key generates traffic at roughly the same rate. The question specifically names the case where that assumption breaks: hot users.
Handling hot users. Two remediations, with different costs:
- Key salting (sub-partitioning): append a bounded suffix to a hot user's key only once their rate crosses a threshold (
user_id#0throughuser_id#K-1), spreading that one user's traffic across up to K partitions. The explicit cost: per-user ordering across the salted shards is no longer guaranteed, only ordering within each shard is. This is a real trade-off to name out loud, not a free win, if the consuming system genuinely needs strict per-user order even for hot users, salting defeats that and a downstream re-sequencing step would be needed, which mostly defeats the purpose of salting in the first place. - Isolating known hot keys: if hot users are identifiable ahead of time (service accounts, bots, high-volume enterprise integrations), route them to a separate topic or partition set with its own capacity plan, keeping the main topic's "roughly uniform" assumption valid for the long tail of ordinary users.
Multi-tenant fairness. Partitioning purely by user_id does nothing to stop one tenant (an account with many users) from occupying a disproportionate share of partitions and consumer capacity relative to a smaller tenant, since partition assignment has no concept of tenant at all. If fairness across tenants is a hard requirement: give large or noisy tenants a dedicated topic or partition range so they cannot starve others' consumer lag budget, and layer per-tenant rate limiting or quota enforcement on the consumer or producer side (the same token-bucket idea used for rate limiting generally), since a static partition assignment cannot adapt to traffic that changes over time the way an explicit quota can.
Worked example
A small, executed illustration of the mechanism (using Python's built-in zlib.crc32 purely as a stand-in deterministic hash function to demonstrate the mod-based mechanic; Kafka's actual default partitioner uses a different, murmur2-based hash internally, the mechanism, not the specific hash, is what this illustrates).
import zlib
num_partitions = 6
for uid in ["user-1042", "user-8831", "user-2207", "user-77519"]:
h = zlib.crc32(uid.encode())
print(uid, "partition=", h % num_partitions)
user-1042 partition= 3
user-8831 partition= 4
user-2207 partition= 4
user-77519 partition= 0
Every event for user-1042 always lands on partition 3, giving that user strict per-user order for free. Now simulate salting the hot user user-8831 into 12 sub-keys (user-8831#0 through user-8831#11):
counts = {}
for shard in range(12):
key = f"user-8831#{shard}"
p = zlib.crc32(key.encode()) % num_partitions
counts[p] = counts.get(p, 0) + 1
print(counts)
{4: 5, 0: 1, 5: 5, 1: 1}
The 12 salted sub-keys land on 4 of the 6 partitions instead of all landing on partition 4 as the unsalted key would. The spread is not perfectly even at this small sample size (partitions 4 and 5 got 5 each, partitions 0 and 1 got 1 each), which is expected statistical noise at n = 12, not a flaw in the technique; at real production volumes for a genuinely hot key, the same law of large numbers that balances ordinary keys also evens out a hot key's salted shards.
Trade-offs and pitfalls
- Salting trades away strict per-user order for the salted user specifically, in exchange for spreading their load; state this cost explicitly when proposing it, do not present it as a free fix.
- Changing the partition count later reshuffles the mod-based assignment for every existing key, not just new ones, since
hash(key) % num_partitionschanges for essentially every key whennum_partitionschanges. Under-provisioning partition count up front makes a later increase a disruptive, coordinated migration, not an incremental capacity add. - "Balanced across the cluster" and "fair across tenants" are different guarantees: the first is about hash uniformity over many keys, the second is about the business meaning of who owns which keys, no hashing scheme fixes tenant fairness on its own.
- Picking the wrong entity granularity for the key, for example partitioning by
session_idwhen the actual ordering requirement is per-user across sessions, silently breaks the guarantee the system was supposed to provide.
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.
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.
Unlock Full Question Bank
Get access to all 25 Event-Driven Architecture and Asynchronous Messaging interview questions and detailed answers.
Sign in to ContinueJoin thousands of developers preparing for their dream job.