StackPractices
advanced Por Mathias Paulenko

Implementar Event Sourcing en Arquitecturas Serverless

Cómo capturar todos los cambios como eventos inmutables usando event sourcing con AWS Lambda, DynamoDB streams y event stores para audit trails y consultas temporales.

Visión general

Los sistemas tradicionales almacenan el estado actual. Una orden está “enviada,” y la fila de base de datos dice status = shipped. Si un usuario pregunta “¿cuándo cambió el estado a enviado?” la base de datos no tiene respuesta — el valor anterior fue sobrescrito. Si un analista pregunta “¿cuántas órdenes fueron canceladas y re-enviadas el mes pasado?” el sistema no puede responder sin agregar columnas de auditoría explícitas que rastreen cada cambio manualmente.

Event sourcing almacena cada cambio de estado como un evento inmutable en un log append-only. El estado actual se computa reproduciendo eventos. El estado de una orden no es una fila — es la secuencia [OrderCreated, ItemAdded, PaymentProcessed, Shipped]. Esto provee un audit trail completo, soporta consultas temporales (“¿cuál era el estado a las 3pm de ayer?”), y permite reconstruir proyecciones desde cero. En arquitecturas serverless, los eventos se capturan vía DynamoDB streams, SQS o EventBridge, y las funciones Lambda proyectan el modelo de lectura. La solucion a continuacion cubre implementación de event sourcing, event stores, proyecciones y consideraciones específicas de serverless.

Cuándo usarlo

Usa esta receta cuando:

  • El historial completo de auditoría de todos los cambios es un requerimiento de negocio. Consulta Event-Driven Functions para arquitecturas event-driven.
  • Necesitas responder preguntas temporales sobre estados pasados
  • Reconstruir modelos de lectura desde cero es una capacidad necesaria. Consulta Serverless Orchestration para gestionar workflows stateful.
  • El modelo de escritura es complejo y el modelo de lectura necesita optimizarse separadamente
  • Requerimientos de compliance o regulatorios mandatan logs de cambio inmutables. Consulta CQRS Pattern para separar modelos de lectura y escritura.

Solución

Event Store con DynamoDB y Streams

interface DomainEvent {
  eventId: string;
  aggregateId: string;
  eventType: string;
  payload: Record<string, unknown>;
  timestamp: string;
  version: number;
}

class OrderEventStore {
  constructor(private tableName: string, private client: DynamoDBDocument) {}

  async appendEvents(aggregateId: string, events: DomainEvent[]): Promise<void> {
    const currentVersion = await this.getCurrentVersion(aggregateId);

    const transactItems = events.map((event, index) => ({
      Put: {
        TableName: this.tableName,
        Item: {
          pk: `ORDER#${aggregateId}`,
          sk: `EVENT#${(currentVersion + index + 1).toString().padStart(10, '0')}`,
          eventId: event.eventId,
          eventType: event.eventType,
          payload: event.payload,
          timestamp: new Date().toISOString(),
          version: currentVersion + index + 1,
        },
        ConditionExpression: 'attribute_not_exists(pk)',
      },
    }));

    await this.client.transactWrite({ TransactItems: transactItems });
  }

  async getEvents(aggregateId: string): Promise<DomainEvent[]> {
    const result = await this.client.query({
      TableName: this.tableName,
      KeyConditionExpression: 'pk = :pk AND begins_with(sk, :sk)',
      ExpressionAttributeValues: {
        ':pk': `ORDER#${aggregateId}`,
        ':sk': 'EVENT#',
      },
      ScanIndexForward: true,
    });

    return (result.Items || []).map(item => ({
      eventId: item.eventId,
      aggregateId,
      eventType: item.eventType,
      payload: item.payload,
      timestamp: item.timestamp,
      version: item.version,
    }));
  }

  private async getCurrentVersion(aggregateId: string): Promise<number> {
    const events = await this.getEvents(aggregateId);
    return events.length > 0 ? events[events.length - 1].version : 0;
  }
}

Lambda Projection Handler

