Core Concept

Event Sourcing & CQRS

Event sourcing stores state changes as an immutable append-only log; CQRS separates write models (commands) from read models (queries) so each side can scale independently.


1. What It Is

When audit history and temporal queries matter more than current state, we store events instead of overwriting rows. CQRS often pairs with event sourcing to separate write and read models.

What:

Event Sourcing records all state updates as a sequence of immutable events; CQRS separates the logic for updating data (Commands) from querying data (Queries).

Primary purpose:

Scaling high-concurrency systems, providing perfect audit histories, and optimizing read/write database workloads separately.

Usually used for:

Financial ledgers, collaborative design tools (Figma), e-commerce checkouts, and high-volume analytics systems.

2. Core Mental Model

Event sourcing and CQRS are independent choices. You can split read/write models without an event log, or keep an event log without CQRS. Most production systems use CQRS alone β€” see the split below before committing to full ES.

🎬 Immutable Event Stream

Never execute database `UPDATE` or `DELETE` statements. State updates are strictly modeled as append-only events (`OrderCreated`, `OrderCancelled`).

πŸ”ͺ Strict Command/Query Split

Commands (writes) return zero data (only validation state). Queries (reads) bypass the write datastore, pulling directly from optimized read replicas.

πŸ“Έ Snapshots & Command Sourcing

Replay from event 0 is O(N). Snapshots store aggregate state at offset K; rebuild loads snapshot + events since K. Command sourcing validates commands against current state before appending β€” the write model holds authoritative state, the log holds the audit trail.

In the room

Event sourcing is powerful but operationally heavy β€” don't propose it for a simple CRUD app. Justify it with requirements: audit trail, replay, temporal queries. Mention snapshotting so replays don't scan millions of events.

3. Why It Matters in HLD

Event sourcing stores state as an append-only log of facts β€” CQRS separates read and write models. Three lenses:

Needed When:

A perfect historical audit trail is critical, or read query volumes are massive and require entirely different data models than writes.

Avoids:

Heavy aggregate queries locking the write database, lost history when rows are overwritten in place, and race conditions on shared mutable state.

Optimizes For:

Write throughput speeds, query latencies, audit trace accuracies, and independent service tier scaling capabilities.

4. Architecture & Data Flow

Walk the CQRS path as interview steps. Step 1 β€” Command: write side validates and appends event to the store. Step 2 β€” Event log: immutable sequence persisted (Kafka or event store). Step 3 β€” Projection: async consumer builds read-optimized views (SQL, Redis, Elasticsearch). Step 4 β€” Query: reads hit the projection, never replay the full log. Step 5 β€” Rebuild: new projection replays events from offset zero.

Loading...

In the room

Do not enable log compaction on a true event-sourcing topic β€” compaction drops intermediate events and destroys the audit trail. That detail separates textbook from production knowledge.

5. Key Characteristics

Event store vs traditional CRUD and sync vs async projections β€” we compare characteristics:

  • CQRS without event sourcing β€” write to a normalized OLTP schema; fan out changes to Elasticsearch or Redis via CDC/outbox. No replayable history, but far less operational overhead. Fits most e-commerce catalog and profile search splits.
  • Event sourcing without CQRS β€” append-only audit log with a single read model rebuilt from replay. Rare alone; usually paired when audit is mandatory.
  • Snapshot cadence β€” snapshot every N events or T minutes; store (aggregate_id, version, state_blob). On cold start: load latest snapshot, replay events with version > snapshot.version. Tune N so replay stays under your recovery SLA (typically <30s).
  • Read projection targets β€” one event stream, multiple query shapes:
Projection ModelTarget DatastoreArchitectural Rationale
Order Database ProjectionPostgreSQL Read ReplicaOptimized to handle basic user order detail point lookups (GET /orders/{id}).
Search Catalog ProjectionElasticsearch FleetEnables complex text search and faceted product filtering queries.
Analytics WarehouseClickHouse / SnowflakeAggregates massive volumes of historical order logs for financial reporting.

