System Design Problem

Design an Event Sourcing System

Commonly Asked By:StripeSquareUberNetflix

Interview Setup

Interview Prompt

Design an event sourcing system where application state is derived from an immutable append only event log. Support CQRS read models, snapshots, sagas, and temporal as-of queries.

Clarifying Questions (ask before designing)

QuestionWhy it matters
Write volume: expected events per second and total aggregate count?Operating at 5B events per day with ~58K average and 200K peak QPS dictates Kafka partition sizing, database sharding, and snapshot frequencies.
How many distinct read models and projections must be derived from one event stream?CQRS fan-out requires each projection to run as an independent consumer group, which directly influences replay throughput and replication lag SLOs.
What ordering guarantees are required: per aggregate or global total ordering?Per-aggregate stream ordering is naturally maintained by partitioning on stream ID, whereas global total ordering requires a sequential global position sequence or a centralized sequencer.
What is the target snapshot interval, and what is the maximum acceptable replay duration?Replaying 10K events across 500M aggregates without intermediate snapshots makes synchronous read operations unacceptably slow.

Scope

In scope

  • Event log as the authoritative immutable source of truth
  • Command Query Responsibility Segregation (CQRS) architecture
  • Materialized read projections and real-time view generation
  • Snapshot checkpoint optimization for low latency state hydration
  • Saga coordination pattern for distributed multi-aggregate workflows
  • Point-in-time temporal queries and state reconstruction

Out of scope (state explicitly)

  • Enterprise data warehouse star-schema modeling and batch analytics pipelines
  • End-to-end exactly-once processing guarantees across heterogeneous third-party external consumers
  • Custom message broker implementations built from scratch

Functional Requirements

When discussing requirements with your interviewer, establish the append only event write path first. Clarify event append semantics, state hydration through sequential replay, independent CQRS projections for read serving, and snapshot intervals to bound reconstruction costs. Confirm whether temporal point in time queries and schema upcasting fall within scope.

In an interview setting, emphasize that when a read model projection bug ships to production, the immutable event log allows complete state reconstruction from offset zero.

  • Append Immutable Events: Persist all state transitions as an immutable, ordered stream of business events rather than updating records in place.
  • Rebuild Aggregate State: Reconstruct the current state of any domain entity at runtime by loading snapshots and replaying subsequent events.
  • Dynamic Event Replay: Replay historical event logs to repair corrupted read models, bootstrap new projections, or recover from software bugs.
  • Materialized Snapshots: Periodically persist aggregate checkpoints to bound replay duration and keep command latency predictable.
  • Independent Projections: Generate and maintain multiple denormalized, read optimized database views derived from a single write stream.
  • Point-in-Time Temporal Queries: Enable historical time-travel queries to inspect the precise state of any entity as of a specific past timestamp.
  • Schema Evolution: Support schema upcasting layers to transform legacy event payloads into contemporary formats during read time deserialization.
  • Publish-Subscribe Streaming: Stream committed domain events to downstream microservices and analytics pipelines through reactive event buses.

Non-Functional Requirements

Immutability and strict per-stream ordering represent mandatory system invariants. When interviewers ask how query performance remains fast while the event log grows continuously, explain how denormalized CQRS projections and periodic snapshots decouple write volume from read latency.

  • Durability: Events serve as the sole source of truth, with an 11 nines durability target within the authoritative primary replication domain and zero accepted logical data loss there. Cross region disaster recovery uses asynchronous replication with an explicitly measured recovery point objective.
  • Immutability: Storage engines must enforce append only constraints, protecting historical events from in-place updates, deletions, or tampering.
  • Strict Per-Stream Ordering: Events belonging to an individual aggregate stream must maintain sequential ordering without missing revisions.
  • Idempotent Commands: Retries of the same logical command must return the original result without appending duplicate events.
  • High Write Throughput: Sustain 58K average and 200K peak event write requests per second.
  • Low Read Latency: Serve queries from denormalized projection stores in under 10 milliseconds at p99.
  • Horizontal Scalability: Scale across billions of cumulative events and 500M distinct aggregate entities without centralized lock bottlenecks.
  • Backward Compatibility: Ensure older serialized event structures can be deserialized and processed by newer application deployments.

