Advanced module · event-driven systems theory
Event-Driven Architecture Deep Dive: the theory underneath your Kafka topics.
Event notification, event-carried-state-transfer, and event sourcing are three distinct patterns interviewers expect you to tell apart — plus CDC, delivery semantics, schema evolution, ordering, and dead-letter handling: the mechanisms this site's mini projects already lean on without ever naming them.
Why this is the missing theory chapter
This site's e-commerce and support-ticket mini projects already publish events, and the payment mini project already builds a transactional outbox to publish them reliably — but none of those pages stop to name the underlying patterns, or explain what "reliable" and "in order" and "exactly-once" precisely mean once you're no longer inside one database's transaction boundary. This module is that missing chapter: the vocabulary and the failure modes that sit underneath every Kafka-based design on this site.
The vocabulary interviewers expect you to use precisely
Where this shows up as a production incident
Jump to a section
Three patterns hiding under one word: "event"
Almost every system design interview eventually says "and then it publishes an event" as if that settles the design. It doesn't — it opens three genuinely different questions: does the event carry data or just a pointer, do consumers keep their own copy of that data or fetch it fresh, and is the event log itself the source of truth or just a delivery mechanism. Those questions have three standard answers.
Event notification — "something happened, go fetch it yourself"
The event carries an id and a type, nothing more: OrderStatusChanged{orderId: 4471}. Any consumer that needs the actual details makes a call back to the owning service to get them. This keeps events tiny, keeps the source service in complete control of what data it exposes and to whom, and avoids ever shipping stale duplicated fields — but it quietly reintroduces a synchronous dependency at read time. If the owning service is slow or down, every consumer trying to react to the notification is blocked on it, which is exactly the runtime coupling event-driven design is usually adopted to remove.
Event-carried state transfer — the "fat event"
The event embeds the full (or materially complete) state a consumer needs: OrderPlaced{orderId, items[], total, shippingAddress, ...}. Consumers never call back to the source at all, which is what actually decouples them at runtime — a consumer can process a backlog of these events even if the producing service is completely down, because everything it needs already arrived with the event. The cost is data duplication: every consumer now holds a copy of data it doesn't own, that copy can drift from the source if an event is missed, and the payload grows every time a consumer needs one more field, coupling the event schema to everyone's read needs at once.
Event sourcing — the log is the source of truth
Here the pattern goes one step further: instead of state living in a table that's occasionally mutated, the sequence of events is the persisted data, and any current-state view (a balance, a status) is a projection derived by replaying that sequence. This buys a complete, free audit trail, the ability to answer "what was true at any past point in time," and the option to build new read models later just by replaying history you already have. It also means every read either replays a potentially long history or depends on a projection someone has to build and keep in sync, and every schema change has to remain readable by a replay running against events written years ago — a much stricter version of the compatibility problem covered later in this module.
| Pattern | Payload | Runtime coupling | Storage growth | Pick it when |
|---|---|---|---|---|
| Event notification | Thin — id + type | Consumer still depends on source being reachable | Minimal | Payload is sensitive/huge, or consumers are rare and can tolerate a callback |
| Event-carried state transfer | Fat — embedded state | None at read time — consumer is self-sufficient | Duplicated across every consumer | Multiple, high-volume, or offline-tolerant consumers — the common default for service integration |
| Event sourcing | The full change history, kept forever | None — state is derived, not fetched | Unbounded, needs snapshotting at scale | Audit/temporal-query value genuinely justifies the replay and schema-forever cost |
OrderPlaced event carrying full order details is event-carried state transfer — inventory, shipping, and notification all need the order immediately and in volume, so a callback-per-consumer model would recreate the exact coupling the event bus exists to remove. The support-ticket triage project's routing events, which announce a decision and let the owning service be queried for full context on demand, lean closer to event notification. Neither page named the pattern; naming it is what turns "I've used Kafka" into "I can defend why I shaped the event the way I did."ledger_entries table resembles an event log, but the wallet's balance is still a directly maintained, directly read column, not something rebuilt from scratch by replaying every entry on every read — reconciliation independently verifies the cached column against the log instead. That hybrid, an authoritative append-only history plus a maintained read-optimized projection, is closer to how event-sourced systems actually run at scale than the textbook "replay everything, always" version anyway.Change Data Capture: reading the log instead of writing twice
The dual-write problem is simple to state and easy to get wrong: a service that writes to its database and then separately publishes an event has two independent operations with no shared transaction, so a crash between them either loses the event (data changed, nobody told) or, if the order is reversed, publishes an event for a write that then fails to commit. The payment mini project's transactional outbox solves this at the application layer by writing the event to the same database, in the same transaction, as the state change itself. CDC solves the identical problem one layer down, in infrastructure instead of application code.
| Transactional outbox | Change Data Capture | |
|---|---|---|
| Where the fix lives | Application code — an outbox table + a relay process, written once per service | Infrastructure — a connector reading the database's own commit log, no app changes |
| Discipline required | Every future code path that changes state must remember to also write the outbox row | None — every commit is captured automatically, nothing to forget |
| Event shape | Exactly what the application chose to write — clean, versioned, intentional | Derived from a raw row diff — usually needs a transform step, can leak internal column names |
| Operational cost | One extra table + one scheduled job, on infrastructure you already run | Kafka Connect cluster, replication slot management, WAL retention pressure on the database |
| Best fit | A service you control, that wants a deliberate, stable event contract | A legacy database you can't or won't modify, or capturing every table's changes without relying on developers to remember |
// Illustrative Debezium Postgres connector config -- NOT used for payment_db,
// which already solves its dual-write problem with the outbox relay.
{
"name": "wallet-db-cdc",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"database.hostname": "postgres-wallet",
"database.dbname": "wallet_db",
"plugin.name": "pgoutput",
"slot.name": "wallet_db_slot",
"table.include.list": "public.ledger_entries",
"topic.prefix": "cdc.wallet",
"tombstones.on.delete": "false"
}
}
kafka-connect/ ├── connectors/ │ ├── wallet-db-cdc.json // hypothetical -- ledger_entries is append-only, a good CDC fit │ └── payment-db-cdc.json // NOT deployed -- payment_db already has a working outbox ├── docker-compose.debezium.yml └── README.md // when to reach for this vs. the outbox already in the payment project
Delivery semantics: what "exactly-once" quietly means
Every message broker has to make a choice at the boundary between "did the consumer finish processing?" and "should I redeliver?" — and that choice produces one of three semantics. Only one of them is achievable without extra work, and it isn't the one most people assume.
| Semantic | Mechanism | Failure behavior | Use it when |
|---|---|---|---|
| At-most-once | Commit the offset before or without waiting on processing | A crash mid-process silently loses the message — no redelivery, ever | Losing an occasional message is cheaper than ever double-processing one (rare) |
| At-least-once | Commit the offset only after processing succeeds | A crash after processing but before commit causes redelivery and possible duplication | The default choice almost everywhere, paired with idempotent consumers |
| "Exactly-once" (effectively-once) | At-least-once delivery + idempotent processing keyed on a stable event id | Duplicates still arrive; the second attempt is a safe no-op instead of a repeat side effect | Any consumer whose side effect must never repeat — the practical default for correctness-sensitive work |
// Idempotent consumer -- same discipline as the payment project's processed_transfers table,
// applied to a Kafka listener instead of an HTTP retry.
@KafkaListener(topics = "payment.completed")
@Transactional
public void onPaymentCompleted(PaymentCompletedEvent event) {
if (processedEventRepository.existsById(event.getEventId())) {
return; // already handled -- redelivery is a no-op, not a bug
}
notificationService.send(event);
processedEventRepository.save(new ProcessedEvent(event.getEventId(), Instant.now()));
// save() shares the transaction with any DB side effects above --
// the ack/commit for this record only happens after this method returns.
}
Schema evolution: backward, forward, and full compatibility
A Kafka topic outlives any single version of the code producing or consuming it — producers and consumers deploy independently and are never all on the same schema version at the same instant during a rollout. Schema compatibility rules exist to make that gap survivable, and Avro and Protobuf both enforce them the same conceptual way even though the wire formats differ.
| Change | Backward compatible? | Forward compatible? | Why |
|---|---|---|---|
| Add optional field with a default | Yes | Yes | Old data parses fine (default fills the gap); old consumers ignore the unknown field |
| Add required field, no default | No | Yes | Old data has no value for it -- a new consumer reading old records has nothing to fill in |
| Remove a field | Depends | No | Backward-safe only if the field had a default; forward always breaks a consumer still expecting it |
| Rename a field | No | No | Treated as remove + add -- breaks both directions unless the format supports an explicit alias |
| Widen a numeric type (int → long) | Yes | Usually | Old values promote cleanly; most serializers handle the widening transparently |
| Narrow a numeric type (long → int) | No | No | Existing large values can silently truncate or fail to parse |
// Avro: adding "priority" with a default is backward AND forward compatible.
// Schema v1
{"type":"record","name":"OrderPlaced","fields":[
{"name":"orderId","type":"string"},
{"name":"total","type":"double"}
]}
// Schema v2 -- safe evolution
{"type":"record","name":"OrderPlaced","fields":[
{"name":"orderId","type":"string"},
{"name":"total","type":"double"},
{"name":"priority","type":"string","default":"STANDARD"} // new field, has a default
]}
// Schema v3 -- BREAKING: orderId removed with no default, rejected by a
// registry enforcing BACKWARD or FULL compatibility before it ever publishes.
{"type":"record","name":"OrderPlaced","fields":[
{"name":"total","type":"double"},
{"name":"priority","type":"string","default":"STANDARD"}
]}
Ordering guarantees: per-partition, not per-topic
Kafka only promises order within a single partition: messages with the same key are always routed to the same partition (by the default hash partitioner) and are always delivered to a consumer in the order they were written. Nothing is promised about order across different partitions, and a topic with more than one partition offers no topic-wide ordering at all — that's not a bug, it's the mechanism that lets a topic be consumed in parallel.
This is exactly why the payment project's downstream events matter to key correctly: keying on walletId (or transferId) guarantees one wallet's history is always processed in commit order by reconciliation and notification, while two different wallets' events can arrive in either relative order with zero consequence, because reconciliation checks each wallet independently. Picking the partition key is picking exactly which entity's order you're willing to guarantee.
| Requirement | Strategy | Trade-off |
|---|---|---|
| Order matters within one entity only | Key by that entity's id (default case — almost always sufficient) | None — this is the normal, scalable design |
| Order matters across every message in the topic | Single partition | Throughput capped at one consumer thread — rarely acceptable at scale |
| Order matters across a few related entities | Coarser key grouping the related entities together | Hot-partition risk concentrated on that key |
| Ordering requirement can be relaxed instead | Make each event self-contained enough to process correctly regardless of arrival order | Pushes toward fatter, state-transfer-style events -- often the better fix |
Dead-letter queues and reprocessing a poison message
A poison message is one a consumer can never successfully process no matter how many times it's redelivered — a malformed payload, a field the code doesn't defensively handle, a downstream dependency that will 500 for this specific record forever. Because Kafka can't skip ahead within a partition without breaking the ordering guarantee for everything behind it, a consumer that keeps re-throwing on the same message doesn't just fail on that message — it blocks every message after it in that partition, indefinitely.
// Bounded retry -> dead-letter, tracked via a header incremented per attempt.
@KafkaListener(topics = "order.events")
public void onOrderEvent(ConsumerRecord<String, String> record) {
int attempt = headerInt(record, "x-retry-count", 0);
try {
process(record.value());
} catch (TransientException e) {
if (attempt >= MAX_ATTEMPTS) {
deadLetterProducer.send(toDeadLetterRecord(record, e, attempt));
// offset for the ORIGINAL record still commits normally here --
// the partition is unblocked once the DLQ write succeeds.
return;
}
throw e; // under MAX_ATTEMPTS -- let the container retry with backoff
} catch (PermanentException e) {
deadLetterProducer.send(toDeadLetterRecord(record, e, attempt)); // don't waste retries on a permanent failure
}
}
Reprocessing has to respect the same discipline as everything else in this module: because consumers are already idempotent under at-least-once delivery, a manual replay of a fixed DLQ message back onto the source topic is safe by construction, not by luck. What isn't safe is a blanket replay of an entire DLQ after a bug fix without checking that every message in it actually failed for the reason that was fixed — and replayed messages land at the tail of the topic with a new position, so anything genuinely order-sensitive relative to other events should be re-injected deliberately, not dumped back in bulk.
Key design decisions and interview talking points
Event notification vs event-carried-state-transfer -- what's the actual difference and why does it matter which one you pick?
An event notification says only that something happened and carries an id -- any consumer that needs details must call back to the source service to fetch them. An event-carried-state-transfer event embeds the full (or materially complete) new state in the payload, so consumers never call back at all. The choice is really a choice about coupling: notification keeps events small and the source service in full control of its data, but it recreates a synchronous dependency at read time; state transfer removes that runtime coupling entirely but means every consumer is now carrying a copy of data it doesn't own, which can drift or leak fields it shouldn't.
When would you choose event sourcing over the other two patterns, and what's the cost?
Event sourcing is the right call when the sequence of changes is itself valuable -- audit trails, temporal queries ("what was the balance at 3pm Tuesday"), or the ability to rebuild entirely new read models from history you didn't know you'd need yet. The cost is real: every read either replays events or depends on a maintained projection, schema evolution has to account for replaying old events forever, and the team has to actually think in events, not rows. Most systems get most of the benefit from a well-designed append-only log plus a couple of good projections, without going all the way to "the log is the only persisted truth."
This site's e-commerce project publishes an OrderPlaced event with the full order payload -- which pattern is that, and why not plain event notification?
That's event-carried-state-transfer: inventory, shipping, and notification services all need order details immediately and in volume, so making every one of them call back to Order Service for the full order on every event would turn a fan-out of independent consumers back into a fan-out of synchronous dependencies on one service -- the exact coupling event-driven design was supposed to remove. Notification would only make sense there if consumers were rare, needed only an id most of the time, or the payload contained something too sensitive or too large to duplicate everywhere.
Why does CDC solve the dual-write problem more elegantly than an application-level outbox in some cases?
The outbox pattern still requires a developer to remember, for every code path that changes state, to also write an outbox row in the same transaction -- it's reliable once done correctly, but it's an application-level discipline that has to be applied everywhere, forever. CDC tails the database's own write-ahead log or binlog, which every committed transaction already passes through whether or not anyone remembered an outbox write, so it captures changes automatically and can't be silently skipped by a new code path that forgot the pattern.
What's the catch with CDC -- is it strictly better than the outbox pattern?
No -- CDC trades application-level discipline for infrastructure complexity and less control over event shape. You're now operating a connector (typically Kafka Connect plus Debezium), managing a replication slot the database must retain WAL for, and deriving your event schema from raw row diffs instead of a clean domain event an application chose to emit -- which tends to leak internal column names and requires a transform step to turn into something consumers should actually depend on. The outbox stays attractive specifically because the application controls exactly what an event looks like, and running it needs nothing beyond the database you already have.
How does Debezium actually get changes out of Postgres without polling the table?
It creates a logical replication slot and registers as a logical replication consumer, the same mechanism Postgres uses to stream changes to a replica -- so it receives a continuous stream of committed row-level changes (insert/update/delete, with before/after images depending on REPLICA IDENTITY) directly from the write-ahead log, in commit order, with no query against the table at all. That's why it can capture every change with low latency and without adding read load to the table itself, which a poll-the-table-for-an-updated_at-column approach cannot do without missing deletes or adding query overhead.
Could you use CDC directly on wallet_db from the payment mini project instead of the outbox table it already has?
Technically yes, but it would be a downgrade there specifically: the outbox table already gives Payment Service an explicit, versioned, clean event shape and needs nothing more than the Postgres it already runs. Swapping in CDC would mean standing up Kafka Connect and Debezium, managing a replication slot against a database that must never fall behind on retention, and deriving payment.completed from a raw wallets/ledger_entries row diff instead of the deliberate event Payment Service already constructs -- more moving parts to solve a dual-write problem the outbox already solved cleanly.
What does "exactly-once" actually mean in Kafka, and why should you be skeptical of vendor claims?
In practice, exactly-once almost always means idempotent processing layered on top of at-least-once delivery -- the message can and sometimes will be delivered more than once, but processing it twice produces the same result as processing it once, so it looks exactly-once from the outside. True exactly-once transport, where a message is guaranteed to cross an unreliable network and land exactly one time with no possibility of duplication or loss, isn't achievable in a general distributed system -- any claim of it should be read as either narrowly scoped (Kafka's own transactions, entirely within Kafka) or marketing.
Concretely, how do you build idempotent consumers when delivery is at-least-once?
Give every event a stable, unique id (an aggregate id plus a version, or a UUID assigned at creation) and record processed ids in a table with a unique constraint before or as part of doing the side-effecting work, in the same transaction where possible -- exactly the discipline the payment project's processed_transfers table already applies to HTTP retries. When the same event arrives again, the insert fails or the lookup finds it, and the consumer returns the prior result instead of repeating the side effect, so redelivery becomes a no-op instead of a bug.
What's the difference between Kafka's idempotent producer and end-to-end exactly-once semantics (transactions)?
The idempotent producer only prevents a single producer's own retries from creating duplicate records within a partition, by tagging each batch with a producer id and sequence number the broker can deduplicate -- it says nothing about consumers or about writes spanning multiple partitions or topics. Kafka transactions build on top of that to atomically commit a set of produced records and a set of consumed offsets together, which is what a read-process-write consumer needs for exactly-once within Kafka -- but the instant that pipeline talks to anything outside Kafka (a database, an HTTP call), you're back to needing idempotency, because that external system isn't part of the transaction.
What is backward compatibility in Avro/Protobuf schema evolution, concretely?
Backward compatible means a consumer running the new schema can correctly read data that was written with the old schema -- so you can deploy consumers first. In practice that means new fields must have defaults (so old records without them still parse) and you must never remove or repurpose a field a still-deployed producer might still be writing. It's the direction that matters most for rolling deploys, because consumers typically need to update ahead of producers to be ready for what's coming.
What is forward compatibility, and which one do you actually need day to day?
Forward compatible means a consumer on the old schema can still read data written with the new schema, which requires new fields to be safely ignorable by code that's never heard of them and forbids repurposing an existing field's meaning. Backward compatibility matters more often in practice because consumers usually deploy ahead of producers during a rolling rollout, but forward compatibility matters whenever a producer might ship before every consumer has upgraded -- which is common enough that most teams should just target full compatibility (both directions) as the default and only relax it deliberately.
What breaks if you remove a required field from a Protobuf message without a default, and a consumer is still on the old schema?
The old consumer's generated code expects that field to be present and typically either fails to parse the message, silently reads a zero-value default it can't distinguish from a legitimate zero, or throws downstream when code that assumed the field exists dereferences a missing value -- none of which is a clean failure. This is exactly why compatibility rules exist: removing a required field without a default breaks backward compatibility outright, and a schema registry configured to enforce compatibility should reject that change at publish time rather than let it ship and fail in production.
Why use a schema registry instead of just versioning topics like order-events-v1 and order-events-v2?
Topic-per-version works but pushes the entire migration cost onto every consumer, who must now dual-subscribe, and every producer, who must dual-write during the transition, and it does nothing to stop an incompatible schema from being published in the first place. A schema registry checks a new schema against the configured compatibility rule (backward, forward, or full) at write time and rejects the write if it would break existing consumers, catching the mistake at the producer before it ever reaches a topic instead of after a consumer crashes in production.
A consumer hits a poison message it can never process -- walk through what happens without a dead-letter queue.
The consumer throws, the offset for that message is never committed, and on the next poll the broker redelivers the exact same message, because Kafka can't skip ahead within a partition without breaking the ordering guarantee for everything behind it. The consumer fails again, retries again, forever -- and because that message is stuck at the head of its partition, every message after it in that partition is also stuck, even ones that would have processed perfectly fine. One bad record halts an entire partition's throughput indefinitely.
What's the risk of blindly replaying everything from a DLQ back onto the main topic?
Replayed messages land at the tail of the topic with today's timestamp and partition assignment, not their original position relative to other events -- so if ordering between that entity's events and others mattered, the replay can process out of the original sequence. It also silently retries everything in the DLQ as if the underlying cause were fixed; if the fix only addressed some of the failures, you've just moved the poison messages back into the live pipeline to fail again, consuming capacity and re-triggering the same alerting, so replay should be selective and verified, not a blanket dump-it-back-in.
If an interviewer asks how you'd prevent a poison message from causing an infinite retry loop, what do you say?
Bound the retries -- a small number of immediate or backed-off attempts to absorb genuinely transient failures, then stop retrying in place and route the message, with its failure reason and attempt history attached, to a dead-letter topic while committing the original offset so the partition keeps moving. The retry count has to be bounded and tracked (a header incremented per attempt) specifically because an unbounded retry loop on a permanent failure -- a downstream reference that was actually deleted, not a network blip -- will spin forever without ever succeeding, burning consumer capacity for no benefit.
Kafka guarantees ordering per partition -- what breaks that guarantee in practice?
Two things: using a key that doesn't consistently map the same logical entity to the same partition (or publishing with no key at all, which round-robins across partitions), and assuming order holds across different keys or across a topic as a whole, which Kafka never promised. A rebalance mid-processing can also complicate redelivery if the consumer isn't careful about committing offsets in order. The fix is almost always making the partition key the entity whose internal event order actually matters, and accepting that different entities' events can and will interleave in any order.
You need strict ordering across multiple entity types that currently land in different partitions -- how do you get it?
The blunt option is a single-partition topic, which guarantees total order at the cost of capping throughput to one consumer thread -- rarely acceptable at scale. A better fix is usually to question whether cross-entity ordering is truly required or whether the consumer can be made to tolerate interleaving by carrying enough state in each event (moving toward event-carried-state-transfer) to process correctly regardless of arrival order. If ordering genuinely can't be relaxed, group the related entities under one coarser partition key so they land together, accepting the resulting hot-partition risk as the trade-off for the ordering guarantee.
How does the payment project's wallet-keyed event partitioning actually keep reconciliation correct?
Keying payment and reconciliation events on walletId guarantees any single wallet's events are always delivered to a consumer in the exact order they were committed, which is all reconciliation needs -- it recomputes and compares one wallet at a time, so it never depends on how that wallet's events are ordered relative to a different wallet's events. That's the ordering guarantee actually required in that system, and it's exactly the guarantee per-partition ordering provides for free once the key is chosen correctly.
If an interviewer asks for the single biggest gap between "I've used Kafka" and "I understand event-driven architecture," what do you say?
Knowing which of the three event patterns you're actually publishing (notification, state transfer, or sourcing) and why, versus just calling everything "an event" -- because that choice quietly determines your coupling, your payload design, and your replay story, and every other topic here (CDC vs outbox, delivery semantics, schema compatibility, ordering, dead-letter handling) is really just the operational machinery for making whichever pattern you picked actually reliable in production.
Related guides
Update these hrefs to your published Blogger post URLs once each page is live.
Post a Comment
Add