Technologies that implement this pattern: Apache Kafka · PostgreSQL · DynamoDB · Elasticsearch · Apache Flink · Redis
Why This Matters#
Most systems store the answer and throw away the question. A row says the account balance is $412.50; it does not say that it got there through 1,840 deposits, withdrawals, reversals and one chargeback that was itself reversed. For a profile page that is fine. For a ledger, an order lifecycle, a trading book or a licence service it is the difference between "we can explain every number on the screen" and "we restored from backup and hope the support tickets stop." Event sourcing keeps the question: the append-only sequence of facts is the system of record, and current state is a fold over it. CQRS — command query responsibility segregation — is the companion move: the model that decides whether a write is allowed is not the model that serves reads, so each can be shaped for its job.
Most candidates treat this as an architecture style question: "use CQRS and event sourcing for scalability," drawn as two boxes and a Kafka topic. Staff engineers treat it as a lifetime-cost question. Event sourcing is a one-way door with three bills that arrive late: every read needs a projection someone owns, every event schema you ever publish must be readable forever, and every rebuild of a read model is a batch job over the entire history. The question is not can we model this as events — you can model anything as events — but does this domain need its history as a first-class asset badly enough to pay for projections, replays and versioning for the next ten years?
The second reframe: CQRS and event sourcing are separable, and most systems that need one do not need the other. Splitting a write model from read models — a normalized PostgreSQL schema feeding an Elasticsearch index and a Redis leaderboard through CDC — is ordinary, cheap and reversible. Making the event log the source of truth is rare, expensive and nearly irreversible. Candidates who weld them together signal they have read about the pattern and not run it.
If you can walk an interviewer from "which decisions need history" to "where the consistency boundary sits" to "how a projection is rebuilt without downtime" to "what happens to an event we published three years ago when the meaning of a field changes," you are answering at Staff level.
The 60-Second Version#
- CQRS without event sourcing is the common case. One write store, several read models fed by CDC or an outbox. Read models lag the write by ~100ms–2s; that lag is the price, and it must be named to product before it surprises a user who cannot see their own edit.
- The aggregate is the consistency boundary, and it should be small. One stream per aggregate (an order, an account), optimistic concurrency on the stream's expected version. Aim for streams under ~1,000 events; past ~10K, load time per command grows into tens of milliseconds and snapshots stop being optional.
- Snapshots are a cache, never a source of truth. Snapshot every ~100–500 events or when load time crosses ~10ms. A snapshot you cannot regenerate from events is a second, unversioned system of record.
- Projections are rebuilt, not repaired. A rebuild replays history into a fresh table and swaps. At 50K events/s of replay throughput, 2 billion events take ~11 hours. That number, not the steady-state write rate, sizes the read side.
- Events are forever; schema evolution is the real job. Additive fields are free. Renames and meaning changes need an upcaster on read or a copy-and-transform of the store. Budget one versioning decision per event type per year.
- Most domains should not be event-sourced. If no one will ever ask "what was the state at 14:02 last Tuesday, and why," a table plus an audit log or CDC stream delivers 90% of the value at ~20% of the lifetime cost.
The Problem#
A brokerage's order service stores orders in a PostgreSQL table: one row per order, status updated in place from NEW to PARTIALLY_FILLED to FILLED or CANCELLED. Three requirements arrive in the same quarter. Compliance wants to reconstruct the exact state of every order at any instant for seven years, including which version of the risk rules accepted it. The client app wants a portfolio screen that joins orders, fills, positions and market prices in under 50ms at 20K reads/s, while the order table takes 4K writes/s and every new index slows matching. And a bug in the fee calculator has been overcharging partial fills for six weeks; finance wants to know exactly which fills, by how much, and wants the corrected balances shown without rewriting history. Updating rows in place destroyed the facts needed for all three. An audit table bolted on with triggers captures columns, not intent — it records that fee changed from 1.20 to 1.35, not that a FillFeeAssessed happened under fee schedule v14. The job is to decide which parts of this domain need their history as the source of truth, keep the decision-making write model small and strongly consistent, serve reads from models shaped for each screen, and make fixing a projection a replay rather than a migration.
Case Studies That Use This Pattern#
- Ledger & Wallet — A double-entry ledger is an event log by construction; balances are projections over immutable postings
- Payments — Payment state machines where every transition must be explainable to a regulator and a disputing customer
- Stock Exchange — A deterministic matching engine that recovers by replaying its input journal from the last snapshot
- Collaborative Docs — Operation logs as the source of truth, with the document as a periodically compacted snapshot
- Hotel Booking — Reservation lifecycles with separate read models for search availability and for the booking itself
- Ad Click Aggregation — Raw click events as the immutable record; aggregates rebuilt by replay when the counting logic changes
- Search Engine — The canonical CQRS read model: an index built from a write store and rebuilt when the mapping changes
- Idempotency — Projections and command handlers both see events more than once; dedupe is part of the design, not an afterthought
Which Problem Are We Solving?#
"Use CQRS/ES" hides four different goals. Name them and commit before drawing a box.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Read/write shape mismatch (search, dashboards, feeds over a normalized store) | Reads 10–100× writes, need denormalized shapes the write schema cannot serve | CQRS only: one write store, read models fed by CDC or an outbox | Read model lag surprises users; rebuilds race live updates | Read models converge within an agreed lag (e.g. p99 < 2s); drift measured |
| History is the product (ledgers, trading, licences, medical records) | Regulators or customers must see any past state and why it changed | Full event sourcing with streams per aggregate, snapshots, projections | Unreadable old events; projection bugs; long rebuilds | Any past state reproducible exactly; events never mutated |
| Audit trail on a CRUD system | "Who changed what, when" for a subset of entities | Table as source of truth plus an append-only audit log or table-level CDC | Audit captures columns, not intent | Every change attributable; no claim to rebuild state from the log |
| Temporal analytics and replays (recompute metrics with new logic) | Logic changes after the fact; raw facts are cheap to keep | Immutable raw event log in a stream or lake; batch and stream recompute | Recompute cost and time; late events | Aggregates reproducible from raw events within a bounded window |
🎯 Staff Move: "I'll treat the order lifecycle as intent two: compliance needs any past state and the reason for it, so orders are event-sourced with one stream per order. Everything else — the portfolio screen, search, the risk dashboard — is intent one: read models projected from those events. I'm not event-sourcing customer profiles; they get a normal table and an audit log."
The Core Tradeoff#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Single model, CRUD | Simplest; read-your-writes for free; one schema to migrate | Read shapes fight write indexes; history lost on update | Product teams waiting on index changes; finance during a "what happened" investigation |
| CRUD + audit log | Cheap history of changes; no new read path | Captures columns, not intent; cannot rebuild state reliably | Auditors interpreting column diffs; nobody until a dispute |
| CQRS, state-stored write model | Read models shaped per screen; write store stays lean; reversible | Eventual consistency between write and read; projection lag; rebuild tooling | Users who refresh and see stale data; the team owning the projectors |
| Event sourcing + CQRS | Complete history; temporal queries; projections rebuilt with new logic; natural fit for event-driven integration | Every read needs a projection; event schemas are forever; rebuilds take hours; steep learning curve | The owning team for the system's lifetime; new hires; whoever runs the 11-hour replay |
| Event sourcing without CQRS (query by loading aggregates) | Fewer moving parts at first | Any list or search query loads thousands of streams; collapses past a demo | Every reader; the event store's IOPS |
| Log as database (Kafka topic as the store, compacted) | Very high write throughput; replay built in | No optimistic concurrency per key; no transactional "append if version = n"; retention and compaction surprises | The team discovering two writers both appended version 8 |
Staff Default Position#
Split reads from writes whenever the shapes diverge, but make the event log the source of truth only for the aggregates where history, explanation or replay is a stated requirement — and then own the projection, snapshot and versioning machinery as a product.
The default stack for intent two: one stream per aggregate in an event store that supports conditional append on expected version (a purpose-built event store, or PostgreSQL with a (stream_id, version) unique key); commands load the aggregate (snapshot plus tail), decide, and append one or more events atomically; projections are idempotent consumers that track their own checkpoint position and upsert by (aggregate_id, version); snapshots are written asynchronously every ~200 events and can be deleted at any time; every event type carries a schema_version and old versions are upcast on read. Read models are disposable: each can be rebuilt into a shadow table and swapped behind a pointer. For intent one, skip the event store entirely — keep the normalized write model and feed read models through a transactional outbox or CDC.
When to Deviate#
- The domain is CRUD with an audit requirement. User profiles, settings, product catalogues. A table plus an append-only audit log (or table-level CDC into a retained topic) answers "who changed what." Say plainly that you cannot rebuild state from it, and that nobody asked you to.
- Throughput per aggregate is extreme. A single stream accepting 5K+ appends/s (a hot order book, a global counter) serializes on its version check. Shard the aggregate, or use a deterministic single-writer process with an input journal — the matching-engine model — rather than per-command optimistic concurrency.
- Personal data that must be erasable. "Events are immutable" collides with deletion rights. Keep PII out of events (store a reference), or encrypt per subject and delete the key ("crypto-shredding"). If neither is acceptable to legal, do not event-source that data.
- The team has never run it and the deadline is a quarter. Event sourcing's costs land in months 6–24. Ship CQRS on a state-stored model with an outbox; the outbox events can become the seed of an event-sourced model later if history proves valuable.
One Question, Three Levels#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Use CQRS and event sourcing for scale: commands to Kafka, a consumer builds the database" | "Which aggregates need history as the source of truth? Those get streams; everything else gets CQRS on a normal table." | "Which teams will own an event-sourced system for ten years, and do we have the platform — store, projector framework, schema registry — to make that survivable?" |
| Consistency | "It's eventually consistent" | Strong inside one aggregate via expected-version append; eventual across read models with a stated lag; read-your-writes by returning the new version and waiting on it | Defines which product surfaces may be eventually consistent and gets product to sign off per surface |
| Read models | One read database updated by a consumer | Each projection owns its checkpoint, is idempotent by version, and can be rebuilt into a shadow and swapped | Projector framework as a paved road; rebuild time is an SLO; projections classified by tier |
| Schema | "We'll version the events" | schema_version on every event; additive by default; upcasters on read; copy-and-transform for breaking changes; never mutate stored events | Event catalogue with owners; deprecation windows; versioning bankruptcy plan for streams that rot |
| Failure | "Replay the events if something breaks" | "A full replay is 2B events at 50K/s — 11 hours. So rebuilds go to a shadow table, run in parallel by partition, and the old projection serves until the swap." | Prices rebuild time against the incident cost; game-days a full rebuild each quarter |
| When not to | Applies it everywhere for consistency of style | Names the aggregates that should stay CRUD and why | Sets an org rule: event sourcing needs a design review and a named long-term owner |
Why "First move" separates levels
The L5 answer is a recognizable architecture and it gets downleveled because it is applied to the whole system without a reason. Putting commands on Kafka and building the database from the topic also loses the one thing event sourcing must keep: an atomic "append only if the stream is still at version n" check. The Staff candidate scopes the pattern to the aggregates that need it and protects the consistency boundary. The Principal candidate asks whether the organization can afford to own it, because the cost of event sourcing is mostly people-years, not servers.
Why "Consistency" separates levels
"Eventually consistent" is true and useless. The interviewer wants to know where strong consistency lives (inside one aggregate, enforced by the version check on append), where it does not (between the event store and every read model), how long the gap is (typically 100ms–2s, worse during rebuilds), and what the user sees in that gap. The Staff answer adds a mechanism: the command returns the new stream version, and the client or API waits until the projection's checkpoint passes it, with a timeout, before reading. That turns a vague property into a design.
Why "Failure" separates levels
"Just replay" is the promise that sells event sourcing, and it is a batch job whose duration grows linearly with history. Seniors say it; Staff engineers have computed it. At 2 billion events and 50K events/s per projector, a rebuild takes ~11 hours — longer than most incidents are allowed to last. The Staff design therefore never rebuilds in place: it builds a shadow, parallelizes by partition, and swaps atomically, while the old (possibly wrong) projection keeps serving or is explicitly marked degraded.
Where the Design Splits#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Event-Sourced vs State-Stored Write Model | Complete, replayable history vs a simpler store with read-your-writes and normal migrations |
| 2 | Synchronous vs Asynchronous Projections | Read-your-writes and one transaction vs independent scaling and failure isolation with lag |
| 3 | Small Aggregates vs Large Aggregates | Low contention and fast loads vs invariants that span more data in one consistency boundary |
| 4 | Upcast on Read vs Copy-and-Transform | Old events stay untouched and every reader pays a translation vs one migration that rewrites the store |
| 5 | Purpose-Built Event Store vs Database or Log You Already Run | Native streams, subscriptions and projections vs operational familiarity and fewer systems |
Fault Line 1: Event-Sourced vs State-Stored Write Model#
An event-sourced write model stores OrderPlaced, OrderPartiallyFilled, OrderCancelled and derives state by folding. A state-stored write model keeps an orders row and emits events through an outbox as a side effect. Both can feed the same read models. Who pays: event sourcing makes the owning team pay forever in versioning, snapshots and projection ownership, and makes every new engineer learn to think in folds; state-stored makes auditors and investigators pay when the history they need was overwritten, and makes "recompute everything with corrected logic" a data migration rather than a replay. Staff default: state-stored with an outbox for most aggregates; event-sourced for the handful where the history is legally or commercially the asset — ledgers, orders in regulated markets, entitlements and licences. Deviate when: the domain is inherently a log already (operations in a collaborative editor, postings in a ledger) — then a state-stored model is the unnatural one.
Fault Line 2: Synchronous vs Asynchronous Projections#
A synchronous projection updates the read table in the same transaction as the event append — possible when both live in one PostgreSQL database. An asynchronous projection consumes the event stream and updates the read model on its own schedule. Who pays: synchronous projections make the write path pay — every append also maintains N read tables, write latency grows by ~1–5ms per projection, and a broken projection blocks commands; asynchronous projections make the user pay a lag window and make the team pay for checkpoints, idempotency and lag alerting. Staff default: asynchronous for everything except at most one "own view" projection the command's caller reads immediately; solve read-your-writes by returning the new version and waiting on the projection checkpoint for up to ~1s. Deviate when: the read model is tiny, co-located and strictly needed for the next command's validation — then it is part of the write model, not a projection.
Fault Line 3: Small Aggregates vs Large Aggregates#
The aggregate is the unit that loads, decides and appends atomically. Making "customer account with all orders" one aggregate lets you enforce "total open exposure ≤ $50K" in one place; it also means every order for that customer contends on one stream version and loads a stream that grows forever. Who pays: large aggregates make throughput and latency pay — conflicts climb past ~5% of commands once a stream sees more than a few writes per second, and load time grows with the stream; small aggregates make cross-aggregate invariants pay, which now need a saga or process manager, a reservation, or acceptance of brief violation with compensation. Staff default: the smallest aggregate that owns one invariant completely — one order, one account's balance — with cross-aggregate rules handled by reservations or a process manager. Deviate when: the invariant is hard and the conflict rate is genuinely low (a few writes per minute per aggregate); then a larger boundary is simpler than a saga.
Fault Line 4: Upcast on Read vs Copy-and-Transform#
When OrderPlaced v1 stored price as a float and v2 stores price_minor as an integer with a currency, you can translate v1 to v2 every time it is read (an upcaster in the deserialization path), or run a migration that copies every stream into a new store with events rewritten to v2. Who pays: upcasting makes every reader pay a small CPU cost and makes the codebase carry a growing chain of translators — after five versions a v1 event passes through four upcasters; copy-and-transform makes the operations team pay a large one-time migration, a cutover with a write freeze or dual-write window, and makes auditors ask why "immutable" events were rewritten. Staff default: additive changes need nothing; non-additive changes get an upcaster; copy-and-transform only for "versioning bankruptcy" — when the chain is long enough or stream boundaries are wrong enough that a one-time migration is cheaper than carrying them. Deviate when: a meaning change (not a shape change) makes old events ambiguous — then add a new event type rather than reinterpreting the old one.
Fault Line 5: Purpose-Built Event Store vs Database or Log You Already Run#
A purpose-built event store offers streams, expected-version appends, category subscriptions and server-side projections out of the box. PostgreSQL gives you the same core with a table — (stream_id, version) unique, an append-only insert, a global sequence for subscriptions — and an operations team that already knows it. A Kafka topic gives replay and throughput but no per-key conditional append. Who pays: a new store makes the platform pay a new tier-0 system and on-call skill; PostgreSQL makes the team pay for building subscriptions, checkpointing and global ordering themselves, and caps throughput at one primary (~5–20K appends/s typical); Kafka-as-store makes correctness pay, because two writers can both append "version 8." Staff default: PostgreSQL (or DynamoDB with a conditional put on (stream_id, version)) for the event store, Kafka downstream as the distribution log for projections and other teams. Deviate when: several teams will run event-sourced services and a dedicated store's projections and subscriptions replace code each team would otherwise write.
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "We'll use CQRS and event sourcing for scalability" | "Pattern-matched from a blog; no reason given" | "Reads and writes have different shapes here, so CQRS. Only the order lifecycle needs history as the source of truth, so only orders are event-sourced." |
| "Commands go to Kafka and a consumer writes the database" | "Lost the optimistic concurrency check" | "Appends are conditional on the stream's expected version, in a store that can enforce it. Kafka distributes events after they're committed." |
| "If a projection is wrong, we just replay" | "Hasn't timed a replay" | "A full replay is about 11 hours at our volume, so rebuilds go to a shadow table in parallel by partition, and we swap when it catches up." |
| "Events are immutable, so schema changes are easy" | "Has never read a three-year-old event" | "Immutable means old shapes live forever. Additive changes are free; anything else gets an upcaster or a new event type." |
| "Snapshots store the state so we don't need old events" | "Snapshots just became the source of truth" | "Snapshots are a cache. I can delete all of them and rebuild from events; if I can't, they're an unversioned second database." |
| "The UI will show the update after the event is processed" | "Hasn't thought about the user who refreshes" | "The command returns the new version. The API waits up to a second for the projection checkpoint to pass it, then reads." |
Quick Reference#
Staff Sentence Templates#
"The consistency boundary is the [aggregate]. A command loads its stream, decides, and appends only if the stream is still at version [n]. Across [read models], consistency is eventual with a p99 lag of [X] ms, and product has signed off on that for [screens]."
"I'm event-sourcing [aggregate] because [regulator / dispute / replay need]. I'm not event-sourcing [other entities]; they get a table and an audit log, because nobody will ask for their state at an arbitrary past instant."
"A full rebuild of [projection] is [N] events at [rate] per second — about [hours]. So rebuilds run into a shadow table, partitioned [P] ways, and swap when the shadow's checkpoint reaches the live head."
"Every event carries a schema version. [Change] is additive, so no action. [Breaking change] gets an upcaster from v[n] to v[n+1]; we never rewrite stored events unless we declare versioning bankruptcy on a stream type."
Implementation Deep Dive#
1. The Event Store on PostgreSQL — Expected-Version Append#
The whole write-side correctness argument is one unique constraint. Two commands that both loaded order 991 at version 7 race to append version 8; the database lets exactly one win.
CREATE TABLE events (
global_seq bigserial PRIMARY KEY, -- subscription order for projectors
stream_id text NOT NULL, -- 'order-991'
stream_version int NOT NULL, -- 1, 2, 3 ... per stream
event_id uuid NOT NULL UNIQUE, -- dedupe for consumers and retried commands
event_type text NOT NULL, -- 'OrderPlaced'
schema_version int NOT NULL, -- payload contract version
payload jsonb NOT NULL,
metadata jsonb NOT NULL, -- causation_id, correlation_id, actor, rule_version
recorded_at timestamptz NOT NULL DEFAULT now(),
UNIQUE (stream_id, stream_version) -- the optimistic concurrency check
);
-- Command handler: load, decide, append
-- 1. load: SELECT ... FROM snapshots WHERE stream_id = $id; -- optional
-- SELECT ... FROM events WHERE stream_id = $id
-- AND stream_version > $snapshot_ver ORDER BY stream_version;
-- 2. decide: new_events = order.handle(command) -- pure function, no I/O
-- 3. append, all-or-nothing:
BEGIN;
INSERT INTO events (stream_id, stream_version, event_id, event_type,
schema_version, payload, metadata)
VALUES ('order-991', 8, $e1, 'OrderPartiallyFilled', 2, $p1, $m1),
('order-991', 9, $e2, 'FillFeeAssessed', 1, $p2, $m2);
COMMIT;
-- unique violation on (stream_id, stream_version) => another command won:
-- reload and re-decide (retry up to 3 times), or return 409 to the caller
The global_seq trap: projectors subscribe by reading WHERE global_seq > $checkpoint. But bigserial values are assigned at insert time, not commit time — transaction A takes seq 100, transaction B takes 101 and commits first, a projector reads 101, advances its checkpoint past 100, and A's event is skipped forever. Fixes: read only up to the oldest in-flight transaction's position (track it with pg_snapshot_xmin(pg_current_snapshot()) and a transaction-id column), serialize appends through a single sequence-assigning writer, or consume through CDC from the WAL, which is in commit order. This is the bug that turns "events are never lost" into "events are occasionally invisible."
🎯 Staff Insight: Put
rule_version(or fee schedule version, model version) in event metadata. When finance asks which fills were charged under the buggy fee schedule, the answer is a query, not an archaeology project.
2. Projections — Checkpoint, Idempotent Upsert, Shadow Rebuild#
A projector is a consumer with three duties: apply each event at most once in effect, remember where it is, and be rebuildable from zero without taking the live read model down.
projector "order_status_v3":
checkpoint = SELECT position FROM projector_checkpoints WHERE name = 'order_status_v3'
loop:
batch = read_events(after = checkpoint, limit = 1000) # in commit order
BEGIN # read model's database
for e in batch:
apply(e) # e.g. INSERT ... ON CONFLICT (order_id) DO UPDATE
# SET status = $s, version = $v
# WHERE order_status_v3.version < $v -- stale or duplicate = no-op
UPDATE projector_checkpoints SET position = last(batch).global_seq
WHERE name = 'order_status_v3'
COMMIT # effect and checkpoint together
rebuild (new logic, or a bug fix):
1. create table order_status_v4; projector 'order_status_v4' starts at position 0
2. run 16 workers, each owning hash(stream_id) % 16, until within 5s of head
3. compare v3 vs v4 on a sample (counts, checksums by status)
4. flip the read pointer (view, alias or config) from v3 to v4; keep v3 for 24h
5. drop v3
Why the checkpoint lives with the read model: if the checkpoint sits in Kafka offsets or a separate store while the effect lives in PostgreSQL, a crash between the two replays a batch (harmless with version-guarded upserts) or skips one (data loss). Same database, same transaction, no gap.
Rebuild throughput is the sizing number. A single-threaded projector doing batched upserts typically manages 5–20K events/s; 16 partitioned workers against a read store that can absorb the writes gets 50–150K/s. For 2B events: ~3.7 hours at 150K/s, ~11 hours at 50K/s, ~4.6 days at 5K/s. Run the back-of-envelope before promising "we'll just replay."
3. Snapshots — A Cache With a Version Number#
Snapshots bound load time for long streams. They are written asynchronously, keyed by the code version that produced them, and are always safe to delete.
snapshot policy:
write when (stream_version - last_snapshot_version) >= 200
or load_time_ms > 10
written by a background worker, never in the command path
snapshots table:
stream_id, stream_version, snapshot_schema, state_blob, created_at
PRIMARY KEY (stream_id) -- keep only the latest
load(stream_id):
s = SELECT * FROM snapshots WHERE stream_id = $id
if s and s.snapshot_schema == CURRENT_SNAPSHOT_SCHEMA:
state = deserialize(s.state_blob); from = s.stream_version
else:
state = initial(); from = 0 # schema changed: ignore old snapshot
for e in events(stream_id, after = from): state = apply(state, upcast(e))
return state
Why snapshot_schema matters: when the aggregate's state shape changes, old snapshots are wrong in a way the deserializer may not notice. Versioning them means a code change silently falls back to a full fold — slower for a day while the snapshotter catches up, but correct. The rule that keeps snapshots honest: the system must pass its tests with the snapshots table truncated.
🎯 Staff Move: "If a stream needs a snapshot to load in under 10 milliseconds, I first ask whether the aggregate is too big. An order with 40 events doesn't need snapshots. An account with five years of daily interest postings does — or it needs a period-close event that starts a new stream."
4. Schema Evolution — Upcasters, New Types, and Bankruptcy#
Every event ever written must remain readable by current code. The tools, in order of preference:
# 1. Additive change: new optional field, readers default it. No version bump needed
# if the format tolerates unknown fields (JSON with a tolerant reader, Avro with defaults).
# 2. Shape change: bump schema_version, register an upcaster
upcasters = {
("OrderPlaced", 1): lambda p: {**p, "price_minor": round(p["price"] * 100),
"currency": p.get("currency", "USD"),
"price": None} | {"_v": 2},
("OrderPlaced", 2): lambda p: {**p, "time_in_force": p.get("tif", "DAY")} | {"_v": 3},
}
def upcast(e):
while (e.type, e.schema_version) in upcasters:
e.payload = upcasters[(e.type, e.schema_version)](e.payload); e.schema_version += 1
return e
# 3. Meaning change: do NOT reinterpret the old event. Add a new type.
# 'OrderCancelled' meant "user cancelled"; risk-initiated cancels now need different handling
# -> introduce 'OrderCancelledByRisk'; old OrderCancelled events keep their original meaning.
# 4. Correction of a wrong fact: append a compensating event, never edit.
# FillFeeAssessed(fee=1.35) was wrong -> FillFeeCorrected(fill_id, delta=-0.15, reason, ticket)
# 5. Versioning bankruptcy (rare): copy-and-transform the store into a new one,
# rewriting old versions to current, with a dual-read window and an auditor-visible record
# of the transformation. Budget: a quarter, a freeze window, sign-off from compliance.
The compensating-event rule is what auditors care about. Accountants do not erase ledger lines; they post a correcting entry. The six-week fee bug is fixed by appending FillFeeCorrected events for each affected fill — found by querying metadata for fee_schedule = v14 — and every projection that shows a balance picks them up. The history now shows both the error and the fix, which is exactly what a regulator wants to see.
Technique Comparison
| Technique | Source of Truth | Read Consistency | History | Rebuild Cost | Ops Burden | Best For |
|---|---|---|---|---|---|---|
| CRUD | Table | Strong | None | N/A | Low | Most entities |
| CRUD + audit log | Table | Strong | Column diffs | N/A | Low | Settings, profiles, catalogues |
| CQRS, state-stored + outbox | Table | Eventual for read models (~0.1–2s) | Events since adoption | Re-snapshot from table | Medium | Read/write shape mismatch |
| Event sourcing + async projections | Event log | Strong per aggregate; eventual reads | Complete | Hours (replay) | High | Ledgers, orders, licences |
| ES + synchronous own-view projection | Event log | Read-your-writes for one view | Complete | Hours | High | One latency-critical view |
| Single-writer journal + snapshots | Input journal | Strong (single thread) | Complete inputs | Minutes from nightly snapshot | High, specialized | Matching engines, exchanges |
Architecture Diagram#
How to narrate it: the only thing a command waits on is the conditional append into the event store. Snapshots are drawn dashed because deleting them changes latency, not correctness. Events leave the store through CDC so subscribers see commit order, not sequence-assignment order, and land in Kafka keyed by stream so each aggregate's events stay ordered on one partition. Every projector owns its own checkpoint and can be rebuilt into a shadow without touching the others. The order team owns the event schemas; each read model's team owns its projector, its lag alert and its rebuild runbook.
Failure Scenarios#
1. The Skipped Event — Sequence Gaps Hide 0.3% of Fills#
A projector reads WHERE global_seq > $checkpoint from a bigserial column under 4K appends/s.
Day 0 Projector polls every 100ms. Concurrent transactions commit out of seq order.
Day 0-21 Transaction holding seq N commits 30ms after N+1. Poll in between sees N+1,
advances checkpoint past N. Event N never projected.
Day 21 Client reports portfolio shows 300 shares; statement shows 400.
Day 21 Reconciliation: 0.3% of fills missing from the portfolio projection, ~180K events.
Day 22 Rebuild from 0 into a shadow: 9 hours. Swap. Root cause found the next morning.
Detection: projection.reconciliation_mismatch — per-aggregate version in the projection vs the event store's max version, sampled hourly; projector.gap_detected when a batch's sequence numbers are non-contiguous.
Blast radius: every projection reading by sequence; customer-visible balances wrong for three weeks.
Mitigation: shadow rebuild; customer communication for affected accounts.
Prevention: consume in commit order (CDC from the WAL) or hold reads below the oldest in-flight transaction; hourly version reconciliation with a page on any aggregate behind by more than 5 minutes.
Owner: event store platform owns subscription semantics; each projection team owns its reconciliation alert.
🎯 Staff Insight: "Append-only" guarantees nothing is lost from the store. It says nothing about whether every reader saw every event. Those are separate properties, and the second one needs its own test.
2. The 30-Hour Replay — A Mapping Change Meets Five Years of History#
A search projection's analyzer changes, requiring a full rebuild. The team rebuilds in place.
t=0 Search index dropped and recreated; projector reset to position 0.
t=0 History: 5.4B events. Single projector, 50K events/s => ~30 hours.
t=+10min Search returns empty for most customers. Support queue triples.
t=+2h Team parallelizes to 8 workers: 300K/s target, Elasticsearch write-rejects at 120K/s.
t=+14h Rebuild complete at ~110K/s effective. 14 hours of degraded search.
Detection: search.result_count_p50 dropping to near zero; projector.lag_events at billions.
Blast radius: every search user for 14 hours; the order write path unaffected (this is the CQRS dividend).
Mitigation: none fast once the old index was dropped; serve a "results incomplete" banner.
Prevention: always rebuild into a shadow index behind an alias; size rebuild time in the design (events ÷ sustainable sink write rate); keep a compacted "latest state per aggregate" topic so rebuilds can start from current state when history is not needed.
Owner: search team owns the rebuild runbook and the shadow-and-swap tooling.
3. The Upcaster That Changed History — Fees Recomputed Wrong#
An upcaster for FillFeeAssessed v1 → v2 defaults a missing fee_currency field to the account's current currency instead of the currency at the time.
Week 0 Upcaster ships. 2% of accounts changed base currency in the past two years.
Week 0 Projections rebuilt for an unrelated reason replay v1 fee events through the upcaster.
Week 0 Historic fees for those accounts now show in the wrong currency, ~40K accounts.
Week 3 Monthly statements disagree with last month's PDFs. Compliance escalates.
Detection: statement.diff_vs_prior_period on closed periods (should be exactly zero); golden-file tests replaying a frozen sample of production events through current upcasters.
Blast radius: every projection rebuilt after the upcaster shipped; closed accounting periods changed, which is a regulatory problem, not just a bug.
Mitigation: fix the upcaster to derive currency from the account's event history at that point; rebuild affected projections; reissue statements.
Prevention: upcasters must be pure functions of the event (and at most its own stream's earlier events) — never of current state; a frozen corpus of real historic events with expected outputs runs in CI for every upcaster change; closed-period checksums alarm on any change.
Owner: the event schema owner (order team) owns upcasters and the golden corpus; finance owns closed-period checks.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Projection lag | projector.lag_ms p99 > 2s; checkpoint not advancing for 60s | Stale reads on one read model | Scale projector; check sink write rejects | Read-model team |
| Skipped events | Version reconciliation mismatch; sequence gap counter | Silent wrong state | Commit-order subscription; shadow rebuild | Event store platform |
| Concurrency conflicts spike | command.conflict_rate > 5% | Higher latency and 409s for a hot aggregate | Smaller aggregate; single-writer for hot streams | Domain team |
| Stream too long | aggregate.load_ms p99 > 10ms; stream length > 10K | Command latency | Snapshot; period-close to start new stream | Domain team |
| Bad upcaster | Golden-corpus test failure; closed-period checksum change | Every rebuilt projection | Fix and rebuild; reissue outputs | Schema owner |
| Rebuild too slow | Rebuild ETA > rebuild SLO | Read model degraded for hours | Shadow-and-swap; partitioned workers; compacted state topic | Read-model team |
| PII erasure request | Erasure queue age > SLA | Legal exposure | Crypto-shred subject key; PII by reference only | Domain team + privacy |
| Event store write outage | append.error_rate; primary health | All commands fail; reads keep serving stale | Failover; commands return 503 | Event store platform |
Beyond Staff: The Principal View#
Why L7 Sees This Problem Differently#
A Staff engineer decides whether one service should be event-sourced and builds it well. A Principal engineer notices that event sourcing is the most expensive architectural decision a team can make that looks cheap on day one. The first quarter feels productive: commands, events, a couple of projections. The bill arrives in year two — a projection rebuild that takes a weekend, an upcaster chain six deep, a new hire who needs a month to stop writing CRUD into a fold, an erasure request that collides with "immutable." Multiply that by every team that read the same blog post and the org has five homegrown event stores with five different subscription bugs. At L7 the work is deciding where the org allows event sourcing, giving those teams a paved road (store, projector framework, schema registry, rebuild tooling), and making everyone else use CQRS on ordinary tables with an outbox — because that captures most of the benefit at a fraction of the lifetime cost.
🧭 Principal Move: "I'd make event sourcing a reviewed decision with a named ten-year owner, not a team-level style choice. Teams that clear the bar get a shared event store and projector framework. Everyone else gets CQRS on their existing database with the outbox library — and most of them will never notice the difference."
The Org-Level Fault Line#
A shared event-sourcing platform vs letting each team build its own vs discouraging the pattern entirely.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each team builds its own event store | Autonomy; fits each domain | Same subscription, snapshot and versioning bugs reinvented N times; no shared rebuild tooling | Every team's on-call; the next reorg that inherits five designs |
| Ban event sourcing; CRUD and outbox only | Simple, consistent, cheap | Domains that genuinely need history (ledger, trading) rebuild it badly with audit tables | Finance and compliance during every investigation |
| Gated platform: review to adopt, shared store and projector framework | Pattern used where it pays; one hardened implementation; rebuilds are a product | Platform must be funded and support a few store backends; review must not become a bottleneck | Platform headcount (2–4 engineers); architecture review time |
The Principal default is the third row: a short adoption review (does history matter, who owns it for ten years, is PII handled), then a shared platform that makes the hard parts — commit-order subscriptions, shadow rebuilds, upcaster test corpora — someone else's solved problem.
Cost Model#
Assumptions: fully loaded engineer ~$25K/month; managed PostgreSQL ~$2–8K/month per primary with replicas; event payload ~800 bytes; 3-year retention online, older in object storage at ~$0.02/GB-month; Kafka distribution included at ~$0.30/GB-month retained.
| Scale | Event Volume | Machinery | Infra Cost | People / On-call | Rough Monthly Total |
|---|---|---|---|---|---|
| One domain | ~50 events/s average, 500 peak; ~4.7B events over 3 years | PostgreSQL event store, 4 projections, hand-built projector | DB ~$3K; ~3.8 TB online; archive negligible | 1 FTE of the domain team's time on ES machinery | ~$3K + ~$25K people |
| Several domains | 5K events/s, 4 teams | Shared store per domain on the platform, projector framework, schema registry, rebuild tooling | DBs ~$20K; Kafka ~$4K; rebuild compute bursts ~$2K | 2 FTE platform; projector on-call per team | ~$26K + ~$50K people |
| Org-wide | 50K events/s, 20 teams | Sharded event stores, commit-order CDC, shadow-rebuild service, golden-corpus CI, erasure tooling | DBs ~$120K; Kafka ~$30K; archive ~$5K | 4-person platform; tier-1 rotation | ~$155K + ~$100K people |
The Principal observation: at every scale the people line dominates. The real cost of event sourcing is not storage — 4.7B events of 800 bytes is under 4 TB — it is the engineering time spent on projections, versioning and rebuilds, forever. That is why the gate matters more than the platform: the cheapest event-sourced system is the one you correctly decided not to build.
The 3-Year Evolution Path#
The Year 2 trigger is the one teams feel first: the first full rebuild that takes longer than anyone promised. That is when shadow-and-swap, partitioned replay and compacted state topics stop being nice-to-haves.
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Event-sourcing an aggregate | One-way | Leaving means freezing history into a table and losing replay; every consumer of the event stream is affected |
| Adding a read model (CQRS) | Two-way | Drop the projection; the write model is untouched |
| Aggregate and stream boundaries | One-way-ish | Changing them later is a copy-and-replace migration of the store |
| Event names and field meanings | One-way | Old events keep their meaning forever; only new types can change it |
| Putting PII inside events | One-way | Erasure requires crypto-shredding or rewriting history |
| Event store technology | One-way-ish | Migration is a full copy with a dual-read window |
| Snapshot frequency and format | Two-way | Delete all snapshots; they regenerate |
| Synchronous vs async projection | Two-way | Move the projector; the events don't change |
The Standard I'd Write#
RFC: Event Sourcing and CQRS (v1)
Scope: Any service that stores events as its source of truth, or maintains read models derived from another store.
MUST:
- Event sourcing MUST be approved in a design review that names the history requirement, the ten-year owning team, and the PII strategy. CQRS on a state-stored model with the outbox library needs no review.
- Appends MUST be conditional on the stream's expected version in a store that enforces it atomically.
- Every event MUST carry
event_id,stream_id,stream_version,schema_versionand causation/correlation metadata. Stored events MUST NOT be updated or deleted except via an approved copy-and-transform or crypto-shredding.- Projections MUST be idempotent by version, keep their checkpoint in the same transaction as their effect, and be rebuildable into a shadow without downtime. Each projection MUST publish its full-rebuild time.
- Upcasters MUST be pure functions of the event and MUST pass a golden corpus of real historic events in CI.
SHOULD: Keep PII out of event payloads (reference by ID); snapshot only when load time exceeds 10ms; reconcile projections against stream versions hourly.
Exceptions: Filed with the architecture group, time-boxed, with a named owner.
Success metrics: zero closed-period changes from replays; projection lag p99 < 2s; every tier-1 projection rebuildable within its published SLO; erasure requests completed within the legal window.
What I'd Tell the VP#
Some of our systems throw away the history of how a number came to be, which is why answering a regulator or a disputing customer takes weeks. For the few systems where that history is the product — the ledger and order handling — I'm proposing we keep every change as a permanent record and build screens from it. It is a long-term commitment: those teams will spend roughly a fifth of their time maintaining this for as long as the system lives, so I'm also proposing we limit it to those systems and provide shared tooling so they don't each build it from scratch. Everyone else keeps their current databases with a lighter version of the same idea. The payoff is that investigations become queries, and fixing a calculation bug becomes a replay instead of a data migration.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Gates the one-way door | "Event sourcing needs a ten-year owner and a reason. CQRS on a table doesn't." |
| Prices people, not storage | "Four terabytes of events is cheap. A team spending 20% of its time on projections is not." |
| Rebuild as an SLO | "Every projection publishes its full-rebuild time, and we game-day it quarterly." |
| Protects history legally | "PII by reference or crypto-shredded. Closed periods are checksummed so a replay can't rewrite them." |
| Paved road for the few | "Four teams need this. One store, one projector framework, one golden corpus." |
Staff answers that L7 interviewers find insufficient:
- "We'll event-source the order service." — Correct locally; silent on who else will copy the pattern and whether the org can support it.
- "Projections can be rebuilt." — True; no rebuild SLO, no game day, no shared tooling.
- "We'll handle GDPR later." — Immutable events plus personal data is a one-way door that must be decided before the first event is written.
How Real Companies Built It#
LMAX: An In-Memory, Event-Sourced Business Logic Processor#
LMAX, a retail financial trading platform, runs all of its business logic — every trade, for every customer, in every market — on a single thread in a Java process that the architecture write-up says handles 6 million orders per second on commodity hardware. The Business Logic Processor holds its state entirely in memory and is event-sourced: its current state is derivable by processing the input events, which are journaled durably before processing. Snapshots are taken nightly during quiet periods, and a full restart — JVM restart, loading a recent snapshot and replaying a day of journaled events — takes under a minute. Several processors (two in the primary data center, one in disaster recovery) consume the same input events with all but one processor's output ignored, so failover takes microseconds, and because the processor is deterministic, engineers can copy the event sequence to a development environment and replay a production problem exactly (The LMAX Architecture).
Staff insight: This is event sourcing at its sharpest: the input log is the truth, the in-memory state is a fold, snapshots bound recovery time, and determinism buys exact replay. It is also the counterexample to per-command optimistic concurrency — when one aggregate is that hot, a single writer with a journal beats conflict retries.
Netflix: Event Sourcing for the Downloads Licensing Service#
When Netflix launched downloads in 2016, its licensing service had to change from stateless to real-time and stateful, because downloaded content must be tracked to enforce business rules such as how long a licence lasts and how many devices a download can live on. The team describes building this on a Cassandra-backed event sourcing architecture, using data versioning and snapshotting to provide flexibility and scale — loading the latest snapshot and replaying only later events — and choosing the approach because the data model was expected to keep changing as product requirements iterated (Scaling Event Sourcing for Netflix Downloads, QCon New York 2017, InfoQ write-up of the talk).
Staff insight: The stated reason is the honest one: requirements were moving fast and the team wanted to derive new views from the same facts rather than migrate a schema each time. Versioning and snapshots are named up front, not discovered later — exactly the two costs this page says to plan for on day one.
EventStoreDB: Server-Side Projections and Their Write Amplification#
EventStoreDB (now Kurrent), a purpose-built event store, ships five system projections — including $by_category, which links events from streams like account-1 into a $ce-account category stream, and $by_event_type — and supports custom projections that react to events by appending or linking to other streams. Its documentation is explicit about the cost: projections run only on the cluster leader, add CPU and I/O load there, are disabled by default, and cause write amplification — enabling the category, event-type and correlation system projections together roughly quadruples writes, and custom projections can amplify further (EventStoreDB projections).
Staff insight: "Projections are free" is the myth; every projection is extra writes somewhere. Whether it runs inside the store or as an external consumer, a projection is a workload with its own throughput ceiling — the same ceiling that decides how long a rebuild takes.
Greg Young: Versioning as the Long-Term Cost#
Greg Young, a long-time practitioner and teacher of CQRS and event sourcing, wrote a book devoted entirely to versioning event-sourced systems. Its table of contents is a map of the costs this page warns about: why you cannot update an event (immutability, consumers, audit), type-based versioning, weak schema, negotiation, versioning snapshots and behaviour, compensating actions for mistakes ("accountants use pens"), copy-and-replace migrations including changing stream boundaries on a running system, copy-transform, and "versioning bankruptcy" (Versioning in an Event Sourced System).
Staff insight: That a whole book exists on versioning is itself the interview point. When a candidate says "events are immutable so it's simple," naming compensating events, upcasting and copy-and-replace — and when each applies — is what separates having run the pattern from having drawn it.
Practice Drill#
Prompt: "A subscription billing platform stores each customer's subscription as a row that is updated in place: plan, seats, status, next invoice date. Finance has found that proration was calculated wrong for mid-cycle seat changes for four months, and wants to know exactly which invoices were affected and to issue corrections. Meanwhile the customer dashboard, the revenue report and the dunning system all query the same table and keep adding indexes, and writes have slowed from 8ms to 40ms p99. The team proposes 'moving everything to CQRS and event sourcing.' What do you do?"
Staff Answer
There are two problems and they want different amounts of the pattern. The read contention is a CQRS problem: the dashboard, revenue report and dunning each want a different shape, and they're all fighting the write table's indexes. I'd keep the subscription table as the write model for now, add an outbox in the same transaction as each change (SeatsChanged, PlanChanged, SubscriptionCancelled, with event_id, subscription_id, version, schema_version), and feed three projections through CDC: a dashboard view in PostgreSQL keyed by customer, a revenue fact table in the warehouse, and a dunning queue. Each projector checkpoints in the same transaction as its upsert and guards on version, so duplicates are no-ops. Expected lag is under 2s; the dashboard reads its own writes by waiting on the returned version for up to 1s. Dropping the extra indexes from the write table should bring writes back near 8ms. The proration problem is a history problem, and that's where event sourcing earns its cost — but only for the billing ledger, not for subscriptions. Invoices and their line items become an append-only ledger: InvoiceIssued, LineItemCharged with proration_rule_version in metadata, and corrections as CreditNoteIssued events that reference the original line — never edits. For the past four months, we reconstruct what happened from invoices, the audit table and the payment provider's records, identify affected line items, and append credit notes; from now on, the rule version on every line makes the next investigation a query. Balances and revenue are projections over the ledger, rebuildable into a shadow table: ~300M ledger events at a 50K/s rebuild rate is under 2 hours, which I'd publish as the rebuild SLO. I would not event-source subscriptions yet — nobody has asked for their state at an arbitrary past instant, and the outbox events give us a seed if they do. Owners: billing team owns the ledger schema and upcasters; each read-model team owns its projector, lag alert and rebuild runbook; finance owns closed-period checksums.
Why this is L6:
- Separates the read-shape problem (CQRS on the existing table) from the history problem (event-source only the ledger), instead of applying the full pattern everywhere.
- Fixes past mistakes with compensating events and makes future ones queryable through rule-version metadata.
- Sizes the rebuild, defines read-your-writes behaviour, and assigns owners for projectors, schemas and closed-period checks.
What L7 adds:
- Treats event-sourcing the ledger as a one-way door needing a review and a ten-year owner, and pushes for a shared projector framework if other teams follow.
- Prices it: the outbox and three projections are weeks of work; the ledger is a permanent ~20% tax on the billing team, justified by the cost of four months of incorrect invoices.
- Sets a policy that closed accounting periods are checksummed so no replay or upcaster can silently change them.
❌ Common L5 Trap
"We'll move to CQRS and event sourcing. Every change becomes an event on a Kafka topic, a consumer builds the subscription table and the read models from the topic, and to fix proration we update the calculation and replay all events to regenerate correct invoices."
Why this misses: It event-sources everything when only billing history needed it, and Kafka-as-store loses the per-subscription conditional append, so concurrent seat changes can both "win." Worse, replaying with new logic to regenerate invoices rewrites issued financial documents — the correction must be new credit-note events appended to history, not a recomputed past. And there is no rebuild time, no read-your-writes plan and no owner for the projections.
Staff Interview Application#
How to Introduce This Pattern#
"First I'll separate two questions: do reads and writes need different shapes, and does any entity need its history as the source of truth? The first gets CQRS — read models fed from the write store. Only entities with a real history requirement get event-sourced, with one stream per aggregate and a version check on append. Then I'll talk about projection lag, rebuild time and how old events stay readable."
Lead with scope (which entities), then the consistency boundary, then projections and their lag, then rebuilds and versioning, then owners.
When NOT to Use This Pattern#
- CRUD domains with an audit need: profiles, settings, catalogues. A table and an audit log answer "who changed what." Event sourcing adds lifetime cost for a question nobody asks.
- Simple read/write ratios: if one well-indexed table and a cache serve reads within SLO, a separate read model is extra lag for no benefit — see Read-Heavy Systems.
- Data that must be truly erasable: unless PII stays out of events or crypto-shredding is accepted by legal, immutable logs and deletion rights collide.
- Teams without time to own it: the costs arrive in years one to three. If nobody will own projections and versioning for that long, ship CQRS on a table with an outbox instead.
- Cross-service transactions: event sourcing inside one service doesn't coordinate several. That's a saga or workflow question.
- Analytics on raw events: keeping immutable click or telemetry logs for recompute is a batch and stream pipeline, not an event-sourced domain model.
Follow-Up Questions to Anticipate#
| Interviewer Asks | What They Are Testing | How to Respond |
|---|---|---|
| "Do you need event sourcing to do CQRS?" | Whether you can separate the two | "No. CQRS is separate read models; they can come from a normal table via CDC. I only event-source entities whose history is the asset." |
| "How do you prevent two commands corrupting one aggregate?" | Consistency boundary | "Conditional append on the expected stream version. The loser reloads and re-decides, or gets a 409." |
| "A user saves and the screen doesn't show it. Why, and what do you do?" | Read-your-writes | "Projection lag. The command returns the new version and the query waits for the projector's checkpoint to pass it, up to a second." |
| "How long does a projection rebuild take?" | Whether you've sized replay | "Events divided by sink write rate. 2B at 50K/s is 11 hours, so it's always a shadow build and a swap." |
| "You need to change an event's structure. How?" | Schema evolution | "Additive needs nothing. Shape changes get an upcaster. Meaning changes get a new event type. Wrong facts get compensating events." |
| "What about GDPR deletion?" | Immutability vs law | "PII stays out of events by reference, or is encrypted per subject and the key deleted. Decided before the first event is written." |
Scorecard#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | Applies CQRS/ES system-wide | Separates CQRS from ES; event-sources only entities with a history requirement | Gates ES org-wide with review and long-term ownership |
| Consistency | "Eventually consistent" | Expected-version append per aggregate; stated read lag; read-your-writes by version wait | Product sign-off per surface on staleness |
| Projections | One consumer builds the read DB | Idempotent by version; checkpoint with effect; shadow rebuild and swap; rebuild time sized | Projector framework and rebuild SLOs as a platform |
| Evolution | "Version the events" | Additive, upcasters, new types for meaning changes, compensating events | Event catalogue, golden-corpus CI, bankruptcy policy |
| Operations | Watches consumer lag | Reconciliation against stream versions; commit-order subscriptions; snapshot discipline | Prices people-years; quarterly rebuild game days |
Strong Hire Signals
| Signal | What It Sounds Like |
|---|---|
| Separates the two patterns | "CQRS for the read shapes; event sourcing only for the ledger." |
| Knows the boundary | "Strong consistency inside one stream, by version check. Eventual everywhere else." |
| Has timed a replay | "Two billion events at 50K a second is 11 hours. Shadow and swap." |
| Treats history legally | "Mistakes are corrected with compensating events. Closed periods never change." |
Lean No-Hire Signals
| Signal | Why It Misses the Bar |
|---|---|
| "Kafka is our event store" with no conditional append | Concurrent commands can both append the same version |
| "Just replay" with no duration | Rebuild becomes a multi-hour outage of the read side |
| Edits stored events to fix bugs | Destroys the audit property that justified the pattern |
Common False Positives: Drawing a command bus and a query bus ≠ knowing where consistency lives. Knowing an event store's API ≠ having a versioning strategy. "Event-driven" ≠ event-sourced — publishing events from a CRUD service is the outbox, not this pattern.
Capacity Planning Quick Reference#
Sizing the Event Store and Rebuilds#
events_total = avg_events_per_sec × 86,400 × 365 × years # 50/s × 3y ≈ 4.7B
event_store_bytes = events_total × avg_event_bytes × ~1.5 (indexes) # 4.7B × 800B × 1.5 ≈ 5.6 TB
rebuild_hours = events_to_replay / sink_write_rate / 3,600 # 2B / 50K/s ≈ 11 h
aggregate_load_ms ≈ events_since_snapshot × ~0.02–0.05 ms # 200 events ≈ 4–10 ms
snapshot_every = events where load_ms crosses ~10 ms # typically 100–500
conflict_rate ≈ concurrent_commands_per_stream × command_time # > 5% => shrink aggregate
projection_lag = commit_to_subscription + batch_wait + apply # 100 ms – 2 s healthy
read_your_writes_wait ≤ 1 s, then serve with a 'still updating' marker
Key Numbers Worth Memorizing#
| Number | Context |
|---|---|
| 100 ms – 2 s | Healthy async projection lag |
| ~1,000 events | Comfortable stream length without snapshots |
| 100–500 events | Typical snapshot interval, or when load time passes ~10 ms |
| > 5% | Command conflict rate that says the aggregate is too large or too hot |
| 5–20K events/s | Single projector replay rate; 50–150K/s partitioned, sink permitting |
| 2B ÷ 50K/s ≈ 11 h | Why rebuilds are shadow-and-swap, never in place |
| ~5–20K appends/s | Typical ceiling for a single PostgreSQL primary as an event store |
| < 1 minute | LMAX's documented restart time: nightly snapshot plus a day of replay |
| 6 million orders/s | LMAX's single-threaded business logic processor, per its architecture write-up |
| ~20% | Rough share of a team's time an event-sourced domain consumes in years 1–3 (illustrative) |
Common Pitfalls Checklist#
- Event sourcing scoped to aggregates with a stated history requirement; everything else is CQRS or CRUD
- Appends conditional on expected stream version, enforced atomically by the store
- Subscriptions read in commit order; no checkpoint can skip an in-flight event
- Projectors idempotent by version, checkpoint stored with the effect
- Every projection has a published full-rebuild time and a shadow-and-swap path
- Snapshots versioned and deletable; tests pass with the snapshot table empty
-
schema_versionon every event; upcasters pure and tested against a golden corpus - Wrong facts corrected with compensating events, never edits
- PII kept out of events or crypto-shredded per subject
- Hourly reconciliation of projection versions against stream heads