Capacity Estimations

Event volume grows continuously over time because records are never deleted. Sizing the append only event store, network ingress, and projection fan-out upfront ensures the design handles both peak transaction bursts and historical backfill loads.

MetricCalculationValue
Events / dayGiven (typical workload assumption)5 Billion
Events / secFrom Events / day ÷ 86,400 (+ peak factor in value)~58K avg (peak 200K)
Avg event sizeGiven (typical workload assumption)500 bytes
Write throughputDerived29 MB/s avg, 100 MB/s peak
Storage / dayDerived from upstream throughput x size2.5 TB
Storage / yearDaily storage x 365~900 TB
Aggregates (entities)Derived500 Million
Avg events per aggregateGiven (typical workload assumption)50
I/O and Bandwidth Derivations:
1. Daily Telemetry and Event Ingestion:
   - 5 Billion events/day x 500 bytes/event = 2.5 TB per day.
2. Peak Write Throughput:
   - 200,000 events/sec peak x 500 bytes = 100 MB/s network ingestion rate.
3. Multi-Year Aggregate Storage Footprint:
   - 2.5 TB/day x 365 days = ~912 TB/year.
   - Handled reliably through distributed object storage tiering and sharded database storage blocks.

Architecture Diagram

The architecture establishes a strict separation between command processing on the write path and projection serving on the query path. The command handler validates business invariants against hydrated aggregate state, appends immutable events to the PostgreSQL event store using optimistic concurrency control, and streams committed changes to Kafka through Change Data Capture. Dedicated projection workers consume events from Kafka and materialize specialized read models across PostgreSQL, ClickHouse, and Redis. Periodic snapshots cap aggregate replay duration so that commands never replay thousands of historical events during execution.

In an interview, frame CQRS as an architectural necessity rather than an optimization because write optimized append logs cannot satisfy low latency multi-dimensional query patterns.

Loading...

Component Deep Dives

The architecture centers around the event store as the sole source of truth, supported by snapshot managers that bound state hydration latency and projection engines that construct specialized query views.

Event Store: Aggregate Streams and Optimistic Concurrency

The event store is the authoritative append only source of truth. Each domain entity is represented as an isolated sequence of versioned events.

Optimistic Concurrency Control checks the expected revision number during append operations, preventing concurrent writers from introducing race conditions without requiring long held database row locks.

JSON
[
  {
    "event_id": "evt-001",
    "aggregate_type": "order",
    "stream_id": "order-12345",
    "schema_version": 1,
    "stream_version": 1,
    "event_type": "OrderCreated",
    "data": { "customer_id": "cust-8821", "items": [{ "product_id": "prod_1", "qty": 2 }], "total": 100 },
    "timestamp": "2026-03-15T10:00:00Z"
  },
  {
    "event_id": "evt-002",
    "aggregate_type": "order",
    "stream_id": "order-12345",
    "schema_version": 1,
    "stream_version": 2,
    "event_type": "PaymentReceived",
    "data": { "amount": 100, "currency": "USD", "method": "card" },
    "timestamp": "2026-03-15T10:02:15Z"
  },
  {
    "event_id": "evt-003",
    "aggregate_type": "order",
    "stream_id": "order-12345",
    "schema_version": 1,
    "stream_version": 3,
    "event_type": "InventoryAllocated",
    "data": { "warehouse_id": "wh-east-1" },
    "timestamp": "2026-03-15T10:04:00Z"
  },
  {
    "event_id": "evt-004",
    "aggregate_type": "order",
    "stream_id": "order-12345",
    "schema_version": 1,
    "stream_version": 4,
    "event_type": "OrderShipped",
    "data": { "tracking_number": "UPS123" },
    "timestamp": "2026-03-15T14:30:00Z"
  }
]