export const handler = async (event: DynamoDBStreamEvent): Promise<void> => {
  for (const record of event.Records) {
    if (record.eventName !== 'INSERT') continue;

    const newImage = unmarshall(record.dynamodb?.NewImage as any);
    const domainEvent: DomainEvent = {
      eventId: newImage.eventId,
      aggregateId: newImage.aggregateId,
      eventType: newImage.eventType,
      payload: newImage.payload,
      timestamp: newImage.timestamp,
      version: newImage.version,
    };

    await projectEvent(domainEvent);
  }
};

async function projectEvent(event: DomainEvent): Promise<void> {
  switch (event.eventType) {
    case 'OrderCreated':
      await createOrderProjection(event.aggregateId, event.payload);
      break;
    case 'ItemAdded':
      await addItemToOrderProjection(event.aggregateId, event.payload);
      break;
    case 'OrderShipped':
      await updateOrderStatus(event.aggregateId, 'shipped');
      break;
  }
}

Reconstrucción de Agregado

class OrderAggregate {
  private status: string = 'pending';
  private items: OrderItem[] = [];
  private total: number = 0;

  applyEvent(event: DomainEvent): void {
    switch (event.eventType) {
      case 'OrderCreated':
        this.status = 'created';
        this.total = event.payload.total as number;
        break;
      case 'ItemAdded':
        this.items.push(event.payload.item as OrderItem);
        this.total += (event.payload.item as OrderItem).price;
        break;
      case 'OrderShipped':
        this.status = 'shipped';
        break;
      case 'OrderCancelled':
        this.status = 'cancelled';
        break;
    }
  }

  static fromEvents(events: DomainEvent[]): OrderAggregate {
    const order = new OrderAggregate();
    for (const event of events) {
      order.applyEvent(event);
    }
    return order;
  }
}

Explicación

  • Event store: el event store es un log append-only. Los eventos nunca se actualizan ni eliminan. Cada evento tiene un ID único, un aggregate ID (la entidad a la que pertenece), un tipo, un payload y una versión.
  • Reconstrucción de agregado: En su lugar, cargas todos los eventos de un agregado y los reproduces en orden. El objeto agregado comienza vacío y aplica cada evento, mutando su estado interno. Esto es determinista — la misma secuencia de eventos siempre produce el mismo estado.
  • Proyecciones (modelos de lectura): los modelos de lectura se construyen suscribiéndose al stream de eventos. Puedes tener múltiples proyecciones para los mismos eventos — una para el dashboard del cliente, otra para analytics, otra para indexación de búsqueda.
  • Snapshots: reproducir miles de eventos para un agregado de larga vida es lento. Los snapshots cachean el estado del agregado en una versión específica. Para reconstruir, carga el último snapshot y reproduce solo eventos después de esa versión. cada 100 eventos) y asíncronamente.

Variantes

EnfoqueStoreProyeccionesMejor para
DynamoDB + StreamsDynamoDBLambdaNativo AWS, escala moderada
EventStoreDBEventStoreDBSubscriptionsAlto volumen, dominios complejos
Kafka + KTablesKafkaKafka StreamsStream processing, replay
S3 + AthenaS3Athena queriesAudit, compliance, analytics
Aurora + OutboxPostgreSQLCDCEvent sourcing relacional

Lo que funciona

  • Versiona cada evento: Esto previene updates perdidos cuando dos usuarios modifican simultáneamente el mismo agregado.
  • Haz los eventos inmutables y autocontenidos: un evento debería llevar todos los datos necesarios para entenderlo, no solo deltas. ” Consumidores futuros no deberían necesitar consultar otros sistemas para interpretar el evento.
  • Usa correlation IDs a través de la cadena de eventos: cuando un evento dispara otro (ej. OrderShipped dispara InventoryDecremented), propaga el correlation ID. Esto habilita tracing end-to-end y debugging a través de cadenas de eventos distribuidas.
  • Implementa proyecciones idempotentes: las funciones Lambda reintentan ante fallas. Una proyección que incrementa un contador en cada invocación sobrecuentará. Diseña proyecciones idempotentes — escribe el event ID en la fila de proyección y salta si ya fue procesado.
  • Archiva eventos viejos a cold storage: DynamoDB es caro para almacenamiento a largo plazo de millones de eventos. Mueve eventos mayores a 90 días a S3 usando TTL de DynamoDB o jobs de export.

