StackPractices
advanced Por Mathias Paulenko

Escalar Cargas de Lectura y Escritura con CQRS

Cómo separar modelos de lectura y escritura usando Command Query Responsibility Segregation para queries optimizadas, event sourcing, y escalado independiente de rutas de lectura y escritura.

Temas: design

Nota para desarrolladores hispanohablantes: Esta guía incluye ejemplos y convenciones de nomenclatura adaptadas a equipos que trabajan en español. Cuando existen diferencias significativas en terminología técnica entre el inglés y el español, se indican explícitamente para facilitar la comunicación en equipos multiculturales.

Visión general

Las aplicaciones CRUD tradicionales usan un único modelo de datos para lectura y escritura. Una tabla relacional sirve queries SELECT para dashboards y operaciones INSERT/UPDATE para envíos de forms. Esta simplicidad funciona para dominios pequeños pero se rompe a escala. Los esquemas optimizados para escritura (normalizados, transaccionales) son lentos para lecturas complejas. Los esquemas optimizados para lectura (desnormalizados, indexados) son costosos de actualizar. A medida que crece el tráfico, ambas cargas de trabajo compiten por los mismos recursos de base de datos.

Command Query Responsibility Segregation (CQRS) divide el modelo de datos en dos: un modelo de escritura optimizado para commands (crear, actualizar, eliminar) y un modelo de lectura optimizado para queries. Los commands mutan estado en el modelo de escritura y publican eventos. Los event handlers actualizan modelos de lectura — proyecciones desnormalizadas adaptadas para patrones de query específicos. Los dos modelos pueden usar diferentes bases de datos, diferentes esquemas, y escalar independientemente. La solucion a continuacion cubre implementación de CQRS con event sourcing y patrones de proyección.

Cuándo usarlo

Usa esta receta cuando:

  • Los volúmenes de lectura y escritura difieren considerablemente (ratios 100:1 de lectura intensiva son comunes). Consulta Réplicas de Lectura para escalado de lecturas.
  • Queries complejas requieren joins entre múltiples aggregates, degradando performance de escritura. Consulta SQL Joins para optimización de queries.
  • Construyendo dashboards en tiempo real, analytics, o búsqueda que necesita datos con forma diferente al modelo transaccional
  • Trabajando con sistemas event-sourced donde el estado se deriva de una secuencia de eventos. Consulta Event Sourcing para patrones de event store.
  • Equipos necesitan optimizar esquemas de lectura y escritura independientemente sin coordinación

Solución

Command Handler con Publicación de Eventos (TypeScript)

interface Event {
  type: string;
  aggregateId: string;
  payload: unknown;
  occurredAt: Date;
}

interface EventStore {
  append(events: Event[]): Promise<void>;
  getEvents(aggregateId: string): Promise<Event[]>;
}

class OrderWriteModel {
  constructor(
    public id: string,
    public customerId: string,
    public items: Array<{ productId: string; quantity: number; price: number }>,
    public status: 'pending' | 'paid' | 'shipped' = 'pending'
  ) {}

  pay(paymentMethod: string): Event[] {
    if (this.status !== 'pending') throw new Error('Order already paid');
    return [{
      type: 'OrderPaid',
      aggregateId: this.id,
      payload: { paymentMethod, total: this.total() },
      occurredAt: new Date(),
    }];
  }

  total(): number {
    return this.items.reduce((sum, item) => sum + item.price * item.quantity, 0);
  }
}

class PayOrderCommand {
  constructor(public orderId: string, public paymentMethod: string) {}
}

class PayOrderHandler {
  constructor(private eventStore: EventStore) {}

  async handle(command: PayOrderCommand): Promise<void> {
    const events = await this.eventStore.getEvents(command.orderId);
    const order = this.rehydrate(events);
    const newEvents = order.pay(command.paymentMethod);
    await this.eventStore.append(newEvents);
  }

  private rehydrate(events: Event[]): OrderWriteModel {
    const order = new OrderWriteModel(events[0].aggregateId, '', []);
    for (const event of events) {
      // Apply each event to mutate order state
    }
    return order;
  }
}

Read Model Projection (SQL)

CREATE TABLE order_summaries (
  order_id UUID PRIMARY KEY,
  customer_name VARCHAR(255),
  total_amount DECIMAL(10,2),
  item_count INT,
  status VARCHAR(20),
  last_updated TIMESTAMP
);

CREATE OR REPLACE FUNCTION project_order_paid()
RETURNS TRIGGER AS $$
BEGIN
  UPDATE order_summaries
  SET status = 'paid', last_updated = NEW.occurred_at
  WHERE order_id = NEW.aggregate_id;
  RETURN NEW;