Compare-And-Swap (CAS) Concurrency Protocol:

When appending an event to stream "order-12345":
  Expected version: 3 (client read version 3 prior to command validation)
  Attempted append version: 4
  
  Concurrency Check:
    IF current stream version != 3:
      Trigger optimistic concurrency conflict (another writer has appended version 4).
      Action: Client must reload latest aggregate state, re-evaluate business invariant rules, and retry append.
    ELSE:
      Append succeeds atomically, establishing version 4.

Key Architectural Guarantees:
  - Compare-And-Swap (CAS) executes at the aggregate stream boundary.
  - Avoids long held application-level locks and removes this source of lock contention and deadlocks.

Event Sourcing vs Traditional CRUD:

SQL
-- Traditional CRUD (Destructive In-Place Mutation):
UPDATE orders SET status = 'shipped' WHERE id = 12345;
-- Previous state is overwritten and lost. Auditability requires separate manual logging.

-- Event Sourcing (Append-Only State Transition):
INSERT INTO events (stream_id, stream_version, event_type, data)
VALUES (
    'order-12345',
    5,
    'OrderCancelled',
    '{"reason": "customer_request", "cancelled_by": "user-789"}'::jsonb
);
-- Complete immutable lineage is preserved. Every state transition is intrinsically auditable.

Materialized Snapshots for Replay Optimization

When an aggregate accumulates thousands of state transitions over its lifetime, hydrating state from version one on every command introduces unacceptable processing delays. The snapshot service checkpoints serialized aggregate state at fixed revision boundaries to cap reconstruction times:

State Reconstruction Workflow:
  1. Fetch latest snapshot: Load snapshot state at stream version 9,900.
  2. Sequential event replay: Replay only delta events from version 9,901 through 10,000 (100 events instead of 10,000).
  3. In-memory convergence: Reconstruct current aggregate state in under 5 milliseconds.

Snapshot Trigger Strategies:
  - Event interval threshold: Snapshot every N events (for example, every 100 events).
  - Data size threshold: Snapshot when cumulative aggregate payload size exceeds a specific memory limit.
  - On-demand hot aggregate trigger: Snapshot dynamically when an aggregate experiences elevated read or command velocity.

CQRS Projections and Denormalized Read Models

Command Query Responsibility Segregation isolates write throughput from read traffic. While the write store prioritizes sequential append latency, projection workers asynchronously ingest events from Kafka to maintain read optimized representations tailored to specific user experiences:

Event Stream Consumer  -->  Projection Worker  -->  Materialized Read Database

Projection 1: Order Details (PostgreSQL Read Model)
  OrderCreated:
    INSERT INTO orders (id, items, total, status) VALUES (event.stream_id, event.data.items, event.data.total, 'created');
  PaymentReceived:
    UPDATE orders SET status = 'paid', payment_method = event.data.method WHERE id = event.stream_id;
  OrderShipped:
    UPDATE orders SET status = 'shipped', tracking = event.data.tracking WHERE id = event.stream_id;

Projection 2: Daily Revenue Dashboard (ClickHouse OLAP Analytics)
  PaymentReceived:
    INSERT INTO daily_revenue (event_date, amount, currency) VALUES (event.metadata.timestamp::date, event.data.amount, event.data.currency);

Projection 3: Real-Time Order Velocity (Redis In-Memory Counter)
  OrderCreated:
    INCR user:{event.data.customer_id}:order_count

API Design

API contracts reflect strict CQRS boundaries where command endpoints append versioned events to the log and query endpoints retrieve denormalized state from specialized projections.

Client API Type Definitions

TypeScript domain models defining immutable event envelopes, command payloads, snapshot records, and event store contracts:

TYPESCRIPT
// Domain models and client contracts for event sourced aggregates, commands, and temporal queries

export type EventId = string;
export type StreamId = string;
export type AggregateType = "order" | "account" | "inventory" | "user";

export interface EventEnvelope<T = Record<string, unknown>> {
  eventId: EventId;
  aggregateType: AggregateType;
  globalPosition: number;
  streamId: StreamId;
  streamVersion: number;
  schemaVersion: number;
  eventType: string;
  data: T;
  metadata: {
    correlationId: string;
    causationId?: string;
    actorId: string;
    commandId?: string;
    timestamp: string;
  };
}

export interface AppendEventsRequest {
  streamId: StreamId;
  expectedVersion: number;
  commandId: string;
  payloadHash: string;
  events: Array<{
    eventType: string;
    data: Record<string, unknown>;
  }>;
}

export interface AppendEventsResponse {
  streamId: StreamId;
  currentVersion: number;
  appendedCount: number;
}

export interface AggregateSnapshot<T = Record<string, unknown>> {
  streamId: StreamId;
  streamVersion: number;
  state: T;
  createdAt: string;
}

export interface TemporalState<T = Record<string, unknown>> {
  streamId: StreamId;
  asOf: string;
  streamVersion: number;
  state: T;
}

export interface EventStoreClient {
  appendEvents(request: AppendEventsRequest): Promise<AppendEventsResponse>;
  readStreamEvents(streamId: StreamId, fromVersion?: number, toVersion?: number): Promise<EventEnvelope[]>;
  readStreamStateAt<T = Record<string, unknown>>(streamId: StreamId, timestamp: string): Promise<TemporalState<T>>;
  getLatestSnapshot<T = Record<string, unknown>>(streamId: StreamId): Promise<AggregateSnapshot<T> | null>;
}

Command Execution (Write Path)

Accepts client intent, evaluates domain business rules, and appends an immutable event to the aggregate stream:

HTTP
POST /api/v1/orders HTTP/1.1
Host: api.eventsourcing.example.com
Authorization: Bearer <token>
Idempotency-Key: cmd-7f3a9c1e
Content-Type: application/json

{
  "items": [{ "product_id": "prod-91", "qty": 2 }],
  "customer_id": "cust-8821"
}

HTTP/1.1 201 Created
Content-Type: application/json

{
  "stream_id": "order-12345",
  "stream_version": 1,
  "event_type": "OrderCreated",
  "status": "created"
}

POST /api/v1/orders/order-12345/ship HTTP/1.1
Host: api.eventsourcing.example.com
Authorization: Bearer <token>
Idempotency-Key: cmd-ship-7f3a9c1e
If-Match-Version: 3
Content-Type: application/json

{
  "tracking_number": "UPS981242"
}

HTTP/1.1 200 OK
Content-Type: application/json

{
  "stream_id": "order-12345",
  "stream_version": 4,
  "event_type": "OrderShipped",
  "status": "shipped"
}

Event Queries (Read Path)

Retrieves materialized views, full historical audit logs, or temporal state reconstructed as of a specific past timestamp:

HTTP
GET /api/v1/orders/order-12345 HTTP/1.1
Host: api.eventsourcing.example.com
Authorization: Bearer <token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "order_id": "order-12345",
  "version": 4,
  "status": "shipped",
  "items": [{ "product_id": "prod-91", "qty": 2 }],
  "tracking_number": "UPS981242",
  "last_updated": "2026-03-15T14:30:00Z"
}

GET /api/v1/orders/order-12345/history HTTP/1.1
Host: api.eventsourcing.example.com
Authorization: Bearer <token>

HTTP/1.1 200 OK
Content-Type: application/json

