Siddhant Deval
Siddhant Deval
backend5 min read

CQRS & Event Sourcing: Scalable Read/Write Architecture for Billing Systems

CQRS separates the read model from the write model — allowing each to be independently optimized. Event Sourcing goes further, storing domain state as an immutable append-only log of Domain Events rather than mutable relational records. Together, they provide complete audit trails, time-travel debugging, and horizontal read scalability for high-throughput billing and subscription systems.

CQRS & Event Sourcing: Scalable Read/Write Architecture for Billing Systems

Domain objects are autonomous state machines with enforced invariants — they are not passive bags of data passed between controller functions. By Part 10, the billing engine's write side is solid: Orders are created through use cases, mutate through guarded Aggregate methods, and emit Domain Events. But reading is still done by loading the full Order Aggregate and projecting it in the use case — the same normalized domain model that enforces write-side invariants is also used to satisfy read queries. This creates a fundamental tension: the model optimized for enforcing invariants is rarely the model optimized for fast, flexible reads.

CQRS (Command Query Responsibility Segregation) resolves this by maintaining two entirely separate models: a write model (the Aggregate) that enforces invariants, and a read model (a denormalized projection) optimized for query performance. Event Sourcing extends this by storing the sequence of Domain Events as the source of truth rather than the current state — making the read model a continuously updated projection of that event stream.

This article implements the full CQRS split for the billing engine, builds the OrderSummaryProjector, and introduces the Outbox Pattern for durable at-least-once event delivery.

Architectural Note

Prerequisites: Part 5 (Domain Events) for the event infrastructure; Part 6 (Application Layer) for the read/write separation already seeded with GetOrderSummaryUseCase; Part 9 (DI Container) for wiring the projector as an event handler.


1. The Unified Model Problem

The GetOrderSummaryUseCase from Part 6 works but has a scaling problem:

TYPESCRIPT
// ❌ Current approach: load full Aggregate to satisfy a read query
async execute(query: GetOrderSummaryQuery): Promise<OrderSummaryView> {
  const order = await this.orders.findById(query.orderId);   // Loads Order + all LineItems
  const customer = await this.customers.findById(order!.customerId); // Second Aggregate load
  // Now project a view model from two loaded Aggregates
  return {
    orderId:      order!.id,
    status:       order!.status,
    customerName: customer?.name ?? 'Unknown',
    total:        order!.total.toString(),
    itemCount:    order!.items.length,
    placedAt:     order!.placedAt?.toISOString() ?? '',
    shippingAddress: /* ... */,
  };
}

Three problems at scale:

  1. Two Aggregate loads per read: every GET /orders/:id loads Order (with N LineItem joins) and Customer. For a dashboard showing 50 recent orders, that is 100 Aggregate loads + N×50 line item joins.
  2. The write model's joins are wrong for reads: the orders + order_items normalized schema is optimized for write-side integrity; the read-side wants a single denormalized row per order with customerName already embedded.
  3. Cross-Aggregate reads require a domain join: reading customerName from Customer for an Order list requires either a cross-table JOIN (coupling read to write schema) or loading the full Customer Aggregate per order.

CQRS solves this by building a separate order_summary_view table that is kept up-to-date by the Domain Event stream.


2. The CQRS Architecture

The write side enforces invariants through the Aggregate. The read side reads directly from the denormalized view — no Aggregate instantiation, no domain joins, a single SELECT statement. The projector keeps them synchronized via the Domain Event stream.


3. The Write Side: No Changes

The write side (Parts 3–10) requires zero modifications. Order.place() already raises OrderPlaced. PlaceOrderUseCase already dispatches it via IEventBus. CQRS is additive on the write side — the Aggregate and use cases are unchanged.


4. The Read Model: order_summary_view

The read model is a denormalized table with one row per order, containing everything a UI list or dashboard needs without any joins:

SQL
-- prisma/migrations/XXXXXXX_add_order_summary_view/migration.sql
CREATE TABLE order_summary_view (
  order_id        TEXT PRIMARY KEY,
  customer_id     TEXT NOT NULL,
  customer_name   TEXT NOT NULL,
  customer_email  TEXT NOT NULL,
  status          TEXT NOT NULL,
  total_cents     INTEGER NOT NULL,
  currency        TEXT NOT NULL,
  item_count      INTEGER NOT NULL,
  placed_at       TIMESTAMPTZ,
  shipped_at      TIMESTAMPTZ,
  cancelled_at    TIMESTAMPTZ,
  shipping_city   TEXT,
  shipping_country TEXT,
  updated_at      TIMESTAMPTZ NOT NULL DEFAULT NOW()
);

CREATE INDEX idx_osv_customer_id ON order_summary_view (customer_id);
CREATE INDEX idx_osv_status ON order_summary_view (status);
CREATE INDEX idx_osv_placed_at ON order_summary_view (placed_at DESC);