6. Strategic Tradeoffs

Auditability and temporal queries trade storage growth and complexity β€” we state both:

BenefitCost
Complete audit trail (every state change is a persisted event you can replay or inspect)Eventual read staleness (projections update asynchronously, so reads may lag writes briefly)
Independent read/write scaling (scale projection tiers without touching the write store)Operational overhead (event brokers, schema versioning, and projection workers to maintain)

7. Failure / Bottleneck Awareness

Projection lag, schema evolution, and snapshot storage β€” we volunteer failure modes:

🐌 The Event Schema Evolution Dilemma

Problem: Business requirements change over time, forcing schema updates on historical events. Since the event store is strictly append-only, modifying historical serialized JSON shapes directly is impossible.

Mitigation: Use upcasting in event processors β€” transform old event payloads to the current schema on replay β€” or run version-specific handlers side by side during migration.

🐒 Projection Lag & Read Staleness Spikes

Problem: Projection engines lagging behind the event stream cause clients to receive stale states (e.g. user updates profile but does not see changes on immediate screen reload).

Mitigation: Partition the event stream by entity key so projections process in order, and return optimistic UI state to the client while the read model catches up.

8. Common HLD Usage

Audit trails, collaborative editing, and financial ledgers use event sourcing patterns:

  • Bank Ledger Transfer (Event Sourcing): Balance is never stored directly. Instead, the balance is dynamically computed by summing the stream of credit/debit transaction event logs over time, ensuring auditability.
  • E-Commerce Catalog Search (CQRS): Product additions write to a transactional database, while search listings route to Elasticsearch projections, preventing catalog locks.

9. Decision Signals

Reach for event sourcing when complete history and replay matter; CQRS when read/write shapes diverge:

🎯 Think Event Sourcing & CQRS When:
  • You are building financial ledgers, transactional ledger sheets, or medical record tracking systems where every state change must remain auditable.
  • Your read query volume is massively disproportionate to writes (e.g. 1000:1 read-to-write ratio) and queries require completely different database structures.
  • You must reconstruct historic system states at specific point-in-time intervals for compliance or debugging.

11. Deep Dive (Optional)

Snapshot Internals & Command Sourcing

Snapshots are not a cache β€” they are a checkpoint that bounds replay time. A bank account with 10M transactions stores a snapshot at version 9,500,000 containing {"balance": 4200.50, "version": 9500000}. Recovery replays only events 9,500,001–10,000,000 to reach current state. Snapshots can live in the same event store or a separate key-value bucket keyed by aggregate_id.

Command sourcing means the write model loads current aggregate state (from snapshot + tail replay), validates the incoming command (Withdraw($500) fails if balance < 500), then appends the resulting event. The command itself is not stored β€” only the validated event β€” keeping the log a pure state-transition history.

Dual-Write Consistency Protection: CDC (Change Data Capture)

In CQRS architectures, updating the read projection database directly from the command handler is a dangerous pattern. If the write to the command DB succeeds, but the write to the read DB fails (or vice versa), the system drifts into permanent split-brain state.

Solution: Change Data Capture (CDC) via transactional outbox

Loading...
  1. The Command Handler writes business state and an outbox row in a single local transaction β€” no direct dual-write to the read model.
  2. A CDC engine (e.g., Debezium) tails the database's binary transaction logs (WAL) asynchronously.
  3. Debezium converts WAL/outbox rows into an event stream and publishes them to a Kafka topic.
  4. Projection processors consume events from Kafka, updating Elasticsearch or read replicas safely.

This completely decouples the write and read tiers, guaranteeing eventual consistency without introducing dual-write failure points.

πŸ’¬Review

Help Us Improve

How helpful was this walkthrough?

Click a star to rate. We actively use this feedback to refine and update our system design content.

Placeholder
Optional but highly appreciated!

Discussion

Share your thoughts, ask questions, or help others.

Loading comments...