[
  { "version": 1, "type": "OrderCreated", "data": { "customer_id": "cust-8821", "items": [{ "product_id": "prod-91", "qty": 2 }], "total": 100 } },
  { "version": 2, "type": "PaymentReceived", "data": { "amount": 100, "currency": "USD", "method": "card" } },
  { "version": 3, "type": "InventoryAllocated", "data": { "warehouse_id": "wh-east-1" } },
  { "version": 4, "type": "OrderShipped", "data": { "tracking_number": "UPS981242" } }
]

GET /api/v1/orders/order-12345/state-at?timestamp=2026-03-15T10:03:00Z HTTP/1.1
Host: api.eventsourcing.example.com
Authorization: Bearer <token>

HTTP/1.1 200 OK
Content-Type: application/json

{
  "order_id": "order-12345",
  "as_of": "2026-03-15T10:03:00Z",
  "reconstructed_version": 2,
  "state": {
    "status": "paid",
    "total": 100
  }
}

Common Error Responses

Standardized HTTP error representations for version conflicts, invariant rejections, and missing aggregates:

400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff

Data Model

The storage architecture divides between the authoritative event log in PostgreSQL and the distributed streaming topology in Kafka.

PostgreSQL Event Store and Snapshot Tables

Relational DDL defining append only event logging with unique compound version constraints alongside cached aggregate snapshots:

SQL
-- Primary event store schema
CREATE TABLE events (
    global_position  BIGSERIAL,             -- Monotonic allocation position within this logical event store, though gaps can occur on rollback and it does not guarantee cross-stream causality
    event_id         UUID NOT NULL,         -- Stable event identity for deduplication and tracing
    aggregate_type   VARCHAR(64) NOT NULL,  -- Aggregate type used by consumers and tooling
    stream_id        VARCHAR(256) NOT NULL, -- Unique aggregate identifier (for example, "order-12345")
    stream_version   INT NOT NULL,          -- Per-aggregate sequence revision number
    schema_version   INT NOT NULL DEFAULT 1, -- Payload schema version, independent of stream ordering
    event_type       VARCHAR(128) NOT NULL, -- Domain event name (for example, "OrderCreated")
    data             JSONB NOT NULL,        -- Immutable event payload data
    metadata         JSONB,                 -- Correlation, causation, and actor tracing metadata
    created_at       TIMESTAMP NOT NULL DEFAULT NOW(),
    PRIMARY KEY (stream_id, stream_version),
    UNIQUE (event_id),
    UNIQUE (global_position)
);
-- Command idempotency table prevents a retried client command from appending duplicate events.
CREATE TABLE processed_commands (
    actor_id         VARCHAR(256) NOT NULL,
    command_id       VARCHAR(256) NOT NULL,
    stream_id        VARCHAR(256) NOT NULL,
    payload_hash     VARCHAR(128) NOT NULL,
    result_version   INT NOT NULL,
    created_at       TIMESTAMP NOT NULL DEFAULT NOW(),
    PRIMARY KEY (actor_id, command_id)
);
-- Reusing a command_id with a different payload_hash is rejected as an idempotency conflict.

-- Snapshot cache table
CREATE TABLE snapshots (
    stream_id        VARCHAR(256) PRIMARY KEY,
    stream_version   INT NOT NULL,
    state            JSONB NOT NULL,
    created_at       TIMESTAMP NOT NULL DEFAULT NOW()
);

Kafka Event Bus Topic Topology

Topic configurations and partition key guarantees that enforce strict message sequencing per aggregate stream:

YAML
topic: domain-events
partitioning:
  partition_key: stream_id
  guarantee: "Strict total ordering within each individual aggregate stream"
schema:
  envelope:
    event_id: uuid
    aggregate_type: string
    global_position: integer
    stream_id: string
    stream_version: integer
    schema_version: integer
    event_type: string
    payload: object
    metadata:
      correlation_id: string
      causation_id: string
      actor: string
      command_id: string
    created_at: timestamp
configuration:
  cleanup_policy: delete
  retention_ms: -1
  retention_strategy: "Retain complete immutable event history indefinitely without log compaction"
  replication_factor: 3
  min_insync_replicas: 2