This table is write-only from the projector and read-only from query handlers. Application code never writes to it directly — only the projector does, in response to Domain Events.

In Prisma schema:

PRISMA
model OrderSummaryView {
  orderId         String   @id @map("order_id")
  customerId      String   @map("customer_id")
  customerName    String   @map("customer_name")
  customerEmail   String   @map("customer_email")
  status          String
  totalCents      Int      @map("total_cents")
  currency        String
  itemCount       Int      @map("item_count")
  placedAt        DateTime? @map("placed_at")
  shippedAt       DateTime? @map("shipped_at")
  cancelledAt     DateTime? @map("cancelled_at")
  shippingCity    String?  @map("shipping_city")
  shippingCountry String?  @map("shipping_country")
  updatedAt       DateTime @updatedAt @map("updated_at")

  @@map("order_summary_view")
}

5. The Projector: Keeping the Read Model Current

The projector is an event handler (IEventHandler<DomainEvent>) that reacts to Domain Events and upserts the read model. It is registered in the DI container alongside other order.placed handlers:

TYPESCRIPT
// src/infrastructure/projectors/OrderSummaryProjector.ts
import { injectable, inject } from 'inversify';
import { PrismaClient } from '@prisma/client';
import { IEventHandler } from '../../application/shared/IEventBus';
import { OrderPlaced } from '../../domain/order/events/OrderPlaced';
import { OrderCancelled } from '../../domain/order/events/OrderCancelled';
import { OrderShipped } from '../../domain/order/events/OrderShipped';
import { ICustomerRepository } from '../../domain/customer/ports/ICustomerRepository';
import { TYPES } from '../container/TYPES';

@injectable()
export class OrderSummaryProjector
  implements
    IEventHandler<OrderPlaced>,
    IEventHandler<OrderCancelled>,
    IEventHandler<OrderShipped> {

  constructor(
    @inject(TYPES.PrismaClient) private readonly prisma: PrismaClient,
    @inject(TYPES.ICustomerRepository) private readonly customers: ICustomerRepository,
  ) {}

  async handle(event: OrderPlaced | OrderCancelled | OrderShipped): Promise<void> {
    switch (event.eventType) {
      case 'order.placed':   return this.onOrderPlaced(event as OrderPlaced);
      case 'order.cancelled': return this.onOrderCancelled(event as OrderCancelled);
      case 'order.shipped':  return this.onOrderShipped(event as OrderShipped);
    }
  }

  private async onOrderPlaced(event: OrderPlaced): Promise<void> {
    const customer = await this.customers.findById(event.customerId);

    await this.prisma.orderSummaryView.upsert({
      where: { orderId: event.orderId },
      create: {
        orderId:        event.orderId,
        customerId:     event.customerId,
        customerName:   customer?.name ?? 'Unknown',
        customerEmail:  customer?.email.value ?? '',
        status:         'PLACED',
        totalCents:     event.orderTotal.amountCents,
        currency:       event.orderTotal.currency,
        itemCount:      event.itemCount,
        placedAt:       event.occurredAt,
        shippingCity:   event.shippingCity ?? null,
        shippingCountry: event.shippingCountry ?? null,
      },
      update: {
        status:    'PLACED',
        totalCents: event.orderTotal.amountCents,
        itemCount:  event.itemCount,
        placedAt:   event.occurredAt,
        updatedAt:  new Date(),
      },
    });
  }

  private async onOrderCancelled(event: OrderCancelled): Promise<void> {
    await this.prisma.orderSummaryView.update({
      where: { orderId: event.orderId },
      data: {
        status:      'CANCELLED',
        cancelledAt: event.occurredAt,
        updatedAt:   new Date(),
      },
    });
  }

  private async onOrderShipped(event: OrderShipped): Promise<void> {
    await this.prisma.orderSummaryView.update({
      where: { orderId: event.orderId },
      data: {
        status:    'SHIPPED',
        shippedAt: event.occurredAt,
        updatedAt: new Date(),
      },
    });
  }
}

5.1 The Projector Is Idempotent by Design

The upsert on OrderPlaced means replaying the event multiple times is safe — the read model converges to the same state regardless of how many times the event is processed. This is critical for the Outbox Pattern (§7).


6. The Read Side: Direct SQL Queries

With the read model in place, GetOrderSummaryUseCase becomes a direct SQL SELECT with zero Aggregate instantiation:

TYPESCRIPT
// src/application/order/GetOrderSummaryUseCase.ts (CQRS version)
import { injectable, inject } from 'inversify';
import { PrismaClient } from '@prisma/client';
import { TYPES } from '../../infrastructure/container/TYPES';
import { GetOrderSummaryQuery, OrderSummaryView } from './queries/GetOrderSummaryQuery';
import { DomainError } from '../../domain/shared/DomainError';