END;
$$ LANGUAGE plpgsql;

Event Handler Actualizando Read Model (Node.js)

class OrderProjection {
  constructor(private readDb, private elasticsearch) {}

  async handleOrderPaid(event) {
    await this.readDb.query(
      'UPDATE order_summaries SET status = $1, last_updated = $2 WHERE order_id = $3',
      ['paid', event.occurredAt, event.aggregateId]
    );

    await this.elasticsearch.update({
      index: 'orders',
      id: event.aggregateId,
      doc: { status: 'paid', paidAt: event.occurredAt }
    });
  }

  async handleOrderCreated(event) {
    await this.readDb.query(
      `INSERT INTO order_summaries (order_id, customer_name, total_amount, item_count, status, last_updated)
       VALUES ($1, $2, $3, $4, $5, $6)`,
      [event.aggregateId, event.payload.customerName, event.payload.total,
       event.payload.items.length, 'pending', event.occurredAt]
    );
  }
}

Explicación

  • Command model: Cada command es validado contra invariantes, muta el modelo de escritura, y produce domain events. El modelo de escritura es normalizado y transaccional — enforce reglas de negocio a costa de complejidad de query.
  • Event sourcing: El modelo de escritura agrega eventos a un event store. El estado se rehidrata reproduciendo eventos. Esto provee historial de auditoría completo y queries temporales.
  • Read model (proyección): una vista desnormalizada y optimizada para queries construida desde eventos. Un read model de customer_orders podría aplanar items de orden, nombres de clientes y estado de envío en una sola tabla con índices apropiados.
  • Consistencia eventual: Hay una breve ventana (milisegundos a segundos) donde el modelo de escritura refleja el cambio pero el read model no. Esto es consistencia eventual — aceptable para la mayoría de sistemas de lectura intensiva.

Variantes

EnfoqueModelo de escrituraModelo de lecturaConsistenciaMejor para
Single DB, vistas separadasRelacionalMaterialized viewsFuerteCQRS simple
Dual DBRelacionalDocumento/BúsquedaEventualAlta escala de lectura
Event sourcingEvent storeMúltiples proyeccionesEventualAuditoría, queries temporales
Réplicas de lecturaDB primariaDB réplicaCasi-fuerteEscalado de lectura sin complejidad

Lo que funciona

  • Mantén read models simples y desechables: un read model es un cache, no una fuente de verdad. Si se corrompe, reconstrúyelo reproduciendo eventos desde el inicio. No pongas lógica de negocio u operaciones de escritura en read models.
  • Versiona tus eventos: a medida que los esquemas evolucionan, proyecciones más antiguas aún deben entender eventos históricos. Esto permite migración gradual sin downtime.
  • Usa proyecciones idempotentes: los event handlers pueden ejecutarse múltiples veces (entrega al-menos-una-vez). Diseña proyecciones para que procesar el mismo evento dos veces produzca el mismo resultado.
  • Monitorea lag de proyección: el delay entre escritura y actualización de read model debe estar acotado. Alerta si el lag de proyección excede tu SLA (ej. 5 segundos). Proyecciones lentas indican backpressure o event handlers ineficientes.
  • Empieza simple, evoluciona a CQRS: no construyas CQRS desde el día uno en un proyecto greenfield. Empieza con un modelo único. Cuando la complejidad o performance de lectura se vuelve un problema, extrae un read model. CQRS prematuro agrega complejidad innecesaria.

Errores comunes

  • CQRS sin razón: si tu aplicación tiene CRUD simple con ratios iguales de lectura/escritura, CQRS agrega complejidad sin beneficio. Úsalo cuando la asimetría de lectura/escritura o complejidad de query justifique la separación.
  • Poner reglas de negocio en read models: los read models son para querying. Si te encuentras validando o mutando estado en una proyección, has violado la separación. Las reglas de negocio pertenecen a los command handlers.
  • Ignorar consistencia eventual en UX: los usuarios pueden enviar un form e inmediatamente refrescar, viendo datos obsoletos.
  • Reproducir eventos desde el inicio en cada deploy: En producción con billones de eventos, esto toma días.

Preguntas frecuentes

Implementación en Python con Event Sourcing
from dataclasses import dataclass, field
from datetime import datetime
from typing import List, Any
from abc import ABC, abstractmethod
import uuid

@dataclass
class Event:
    type: str
    aggregate_id: str
    payload: Any
    occurred_at: datetime = field(default_factory=datetime.utcnow)