Fault Tolerance

Resilience strategies address dual-write inconsistencies, projection worker outages, read model corruption, and asynchronous projection lag.

ConcernSolution
Event Data LossSynchronous PostgreSQL primary-standby replication protects the authoritative event store within the primary replication domain, while Kafka replication factor of 3 and min.insync.replicas of 2 protect downstream event delivery. Cross region recovery remains bounded by asynchronous database replication lag.
Projection Builder FailureEach worker resumes from its committed Kafka consumer offset or another durable checkpoint without skipping successfully processed events.
Projection DB CorruptionRebuild the corrupted read model from scratch by replaying historical events from offset zero into a clean table.
Event Schema DriftImplement dynamic upcaster layers that translate legacy event formats into contemporary structures during read time deserialization.

Resolving the Dual-Write Delivery Gap

The event store commit and Kafka publication need a reliable bridge so a successful database commit cannot diverge from downstream event delivery:

  • Transactional Outbox Pattern: Persists the domain event and an outbox record within a single database transaction, allowing a background relay worker to poll and publish to Kafka reliably.
  • Change Data Capture (CDC): Utilizes Debezium to stream commits directly from the PostgreSQL Write-Ahead Log (WAL) to Kafka topics without requiring dual application writes.

Mitigating Projection Eventual Consistency Lag

Because projections update asynchronously, users may query their views before background workers ingest recent commits:

  1. Read-Your-Own-Writes: Returns the computed aggregate mutation directly in the command response so client applications can update their local state optimistically.
  2. Version-Assertion Queries: Command responses return the newly assigned stream revision, allowing read gateways to poll projections until the materialized version matches or exceeds that revision.

Additional Considerations

Advanced architectural patterns covering schema migration over immutable logs, boundary guidelines for event sourcing adoption, and distributed transaction management.

Dynamic Schema Evolution with Upcasters

As business requirements evolve, domain event schemas inevitably change through added fields, deprecated flags, or structural renames. Because historical events are immutable and cannot be rewritten in storage, the system applies read time upcasters to transform legacy representations on the fly:

TYPESCRIPT
// Dynamic event upcasting layer applied during read time deserialization
// schemaVersion describes payload shape, while streamVersion represents the aggregate revision
interface StoredEventV1 {
  schemaVersion: 1;
  type: "OrderCreated";
  data: { items: Array<{ productId: string; qty: number }>; total: number };
}

interface StoredEventV2 {
  schemaVersion: 2;
  type: "OrderCreated";
  data: { items: Array<{ productId: string; qty: number }>; subtotal: number; tax: number; total: number };
}

export function upcastOrderCreated(event: StoredEventV1): StoredEventV2 {
  return {
    ...event,
    schemaVersion: 2,
    data: {
      items: event.data.items,
      subtotal: event.data.total,
      tax: 0,
      total: event.data.total,
    },
  };
}

When to Avoid Event Sourcing

Event sourcing introduces notable engineering complexity through eventual consistency, projection maintenance, and schema migration overhead. Teams should avoid this pattern when designing straightforward CRUD applications without historical audit needs, when entity schemas change too unpredictably to maintain upcasters, or when strict synchronous read consistency is a firm architectural requirement.

Saga Choreography vs Orchestration

Managing distributed transactions across multiple aggregate boundaries requires choosing between choreography and orchestration.

In choreographed sagas, services listen to domain events and react autonomously. For example, an OrderCreated event triggers payment processing, which subsequently emits a payment confirmation event. This approach provides loose coupling for simple two-to-three step workflows, but becomes difficult to trace when business flows branch conditionally.

In orchestrated sagas, a dedicated saga manager coordinates the entire workflow by dispatching explicit commands and tracking intermediate state transitions in a dedicated saga event stream. Orchestration provides centralized observability and cleaner failure management when coordinating five or more steps. Regardless of the coordination model, sagas maintain immutability by appending forward-compensating events such as VoidPayment or ReleaseInventory rather than deleting historical records.