@injectable()
export class GetOrderSummaryUseCase {
  constructor(
    @inject(TYPES.PrismaClient) private readonly prisma: PrismaClient,
  ) {}

  async execute(query: GetOrderSummaryQuery): Promise<OrderSummaryView> {
    // ✅ Single SELECT — no Aggregate load, no domain join, no N+1
    const row = await this.prisma.orderSummaryView.findUnique({
      where: { orderId: query.orderId },
    });

    if (!row) throw new DomainError(`Order ${query.orderId} not found`);

    return {
      orderId:      row.orderId,
      status:       row.status,
      customerName: row.customerName,
      total:        `${row.currency} ${(row.totalCents / 100).toFixed(2)}`,
      itemCount:    row.itemCount,
      placedAt:     row.placedAt?.toISOString() ?? '',
      shippingAddress: row.shippingCity
        ? { city: row.shippingCity, countryCode: row.shippingCountry! }
        : null,
    };
  }
}
TYPESCRIPT
// src/application/order/ListCustomerOrdersUseCase.ts — trivially fast with read model
@injectable()
export class ListCustomerOrdersUseCase {
  constructor(@inject(TYPES.PrismaClient) private readonly prisma: PrismaClient) {}

  async execute(customerId: CustomerId, page = 1, limit = 20): Promise<OrderSummaryView[]> {
    const rows = await this.prisma.orderSummaryView.findMany({
      where: { customerId },
      orderBy: { placedAt: 'desc' },
      skip: (page - 1) * limit,
      take: limit,
    });

    return rows.map(row => ({
      orderId:      row.orderId,
      status:       row.status,
      customerName: row.customerName,
      total:        `${row.currency} ${(row.totalCents / 100).toFixed(2)}`,
      itemCount:    row.itemCount,
      placedAt:     row.placedAt?.toISOString() ?? '',
      shippingAddress: null,
    }));
  }
}

A 50-order dashboard page is now one SELECT ... LIMIT 50 against a single indexed table. Response time drops from O(N × Aggregate load) to O(1 indexed SELECT).

Crucial Requirement

The read side uses PrismaClient directly — it does not go through IOrderRepository. This is intentional and correct: the read model is infrastructure, not a domain concern. The GetOrderSummaryUseCase is the one legitimate place in the Application layer where a direct infrastructure dependency is acceptable, because it is a pure read with no domain invariants to enforce.


7. The Outbox Pattern: Durable At-Least-Once Delivery

The post-commit dispatch from Part 5 has a gap: if the server crashes between repository.save() (successful DB commit) and eventBus.publish(), the events are lost. For a billing system, "lost OrderPlaced event" means the inventory is never reserved and the confirmation email is never sent.

The Outbox Pattern closes this gap by writing events to an outbox table in the same database transaction as the Aggregate save. A separate background worker reads unpublished events and dispatches them:

7.1 The Outbox Table

PRISMA
model OutboxEvent {
  id          String   @id @default(uuid())
  eventType   String   @map("event_type")
  aggregateId String   @map("aggregate_id")
  payload     Json
  occurredAt  DateTime @map("occurred_at")
  publishedAt DateTime? @map("published_at")  // null = unpublished
  retryCount  Int      @default(0) @map("retry_count")

  @@index([publishedAt])
  @@map("outbox_events")
}

7.2 Modified PrismaOrderRepository.save(): Atomic Write + Outbox

TYPESCRIPT
// src/infrastructure/persistence/PrismaOrderRepository.ts
async save(order: Order): Promise<void> {
  const events = order.pullDomainEvents(); // Pull events BEFORE the transaction

  await this.prisma.$transaction(async (tx) => {
    // 1. Save the Aggregate (same as before)
    await tx.order.upsert({
      where: { id: order.id },
      create: { id: order.id, customerId: order.customerId, status: order.status },
      update: { status: order.status, version: { increment: 1 } },
    });

    await tx.orderItem.deleteMany({ where: { orderId: order.id } });
    await tx.orderItem.createMany({
      data: order.items.map(item => ({ /* ... */ })),
    });

    // 2. Write events to outbox IN THE SAME TRANSACTION — atomic
    if (events.length > 0) {
      await tx.outboxEvent.createMany({
        data: events.map(event => ({
          id:          event.eventId,
          eventType:   event.eventType,
          aggregateId: event.aggregateId,
          payload:     JSON.stringify(event),
          occurredAt:  event.occurredAt,
          publishedAt: null, // Unpublished
        })),
      });
    }
  });
  // If the transaction commits, events are in the outbox.
  // If the transaction rolls back, events are NOT in the outbox.
  // Either way: consistent.
}

7.3 The Outbox Poller

A background worker (run as a separate process or a setInterval in the same process) polls for unpublished events and dispatches them:

TYPESCRIPT
// src/infrastructure/messaging/OutboxPublisher.ts
@injectable()
export class OutboxPublisher {
  private readonly BATCH_SIZE = 50;
  private readonly POLL_INTERVAL_MS = 1000;

  constructor(
    @inject(TYPES.PrismaClient) private readonly prisma: PrismaClient,
    @inject(TYPES.IEventBus) private readonly eventBus: IEventBus,
  ) {}

  start(): void {
    const poll = async () => {
      try {
        await this.processBatch();
      } catch (err) {
        console.error('Outbox poll error:', err);
      } finally {
        setTimeout(poll, this.POLL_INTERVAL_MS);
      }
    };
    poll();
  }

  private async processBatch(): Promise<void> {
    const unpublished = await this.prisma.outboxEvent.findMany({
      where: { publishedAt: null, retryCount: { lt: 5 } },
      orderBy: { occurredAt: 'asc' },
      take: this.BATCH_SIZE,
    });

    for (const record of unpublished) {
      try {
        // Deserialize the event and publish it
        const event = this.deserialize(record);
        await this.eventBus.publish(event);

        // Mark as published — idempotent: re-processing is safe because handlers are idempotent
        await this.prisma.outboxEvent.update({
          where: { id: record.id },
          data: { publishedAt: new Date() },
        });
      } catch (err) {
        await this.prisma.outboxEvent.update({
          where: { id: record.id },
          data: { retryCount: { increment: 1 } },
        });
      }
    }
  }

  private deserialize(record: OutboxEventRecord): DomainEvent {
    const payload = JSON.parse(record.payload as string);
    // Reconstruct the correct event type based on eventType discriminator
    switch (record.eventType) {
      case 'order.placed':    return Object.assign(new OrderPlaced('', '' as any, Money.of('USD', 0), 0), payload);
      case 'order.cancelled': return Object.assign(new OrderCancelled('', ''), payload);
      case 'order.shipped':   return Object.assign(new OrderShipped('', '' as any), payload);
      default: throw new Error(`Unknown event type: ${record.eventType}`);
    }
  }
}
Pro Tip & Optimization

At-least-once vs. exactly-once: The Outbox Pattern guarantees at-least-once delivery — an event may be dispatched more than once if the poller crashes after publishing but before marking publishedAt. This is why projectors must be idempotent (upsert on order_summary_view). Design all event handlers to be safely re-runnable with the same event payload.


8. Event Versioning: Evolving the Event Schema

As the billing engine matures, OrderPlaced v1 may need a new field (promotionCode). Events already stored in the outbox or Kafka are v1. The projector must handle both versions:

TYPESCRIPT
// src/domain/order/events/OrderPlaced.ts
export class OrderPlaced extends DomainEvent {
  constructor(
    public readonly orderId: OrderId,
    public readonly customerId: CustomerId,
    public readonly orderTotal: Money,
    public readonly itemCount: number,
    public readonly promotionCode: string | null = null, // ← v2 field, nullable for v1 compat
  ) {
    super(orderId, 2); // ← version bumped to 2
  }

  get eventType() { return 'order.placed' as const; }
}

The projector's onOrderPlaced handles both via eventVersion:

TYPESCRIPT
private async onOrderPlaced(event: OrderPlaced): Promise<void> {
  const promoCode = event.eventVersion >= 2 ? event.promotionCode : null;
  await this.prisma.orderSummaryView.upsert({ /* ... promoCode included if v2 */ });
}

Summary

Concept Domain Rule
CQRS Write model enforces invariants via Aggregates; read model is a denormalized projection optimized for queries
Projector An event handler that upserts the read model; must be idempotent; lives in Infrastructure
Read model A dedicated table/view with pre-joined data; read handlers use direct SQL, never Aggregate loading
Outbox Pattern Events written to an outbox table in the same DB transaction as the Aggregate save; eliminates the post-commit crash gap
At-least-once delivery Outbox polling can dispatch an event more than once; all handlers must be idempotent
Event versioning New fields are nullable with defaults; eventVersion discriminates v1 vs. v2 handlers
GetOrderSummaryUseCase Direct PrismaClient query on order_summary_view — zero Aggregate load, O(1) indexed SELECT

What's Next

In Part 12, we build the full testing strategy — defining the test pyramid for this architecture, implementing contract tests for all Ports, and assembling the fast (~6s) full test suite that validates 237 behaviors across domain, application, and infrastructure layers. Part 12: Testing Strategy →

Research & Synthesis Note

This article was developed with AI-assisted deep search, specification cross-referencing, and technical research synthesis.

#CQRS#Event Sourcing#Domain-Driven Design#TypeScript#Node.js#Architecture#Event-Driven
Siddhant Deval

Written by Siddhant Deval

Senior Full-Stack Engineer building high-scale architectures, browser performance engineering systems, and SaaS platforms.