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 customer needs an immutable audit trail and the ability to rebuild multiple read models quickly. Compare event sourcing + CQRS against a traditional relational database augmented with Change Data Capture (CDC). Discuss complexity, operational cost, replayability, schema evolution, developer ergonomics, and scenarios where event sourcing is or is not justified.
Sample Answer
Direct answer
For an immutable audit trail with fast rebuild of multiple read models, both event sourcing plus Command Query Responsibility Segregation (CQRS) and a relational database augmented with change-data-capture (CDC) can work, but they solve it differently: event sourcing stores domain intent as the source of truth and replays it deterministically, while CDC turns an existing relational database's row-level changes into a changelog that downstream consumers materialize into views. Pick event sourcing when the business needs the "why," not just the "what," and needs a canonical replayable log; pick relational-plus-CDC when the relational model already fits the domain and the team wants a lighter operational footprint.
Structured elaboration
Complexity
- Event sourcing + CQRS: higher architectural complexity. Domain events are the source of truth; the team must implement an append-only event store, event versioning, snapshotting, and one or more projections, and must handle at-least-once delivery into those projections.
- Relational + CDC: lower incremental complexity if the team already runs a relational database. CDC (typically reading the database's write-ahead log, for example via Debezium) streams committed transactional changes to downstream systems without changing the core domain model.
Operational cost
- Event sourcing needs operational maturity for the event store itself: scaling, backups, compaction, and snapshotting are all new operational surfaces, plus projection workers and a message bus.
- CDC leans on existing database tooling and log-shipping infrastructure; fewer bespoke components, generally lower operational overhead, but the pipeline is now coupled to the source database's internal log format and retention.
Replayability (folding the append-only-log / materialized-view nuance)
- Event sourcing's replay is native and deterministic: any read model or a brand-new projection can be rebuilt by replaying the event log from the start (or from a snapshot). The event log is intentionally an append-only log of domain intent, and every read model is a materialized view derived from it.
- CDC is structurally similar (the database's write-ahead log is also an append-only log, and CDC consumers building denormalized tables are also building materialized views), but the two differ in what the log records: CDC's changelog captures row-level state deltas ("column X became Y"), not business intent ("customer upgraded their plan"). Reconstructing intent from a stream of row deltas is lossy and sometimes ambiguous, and a full CDC-based rebuild requires retaining the source database's change history for as long as you might need to replay it, which is a weaker retention guarantee than an event store's log is designed to give.
Immutable audit trail and intent
- Event sourcing stores intent explicitly as first-class domain events; the audit trail directly answers "what business action happened and why."
- CDC provides a factual log of persisted state changes, useful for audit of what changed, but not necessarily why, unless the application already wrote that intent into the row (e.g., an explicit
reasoncolumn).
Schema evolution
- Event sourcing requires an explicit event-versioning strategy (adding fields with defaults, or introducing a new event type for a breaking change, plus upcasters that translate old event versions when replayed).
- CDC schema evolution tracks the source table's schema; adding a nullable column is generally safe, but column renames, type changes, or table restructuring can break the CDC pipeline's mapping and every downstream consumer of it.
Developer ergonomics
- Event sourcing has a steeper learning curve: developers must think in events, eventual consistency, and projection design, but gain very clear auditability and strong support for temporal questions ("what did this account look like on March 3rd?").
- CDC lets developers keep writing familiar create/read/update/delete (CRUD) code against the relational schema; projections consume the CDC stream with comparatively little domain-model rework.
When event sourcing + CQRS is justified
- A complex domain with rich auditability or regulatory replay requirements, where the business needs to reconstruct exact past system state.
- Many read models that change frequently and need fast, correct rebuilds.
- The business explicitly wants to capture intent, not just state, for analytics or downstream machine learning use.
When it is not justified
- A straightforward create/read/update/delete domain where the relational schema already models the business well and the database's own transaction log already satisfies audit and retention requirements.
- Limited engineering bandwidth or a need for fast delivery where the extra event-sourcing machinery would slow the team down for no corresponding benefit.
Worked example
A payments team is choosing between the two approaches for a ledger requiring 7 years of audit retention and 3 read models (customer statement, fraud-review queue, regulatory export).
- Event sourcing: the event store holds roughly 40 million domain events over 7 years (about 15,600 events/day on average for a mid-size ledger). A new read model, say a fourth "tax reporting" view added in year 5, is built by replaying those 40 million events once, deterministically, against the new projection logic; correctness is verifiable because the same event log produces the same output every time it is replayed.
- Relational + CDC: the same 7-year retention means either keeping 7 years of database write-ahead log history available to the CDC pipeline (expensive and often beyond what most databases retain by default) or accepting that a full historical rebuild of a new view is not actually possible from CDC alone, only from that point forward. This is the concrete cost of CDC's weaker replay guarantee versus a purpose-built event store: it shows up exactly when the business asks for a new view of old data.
Trade-offs and pitfalls
- Do not choose event sourcing for its audit-trail marketing value alone; a relational database with proper CDC and immutable audit columns can satisfy many audit requirements at a fraction of the operational cost.
- Do not underestimate CDC's replay ceiling: if "rebuild any read model from any point in history" is a hard requirement, verify the source database's log retention actually supports it before committing to CDC as the long-term answer.
- Senior signal: naming the retention and rebuild requirement in concrete terms (how far back, how many read models, how often they change) before picking a side, rather than treating this as a purely stylistic architecture preference.
How would you choose a Kafka partitioning key for a user-events topic such that ordering is preserved per user but partitions are balanced across the cluster? Discuss hashing strategies, handling hot users, and approaches for multi-tenant fairness when some tenants generate far more events than others.
Sample Answer
Direct answer
Use user_id itself as the partition key so Kafka's partitioner (a mod-based hash of the key) always routes every event for a given user to the same partition, which is what gives per-user ordering, it falls directly out of using the same key for the same entity, not something built separately. Balance across the cluster then comes from having enough distinct keys and a reasonably uniform hash, with a deliberate, separate remediation for the minority of disproportionately hot keys.
Structured elaboration
Hashing strategy. Kafka's default partitioner computes hash(key) % num_partitions (using a murmur2-based hash internally); the exact hash function matters less than the property that it is deterministic, the same key always maps to the same partition number for a fixed partition count. Because the partition assignment is a pure function of the key, ordering per key is automatic once you commit to using that key consistently, not an extra mechanism you build on top.
Why "balanced" is a statistical property, not a guarantee. With many distinct user ids and a well-distributed hash, the law of large numbers makes total load spread roughly evenly across partitions, if every key generates traffic at roughly the same rate. The question specifically names the case where that assumption breaks: hot users.
Handling hot users. Two remediations, with different costs:
- Key salting (sub-partitioning): append a bounded suffix to a hot user's key only once their rate crosses a threshold (
user_id#0throughuser_id#K-1), spreading that one user's traffic across up to K partitions. The explicit cost: per-user ordering across the salted shards is no longer guaranteed, only ordering within each shard is. This is a real trade-off to name out loud, not a free win, if the consuming system genuinely needs strict per-user order even for hot users, salting defeats that and a downstream re-sequencing step would be needed, which mostly defeats the purpose of salting in the first place. - Isolating known hot keys: if hot users are identifiable ahead of time (service accounts, bots, high-volume enterprise integrations), route them to a separate topic or partition set with its own capacity plan, keeping the main topic's "roughly uniform" assumption valid for the long tail of ordinary users.
Multi-tenant fairness. Partitioning purely by user_id does nothing to stop one tenant (an account with many users) from occupying a disproportionate share of partitions and consumer capacity relative to a smaller tenant, since partition assignment has no concept of tenant at all. If fairness across tenants is a hard requirement: give large or noisy tenants a dedicated topic or partition range so they cannot starve others' consumer lag budget, and layer per-tenant rate limiting or quota enforcement on the consumer or producer side (the same token-bucket idea used for rate limiting generally), since a static partition assignment cannot adapt to traffic that changes over time the way an explicit quota can.
Worked example
A small, executed illustration of the mechanism (using Python's built-in zlib.crc32 purely as a stand-in deterministic hash function to demonstrate the mod-based mechanic; Kafka's actual default partitioner uses a different, murmur2-based hash internally, the mechanism, not the specific hash, is what this illustrates).
import zlib
num_partitions = 6
for uid in ["user-1042", "user-8831", "user-2207", "user-77519"]:
h = zlib.crc32(uid.encode())
print(uid, "partition=", h % num_partitions)
user-1042 partition= 3
user-8831 partition= 4
user-2207 partition= 4
user-77519 partition= 0
Every event for user-1042 always lands on partition 3, giving that user strict per-user order for free. Now simulate salting the hot user user-8831 into 12 sub-keys (user-8831#0 through user-8831#11):
counts = {}
for shard in range(12):
key = f"user-8831#{shard}"
p = zlib.crc32(key.encode()) % num_partitions
counts[p] = counts.get(p, 0) + 1
print(counts)
{4: 5, 0: 1, 5: 5, 1: 1}
The 12 salted sub-keys land on 4 of the 6 partitions instead of all landing on partition 4 as the unsalted key would. The spread is not perfectly even at this small sample size (partitions 4 and 5 got 5 each, partitions 0 and 1 got 1 each), which is expected statistical noise at n = 12, not a flaw in the technique; at real production volumes for a genuinely hot key, the same law of large numbers that balances ordinary keys also evens out a hot key's salted shards.
Trade-offs and pitfalls
- Salting trades away strict per-user order for the salted user specifically, in exchange for spreading their load; state this cost explicitly when proposing it, do not present it as a free fix.
- Changing the partition count later reshuffles the mod-based assignment for every existing key, not just new ones, since
hash(key) % num_partitionschanges for essentially every key whennum_partitionschanges. Under-provisioning partition count up front makes a later increase a disruptive, coordinated migration, not an incremental capacity add. - "Balanced across the cluster" and "fair across tenants" are different guarantees: the first is about hash uniformity over many keys, the second is about the business meaning of who owns which keys, no hashing scheme fixes tenant fairness on its own.
- Picking the wrong entity granularity for the key, for example partitioning by
session_idwhen the actual ordering requirement is per-user across sessions, silently breaks the guarantee the system was supposed to provide.
Explain the difference between publish-subscribe and point-to-point (producer-consumer) messaging patterns. Provide concrete scenarios where pub/sub is a better fit (e.g., notifications, analytics) and where queues are preferable (e.g., work queues, task processing), particularly in multi-tenant SaaS and event-driven microservice architectures.
Sample Answer
Direct answer
Publish-subscribe (pub/sub) delivers each message to every interested subscriber, so it fits situations where multiple independent parties need to know the same fact happened, such as notifications or analytics. Point-to-point (producer-consumer, work-queue style) delivers each message to exactly one consumer among a pool, so it fits situations where a unit of work must be done exactly once by whichever worker picks it up, such as background task processing. The distinction is about fan-out (one-to-many awareness) versus load distribution (one-of-many execution), and both patterns are commonly implemented on the same underlying broker.
Structured elaboration
Publish-subscribe. A publisher emits an event to a topic; every subscriber with an active subscription receives its own copy. Subscribers are typically unaware of each other, can be added or removed without changing the publisher, and each one processes the event for its own purpose. This is the right model whenever "N different systems need to react to the same fact" is the actual requirement: a "UserSignedUp" event might be consumed by an email-welcome service, an analytics pipeline, and a fraud-scoring service simultaneously, with none of them competing for the message.
Point-to-point (work queues). A producer places a task on a queue; a pool of competing consumers pulls from the same queue, and each task is handled by exactly one consumer. This is the right model for "this unit of work needs to happen once, by whichever worker is free," such as resizing an uploaded image or sending a single transactional email: you do not want three workers all resizing the same image.
Where pub/sub is the better fit. Notifications is the clearest case: a single "OrderShipped" event needs to reach a push-notification service, an SMS service, and an in-app activity feed, each independently, and adding a fourth channel later should not require touching the producer. Analytics is the same shape: every business event (page view, purchase, signup) typically needs to reach an analytics pipeline in addition to whatever else consumes it, without competing with those other consumers for the message.
Where queues are preferable. Work queues and task processing are the clear case: a video-transcoding job, a report-generation job, or an outbound-email send should be picked up and completed by exactly one worker, with the queue's competing-consumers model providing natural load balancing and horizontal scaling (add more workers, they compete for the same backlog) without any risk of duplicate execution beyond what at-least-once delivery already requires the consumer to handle idempotently.
Multi-tenant SaaS (software as a service) and event-driven microservices. In a multi-tenant SaaS system, pub/sub is what lets independently-owned services (billing, usage-metering, audit logging) all react to the same tenant-level event, such as "SubscriptionUpgraded," without the team that owns the upgrade flow needing to know or coordinate with every downstream consumer; new consumers subscribe without any change to the publisher. Point-to-point queues, by contrast, are what those same microservices use internally for their own background work, such as a billing service's queue of pending invoice-generation tasks, where exactly-once-effective execution by one worker in the pool is the requirement, not fan-out to observers.
Worked example
A multi-tenant SaaS platform publishes a "TenantUpgraded" event when a customer moves from a free to a paid plan. Three independent subscribers exist on this topic: a billing service that starts metered invoicing, a feature-flag service that unlocks paid features, and a customer-success service that triggers an onboarding email sequence. All three receive their own copy of the same event; the team that owns the upgrade flow never had to know these three consumers existed. Separately, the feature-flag service's own onboarding-email trigger enqueues an actual "send welcome email" task onto a point-to-point work queue consumed by a pool of 5 worker processes; only one of those 5 workers ends up sending that specific email, because the queue hands each task to a single competing consumer, not to all 5.
Trade-offs and pitfalls
The common mistake is using a work queue where pub/sub was needed: if a "TenantUpgraded" task were placed on a single point-to-point queue instead of published to a topic, only one of billing, feature-flags, or customer-success would ever see it, and the other two would silently never fire, which is a subtle and easy-to-miss integration bug. The opposite mistake is using pub/sub where a work queue was needed for a task that must be done exactly once: if "resize this uploaded image" were published to a topic with multiple subscribed workers, every worker would independently resize the same image, wasting resources and, if the workers write to the same output path, potentially racing each other. A senior answer names this fan-out-versus-load-distribution distinction explicitly, rather than treating "pub/sub" and "queue" as interchangeable synonyms for "asynchronous messaging."
Explain the difference between message queues and event streams using: (1) simple high-level definitions, (2) step-by-step architecture differences (point-to-point vs publish-subscribe, retention and consumer offsets), (3) concrete use cases (task queue for background jobs vs event log for analytics), (4) discuss delivery semantics (at-least-once, at-most-once, exactly-once) and operational trade-offs.
Sample Answer
Direct answer
A message queue is a transient work-distribution mechanism: a message is delivered to one consumer and then typically removed. An event stream is a durable, ordered, append-only log that many consumers can read and re-read independently, each tracking their own position. Use a queue for background jobs that need to happen once; use an event stream when multiple consumers need to replay or independently process the same history of events, such as for analytics.
Structured elaboration
(1) Simple high-level definitions. A message queue holds discrete units of work; once a consumer successfully processes (and acknowledges) a message, it is generally gone from the queue. An event stream holds an ordered, append-only sequence of events that persists for a configured retention period regardless of whether any consumer has read them yet, and multiple independent consumers can read the same stream at their own pace.
(2) Step-by-step architecture differences.
- Delivery model: a queue is point-to-point: competing consumers pull from the same backlog, and each message goes to exactly one of them. An event stream is publish-subscribe at the storage layer: every consumer group gets its own full view of the log, so 5 independent consumer groups can each read every event, not compete for it.
- Retention: a queue typically discards a message once acknowledged (or after a short retention window for safety); an event stream retains events for a configured duration (hours to indefinitely) independent of consumption, which is what makes replay possible.
- Consumer offsets: a queue broker tracks per-message state (delivered, in-flight, acknowledged) on the broker's behalf; a consumer of an event stream tracks its own offset (a pointer into the log marking "I have processed up through here"), which the consumer can reset backward to reprocess history or forward to skip it, something a queue's consume-once model does not support.
(3) Concrete use cases. A task queue for background jobs (resize an uploaded image, send a single email, generate a PDF) is the queue use case: each job runs once, and once done there is no reason to revisit it. An event log for analytics (every purchase event, retained for 30+ days, read independently by a real-time dashboard consumer, a daily batch-aggregation consumer, and a fraud-detection consumer, each at their own offset) is the event-stream use case: the same history needs to serve multiple, independently-paced readers, and a new consumer added next month should be able to read events from before it existed.
(4) Delivery semantics and operational trade-offs. Both queues and streams commonly offer at-least-once delivery (a message may be delivered more than once, so consumers must be idempotent) as the practical default, since it is the easiest guarantee to build reliably. At-most-once (a message may be silently lost but never duplicated) is rarely chosen deliberately; it falls out of not retrying failed deliveries and is generally the wrong default for anything that matters. Exactly-once, from the consumer's experience, is achievable at the application layer by combining at-least-once delivery with a consumer-side idempotency check (a deduplication store keyed by message ID with a time-to-live, TTL), rather than relying on a broker guarantee alone; broker-level exactly-once semantics (EOS) protocols exist on some platforms but their internals are a streaming-platform implementation detail, not something the consuming application needs to reason about. Operationally, queues are simpler to run (bounded backlog, straightforward monitoring of queue depth) but lose replay ability; event streams give you replay and multiple independent consumer groups at the cost of needing retention-size planning, offset-lag monitoring per consumer group, and generally a heavier broker to operate.
Worked example
An e-commerce platform's "PurchaseCompleted" event needs to reach a real-time inventory-decrement consumer (must process within seconds), a nightly analytics-aggregation job (reads the whole day's events once, in batch), and a newly-added fraud-review consumer (needs to backfill and reprocess the last 7 days of purchases to retrain a rule set). This is only possible on an event stream: the inventory consumer reads from the current offset forward, the analytics job resets its offset to the start of each day, and the fraud-review consumer resets its offset back 7 days on rollout, all three reading the same retained log independently. If "PurchaseCompleted" had instead been placed on a point-to-point queue, only one of the three would ever have received a given event, and none of them could replay history that had already been consumed.
Trade-offs and pitfalls
The most common pitfall is reaching for an event stream by default because it seems more capable, then paying its operational cost (retention planning, per-consumer-group offset-lag monitoring, a heavier broker) for a workload that was really a simple once-only task queue. The opposite pitfall is placing something that needed independent replay, like the fraud-review backfill above, onto a queue and discovering months later that the history needed to replay it was already gone. A senior answer treats "does more than one independent party need to read this, possibly at different times or replaying history" as the deciding question, not "which technology is newer or more scalable."
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.
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.