class EventStore(ABC):
    @abstractmethod
    async def append(self, events: List[Event]) -> None:
        ...

    @abstractmethod
    async def get_events(self, aggregate_id: str) -> List[Event]:
        ...

@dataclass
class OrderItem:
    product_id: str
    quantity: int
    price: float

class OrderWriteModel:
    def __init__(self, order_id: str, customer_id: str, items: List[OrderItem]):
        self.id = order_id
        self.customer_id = customer_id
        self.items = items
        self.status = 'pending'

    def pay(self, payment_method: str) -> List[Event]:
        if self.status != 'pending':
            raise ValueError('Order already paid')
        self.status = 'paid'
        return [Event(
            type='OrderPaid',
            aggregate_id=self.id,
            payload={'payment_method': payment_method, 'total': self.total()}
        )]

    def total(self) -> float:
        return sum(item.price * item.quantity for item in self.items)

class PayOrderHandler:
    def __init__(self, event_store: EventStore):
        self._event_store = event_store

    async def handle(self, order_id: str, payment_method: str) -> None:
        events = await self._event_store.get_events(order_id)
        order = self._rehydrate(events)
        new_events = order.pay(payment_method)
        await self._event_store.append(new_events)

    def _rehydrate(self, events: List[Event]) -> OrderWriteModel:
        if not events:
            raise ValueError('No events found for order')
        order = OrderWriteModel(events[0].aggregate_id, '', [])
        for event in events:
            if event.type == 'OrderCreated':
                order.customer_id = event.payload['customer_id']
                order.items = [OrderItem(**i) for i in event.payload['items']]
            elif event.type == 'OrderPaid':
                order.status = 'paid'
        return order
Snapshotting para Event Streams Grandes
interface Snapshot {
  aggregateId: string;
  version: number;
  state: unknown;
}

class SnapshotStore {
  async save(snapshot: Snapshot): Promise<void> {
    // Persistir snapshot a base de datos
  }

  async load(aggregateId: string): Promise<Snapshot | null> {
    // Cargar último snapshot
    return null;
  }
}

class PayOrderHandlerWithSnapshots {
  constructor(
    private eventStore: EventStore,
    private snapshots: SnapshotStore
  ) {}

  async handle(command: PayOrderCommand): Promise<void> {
    const snapshot = await this.snapshots.load(command.orderId);
    let order: OrderWriteModel;
    let fromVersion = 0;

    if (snapshot) {
      order = snapshot.state as OrderWriteModel;
      fromVersion = snapshot.version;
    } else {
      const events = await this.eventStore.getEvents(command.orderId);
      order = this.rehydrate(events);
      fromVersion = events.length;
    }

    // Aplicar solo eventos después del snapshot
    const recentEvents = await this.eventStore.getEventsAfter(
      command.orderId, fromVersion
    );
    for (const event of recentEvents) {
      this.apply(order, event);
    }

    const newEvents = order.pay(command.paymentMethod);
    await this.eventStore.append(newEvents);

    // Guardar snapshot cada 100 eventos
    if (fromVersion + recentEvents.length + newEvents.length >= 100) {
      await this.snapshots.save({
        aggregateId: command.orderId,
        version: fromVersion + recentEvents.length + newEvents.length,
        state: order,
      });
    }
  }

  private apply(order: OrderWriteModel, event: Event): void {
    // Aplicar evento a la orden
  }

  private rehydrate(events: Event[]): OrderWriteModel {
    return new OrderWriteModel(events[0]?.aggregateId ?? '', '', []);
  }
}
Base de Datos de Lectura Separada con Redis
class OrderReadModelRedis {
  constructor(private redis: RedisClient) {}

  async getOrderSummary(orderId: string): Promise<OrderSummary | null> {
    const data = await this.redis.hgetall(`order:${orderId}:summary`);
    if (!data || Object.keys(data).length === 0) return null;

    return {
      orderId,
      customerName: data.customerName,
      totalAmount: parseFloat(data.totalAmount),
      itemCount: parseInt(data.itemCount, 10),
      status: data.status,
    };
  }

  async handleOrderCreated(event: Event): Promise<void> {
    await this.redis.hset(`order:${event.aggregateId}:summary`, {
      customerName: event.payload.customerName,
      totalAmount: event.payload.total.toString(),
      itemCount: event.payload.items.length.toString(),
      status: 'pending',
    });
  }

  async handleOrderPaid(event: Event): Promise<void> {
    await this.redis.hset(`order:${event.aggregateId}:summary`, {
      status: 'paid',
    });
  }
}