Errores comunes

  • Almacenar estado actual junto a eventos: si mantienes tanto un log de eventos como una tabla de estado actual, pueden divergir. Un bug en la proyección escribe estado A mientras el log contiene eventos para estado B. La fuente de verdad es el event store; las proyecciones son derivadas. No trates la proyección como estado primario.
  • Exponer tipos de evento a sistemas externos: los consumidores externos no deberían depender de schemas internos de eventos. OrderConfirmed) y mapea eventos internos a públicos. El refactoring interno de tipos de evento no debería romper integraciones externas.
  • No manejar la evolución de schema de eventos: cuando un tipo de evento cambia (agregando un campo), eventos viejos en el log no tienen el nuevo campo.
  • Reproducir eventos desde el inicio para cada query: siempre usa snapshots para agregados con historias largas. Reproducir 10,000 eventos para cada GET /order/123 destruye el rendimiento. Toma snapshots asíncronamente y carga desde ellos.

Lectura Adicional

  • Documentación oficial: consulta la referencia actualizada del framework o herramienta utilizada.
  • Guías relacionadas: explora las guías de serverless y cqrs para profundizar.
  • Patrones complementarios: revisa los patrones de diseño aplicables a tu stack tecnológico.
  • Postmortems públicos: estudia incidentes reales de equipos que enfrentaron problemas similares en producción.

Notas de Producción

  • Despliega gradualmente usando canary o blue-green para detectar regresiones temprano.
  • Configura alertas para errores, latencia p99 y tasa de fallos antes de habilitar en producción.
  • Documenta el rollback en el runbook; prueba el procedimiento en staging al menos una vez por trimestre.
  • Revisa logs estructurados con correlation IDs para trazar requests end-to-end en incidentes.

Puntos Clave

  • Aplica implementar event sourcing en arquitecturas serverless cuando necesites una solución práctica para tu caso de uso.
  • Monitorea el rendimiento después de implementar; mide latencia, errores y uso de recursos antes y después.
  • Revisa la sección de Troubleshooting ante errores comunes; la mayoría tienen causa raíz documentada con solución.
  • Mantén dependencias actualizadas y ejecuta tests en CI para prevenir regresiones en producción.

Troubleshooting

  • Cold start latency is high: increase provisioned concurrency, reduce package size, and avoid initializing heavy clients per invocation.
  • Function times out: check downstream dependencies, memory allocation, and retry logic. Increase timeout only after optimizing the code.
  • State lost between invocations: serverless functions are stateless. Persist state in a database, cache, or durable queue.
  • Deployment package too large: exclude dev dependencies and unused assets.
  • Event ordering issues: many event sources are at-least-once and unordered. Design for idempotency and explicit sequencing.

Errores Comunes en Producción

  • Copiar el ejemplo sin adaptarlo a volúmenes y modos de fallo reales.
  • Saltar tests de carga e inyección de errores antes del primer despliegue productivo.
  • Codificar valores fijos que deberían ser configurables por entorno.
  • Olvidar agregar logging y monitoreo en cada paso.
  • Desplegar sin plan de rollback ni estrategia de backup probada.
  • Asumir que el ejemplo mínimo escalará sin agregar caché o procesamiento por lotes.
  • No documentar la versión y configuración usadas en producción.
  • Dejar la receta sin cambios cuando evolucionan las dependencias o la escala.

Preguntas frecuentes

¿Es event sourcing más complejo que CRUD?

Sí. Agrega conceptos (agregados, proyecciones, versionado de eventos) e infraestructura (event stores, stream processors). Úsalo solo cuando los beneficios (auditoría, consultas temporales, capacidad de reconstrucción) justifiquen la complejidad. Para CRUD simple sin requerimientos de auditoría, el almacenamiento tradicional de estado es suficiente.

¿Cómo elimino datos bajo GDPR si los eventos son inmutables?

