PrepZone Logo
PrepZone

CQRS and Event Sourcing

Separate read and write models with an append-only event log as the source of truth.

Why this matters

  • StreamHub's creator analytics dashboard needs denormalised read models (views per hour, revenue by region) while the write path handles subscriptions and payments with strong consistency.
  • CQRS lets you scale reads and writes independently — a read-heavy dashboard does not slow down payment processing.
  • Event sourcing provides a complete audit trail: replay events to rebuild state, debug production issues, or migrate to a new schema.

CQRS + event sourcing components

  • Command handler — validates and executes business logic; emits domain events.
  • Event store — append-only log (Kafka, EventStoreDB) holding every state change.
  • Projections — consumers that build read models from the event stream.
  • Read store — denormalised views optimised for queries (Elasticsearch, Redis, materialised views).
  • Snapshots — periodic state checkpoints to avoid replaying millions of events on startup.

Architecture

CQRS on AWS

writeCDCqueryCOMPUTE
Command APIEKS
DATABASE
RDS write mod…normalised
INTEGRATION
Amazon MSKdomain events
ANALYTICS
OpenSearchread projection
CLIENT
Dashboardquery API
Commands write RDS; MSK events update OpenSearch read projections.

Command and event flow

Java
// Command: user subscribes to a creator
{
  "command": "SubscribeToCreator",
  "user_id": "usr_42",
  "creator_id": "crt_99",
  "tier": "premium",
  "idempotency_key": "sub_8f2a_20260401"
}
Java
// Event emitted to event store
{
  "event_type": "SubscriptionCreated",
  "aggregate_id": "sub_8f2a",
  "version": 1,
  "timestamp": "2026-04-01T12:00:00Z",
  "payload": {
    "user_id": "usr_42",
    "creator_id": "crt_99",
    "tier": "premium",
    "amount_cents": 999
  }
}
Java
// Projection updates read model
{
  "projection": "creator_subscriber_count",
  "creator_id": "crt_99",
  "action": "increment",
  "new_count": 14823
}

When to use CQRS

AspectGood fitPoor fit
Complex domain logicOrder fulfilment, billingSimple CRUD blog
Different read/write shapesDashboard vs transaction processingSame schema for reads and writes
Audit requirementsFinancial, healthcare, complianceNo audit trail needed
Temporal queriesWhat was the balance on March 1?Only current state matters
Team maturityExperienced with distributed systemsSmall team, simple product
  • Complex domain logic

    Good fitOrder fulfilment, billing
    Poor fitSimple CRUD blog
  • Different read/write shapes

    Good fitDashboard vs transaction processing
    Poor fitSame schema for reads and writes
  • Audit requirements

    Good fitFinancial, healthcare, compliance
    Poor fitNo audit trail needed
  • Temporal queries

    Good fitWhat was the balance on March 1?
    Poor fitOnly current state matters
  • Team maturity

    Good fitExperienced with distributed systems
    Poor fitSmall team, simple product

CQRS adds complexity. Use it when the read/write asymmetry or audit requirements justify the overhead.

Event store with Kafka

Java
def handle_subscribe_command(command: dict) -> None:
    subscription = Subscription.create(
        user_id=command["user_id"],
        creator_id=command["creator_id"],
        tier=command["tier"]
    )
    event = SubscriptionCreated(
        aggregate_id=subscription.id,
        version=1,
        payload=subscription.to_dict()
    )
    event_store.append(
        topic="streamhub.subscriptions",
        key=subscription.id,
        event=event
    )

Projections consume from the same topic:

Java
def projection_handler(event: dict) -> None:
    if event["event_type"] == "SubscriptionCreated":
        redis.hincrby(f"creator:{event['payload']['creator_id']}", "subscribers", 1)
        elasticsearch.index(
            index="subscriptions",
            id=event["aggregate_id"],
            body=event["payload"]
        )

Snapshots for performance

Replaying 500K events to rebuild one subscription aggregate is slow. Take periodic snapshots.

Java
{
  "aggregate_id": "sub_8f2a",
  "snapshot_version": 500,
  "state": {
    "user_id": "usr_42",
    "creator_id": "crt_99",
    "tier": "premium",
    "status": "active",
    "renewed_count": 12
  }
}

On load: read latest snapshot, then replay only events after snapshot_version.

Quick recall

Everything you need if you only revisit this box.

  • CQRS splits write (commands) and read (queries) models for independent scaling and optimisation.
  • Event sourcing stores state changes as immutable events in an append-only log.
  • Projections build denormalised read models from the event stream.
  • Snapshots avoid replaying entire event history on aggregate load.
  • Use when read/write shapes differ significantly or audit trails are required — not for simple CRUD.

Test yourself

Answer these before moving on — recall is what makes it stick.