Interview Walkthrough

A structured pacing guide for navigating a 45-minute event sourcing architecture interview:

  • 25-minute cut

    Focus on the write path, optimistic concurrency, and CQRS projection mechanics before introducing advanced replay features.

    • Contrast CRUD overwrites with append only events and justify when auditability and replay warrant the complexity (5 min)
    • Draw the CQRS separation: event store for writes and asynchronous projection builders for denormalized reads (6 min)
    • Explain optimistic concurrency control using expected stream revisions on every append (5 min)
    • Cover snapshot checkpoints to cap state hydration latency and establish Kafka as the durable event backbone (5 min)
    • Staff depth: schema upcasting, point in time temporal queries, and blue green projection rebuilds (4 min)
  • Open by contrasting destructive CRUD overwrites with append only event streams, establishing early that event sourcing preserves complete audit lineage.
  • Draw the CQRS boundary clearly, positioning the PostgreSQL event store on the command side and Kafka-driven projection builders on the query side.
  • Detail optimistic concurrency control using expected stream versions to reject concurrent write conflicts without holding database row locks.
  • Introduce materialized snapshots for long lived aggregates so command handlers replay only recent delta events rather than entire multi-year histories.
  • Solve the database to Kafka delivery gap using either the Transactional Outbox pattern or Debezium CDC so committed events are published reliably without an application-level dual-write race.
  • Explain schema evolution through read time upcasters because stored historical events are immutable and cannot be modified in place.
  • Quantify write scale early: 5B events per day at 500 bytes per event generates 2.5 TB daily, requiring uniform stream partitioning across Kafka brokers.
  • Highlight the critical architectural pitfall of attempting to serve user-facing queries by replaying raw event streams on every request, which leads to unacceptable query latency.

Related Problems

Event sourcing and CQRS patterns connect directly to these specialized distributed systems architectures and core infrastructure concepts:

Engineering Trade-offs

Evaluating storage engine trade-offs, consistency models, and operational verification strategies across event sourced systems.

Event Store Storage Engine Selection

Evaluating relational databases, specialized event stores, and distributed message brokers for immutable event logging:

FeaturePostgreSQLEventStoreDBKafka
Optimistic ConcurrencySupported via UNIQUE (stream_id, stream_version) constraintsSupported via native revision version checksNot supported (lacks partition-key version checks)
Stream SubscriptionsPolling or LISTEN/NOTIFY notification, with polling required for durable catch-upSupported via native catch-up client subscriptionsSupported via consumer group offsets
Ordering GuaranteePer-stream ordering plus a monotonic global position within one logical storePer-stream ordering plus global positionStrictly partition-level ordering
Query FlexibilityFull relational SQL queries and indexesLimited secondary indexingOffset-based sequential scans only
Scalability & IngestionRequires manual active shardingClustered replicasHigh throughput across partitioned streams
Operational MaturityExtremely matureSpecialized niche storage engineExtremely mature

Recommendation: Use PostgreSQL while the sustained append rate fits the chosen primary's measured capacity because its transactional maturity and relational SQL tooling minimize operational complexity. At sustained rates above roughly 100,000 events per second, evaluate partitioning or sharding the event store or adopting a specialized event store such as EventStoreDB. Leverage Kafka primarily as a high throughput event delivery bus rather than an authoritative event store, because Kafka lacks stream-level optimistic concurrency checks on partition keys.

Snapshot Consistency and Verification

If a snapshot is written with state corrupted by a software defect, all subsequent state hydrations will inherit that defect. Teams protect state integrity by running asynchronous verification jobs that periodically replay raw historical events from stream version zero and compare the resulting state against existing snapshot checkpoints, rebuilding any snapshots that exhibit drift.

💬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...