Implementa crypto-shredding: encripta payloads de eventos con una clave por usuario. Para "eliminar" los datos de un usuario, borra su clave de encriptación. Los eventos permanecen pero son ilegibles. Alternativamente, almacena PII en un store mutable separado y referéncialo desde los eventos.

¿Puedo usar event sourcing con bases de datos relacionales?

Sí — usa el outbox pattern. Escribe eventos a una tabla outbox en la misma transacción que los cambios de datos de negocio. Un proceso CDC (change data capture) sondea el outbox y publica eventos. Esto te da garantías ACID con semántica de event sourcing.

¿Cómo consulto a través de agregados?

No consultes el event store directamente para queries cross-aggregate. Construye proyecciones de modelo de lectura que desnormalicen datos para eficiencia de query. El event store es el modelo de escritura; las proyecciones son el modelo de lectura. Esta separación es CQRS.

¿Cómo manejo la evolución del schema de eventos?

Versiona eventos explícitamente: incluye un campo version en cada evento. Usa upcasters (transformadores que convierten versiones viejas de eventos a nuevas) al leer eventos del store. Nunca modifiques clases de evento existentes — crea una nueva versión y escribe un upcaster. Para protobuf, usa campos reserved y agrega nuevos campos con nuevos números. Para eventos JSON, usa evolución de json-schema con cambios additive-only.

¿Cómo manejo eventos duplicados en serverless?

Usa idempotency keys: incluye un event ID único (UUID) y trackea IDs procesados en una tabla de deduplicación. En AWS Lambda, usa DynamoDB conditional writes para marcar atómicamente un evento como procesado. Setea un TTL en la tabla de deduplicación (ej., 7 días) para limitar storage. Para Kinesis, usa el sequence number como key de deduplicación. Procesa eventos idempotentemente para que reprocesar el mismo evento produzca el mismo resultado.

¿Cómo reproceso eventos para reconstruir read models?

Lee todos los eventos del event store en orden, aplica cada uno al handler de proyección, y escribe el read model actualizado. Usa una tabla de checkpoint para trackear el último sequence number procesado. Para event stores grandes, reprocesá en batches (ej., 1000 eventos a la vez) para evitar memory pressure. Corre el replay como una Lambda function separada o batch job. Pausa el handler de proyección real-time durante el replay para evitar conflictos, luego resume desde el checkpoint.

¿Cómo testeo sistemas de event sourcing?

Testea agregados reproduciendo eventos y asertando sobre el estado resultante. Testea proyecciones alimentando una secuencia de eventos conocida y asertando sobre el output del read model. Usa event fixtures: una lista de eventos que producen un estado de agregado conocido. Para tests de integración, usa un event store in-memory y verifica el ciclo completo: command → events → projection. Testea versionado de eventos reproduciendo eventos de versión vieja a través de upcasters y asertando que el payload upcasted coincide con el nuevo schema.

¿Cómo manejo writes concurrentes al mismo agregado?

Usa optimistic concurrency control. Incluye el número de versión esperado en el write request. El event store rechaza el write si la versión actual no coincide con la versión esperada. En DynamoDB, usa una conditional expression: attribute_not_exists(version) OR version = :expected_version. En conflicto, reintenta cargando los últimos eventos, reaplicando el command, y escribiendo de nuevo. Para agregados de alta contención, considera usar un saga o process manager para serializar writes. No uses pessimistic locking en serverless — las Lambda functions son stateless y no pueden mantener locks.

¿Cómo implemento snapshots para agregados con historias de eventos largas?

Periódicamente guarda el estado completo del agregado como snapshot. Almacena snapshots en una tabla separada con el aggregate ID y número de versión. Al cargar, fetchea el último snapshot y reproduce solo los eventos después de la versión del snapshot. Toma snapshots cada N eventos (ej., cada 100) o después de un intervalo de tiempo. En DynamoDB, almacena snapshots en una partición separada: PK = AGGREGATE#123, SK = SNAPSHOT#42. La creación de snapshots debería ser async — no bloquees el write path. Si un snapshot falla, el sistema continúa funcionando reproduciendo desde el principio.