CultivaMarket — Proyecciones CQRS / Event Sourcing

Modelos de lectura optimizados para el marketplace B2B de servicios de IA

Python + asyncpg PostgreSQL Event Store Elasticsearch 8 4 Proyecciones Idempotente
Arquitectura de Proyecciones
🗄️
Event Store
PostgreSQL
stream_id event_type data jsonb global_position
eventos
⚙️
Projector
Orquestador
checkpoint batch_size retry logic monitoring
aplica
🔀
Proyecciones
4 handlers
order_summary product_search daily_sales customer_activity
escribe
📊
Read Models
Optimizados
PG Read DB Elasticsearch Redis Cache
Tipos de Proyección

Live

Suscripción en tiempo real al stream de eventos. Procesa eventos a medida que se publican.

CultivaMarket usa: order_summary, customer_activity

Catchup

Procesa eventos históricos desde una posición determinada. Ideal para reconstruir modelos.

CultivaMarket usa: daily_sales rebuild, nuevas proyecciones

Persistent

Almacena un checkpoint para continuar desde el último punto tras reinicios o fallos.

CultivaMarket usa: Todas — tabla projection_checkpoints

Inline

Ejecutada en la misma transacción que el write. Garantiza consistencia fuerte.

CultivaMarket usa: customer_activity (multi-tabla)
Las 4 Proyecciones de CultivaMarket
🛒

Order Summary

order_summary
OrderCreated OrderItemAdded OrderItemRemoved OrderShipped OrderCompleted OrderCancelled

Vista desnormalizada de pedidos con estado actual, importe total e item count. Permite filtrado rápido por agencia, estado y rango de fechas sin JOINs.

Target: tabla order_summaries en PostgreSQL Read DB
🔍

Product Search

product_search
ProductCreated ProductUpdated ProductPriceChanged ProductDeleted

Índice Elasticsearch de recursos digitales (prompts, plantillas, automatizaciones) con búsqueda full-text, filtros de categoría y rango de precios.

Target: índice cultivamarket-products en Elasticsearch 8
📈

Daily Sales

daily_sales
OrderCompleted OrderRefunded

Agregación diaria de ventas con upsert idempotente (ON CONFLICT). Alimenta el dashboard de analytics con total_orders, total_revenue y total_refunds por día.

Target: tabla daily_sales en PostgreSQL Read DB
👤

Customer Activity

customer_activity
CustomerCreated OrderCompleted ReviewSubmitted CustomerTierChanged

Perfil de actividad de agencias compradoras. Actualiza 3 tablas en una transacción: customers, customer_activity_summary, customer_order_history.

Target: 3 tablas en transacción (Inline projection)
Implementación: Clase Base y Projector
projections/base.py
Core
from abc import ABC, abstractmethod
from dataclasses import dataclass
from typing import List, Optional
import asyncio
import logging

logger = logging.getLogger("cultivamarket.projections")


@dataclass
class Event:
    stream_id: str
    event_type: str
    data: dict
    version: int
    global_position: int


class Projection(ABC):
    """Clase base para todas las proyecciones de CultivaMarket."""

    @property
    @abstractmethod
    def name(self) -> str:
        """Nombre único para checkpointing."""
        pass

    @abstractmethod
    def handles(self) -> List[str]:
        """Tipos de eventos que maneja esta proyección."""
        pass

    @abstractmethod
    async def apply(self, event: Event) -> None:
        """Aplica el evento al modelo de lectura."""
        pass


class Projector:
    """
    Orquestador de proyecciones con checkpointing persistente.
    Garantiza at-least-once delivery y manejo de errores por proyección.
    """

    def __init__(self, event_store, checkpoint_store):
        self.event_store = event_store
        self.checkpoint_store = checkpoint_store
        self.projections: List[Projection] = []

    def register(self, projection: Projection):
        """Registra una proyección para ser ejecutada por el Projector."""
        self.projections.append(projection)
        logger.info(f"Proyección registrada: {projection.name}")

    async def run(self, batch_size: int = 100):
        """Ejecuta todas las proyecciones en loop continuo."""
        logger.info("Iniciando Projector de CultivaMarket...")
        while True:
            for projection in self.projections:
                try:
                    await self._run_projection(projection, batch_size)
                except Exception as e:
                    logger.error(f"Error en proyección {projection.name}: {e}")
            await asyncio.sleep(0.1)

    async def _run_projection(self, projection: Projection, batch_size: int):
        checkpoint = await self.checkpoint_store.get(projection.name)
        position = checkpoint or 0

        events = await self.event_store.read_all(position, batch_size)

        for event in events:
            if event.event_type in projection.handles():
                await projection.apply(event)

            # Checkpoint SIEMPRE — incluso para eventos no manejados
            await self.checkpoint_store.save(
                projection.name,
                event.global_position
            )

    async def rebuild(self, projection: Projection):
        """Reconstruye una proyección desde cero (Catchup completo)."""
        logger.warning(f"Iniciando rebuild de {projection.name}...")
        await self.checkpoint_store.delete(projection.name)
        await self._run_projection(projection, batch_size=1000)
        logger.info(f"Rebuild de {projection.name} completado.")
Proyección 1: Order Summary
projections/order_summary.py
PostgreSQL Read DB
class OrderSummaryProjection(Projection):
    """
    Proyecta eventos de pedidos a una vista desnormalizada.
    Optimizada para el panel de pedidos de las agencias en CultivaMarket.
    """

    def __init__(self, db_pool: asyncpg.Pool):
        self.pool = db_pool

    @property
    def name(self) -> str:
        return "order_summary"

    def handles(self) -> List[str]:
        return [
            "OrderCreated", "OrderItemAdded", "OrderItemRemoved",
            "OrderShipped", "OrderCompleted", "OrderCancelled"
        ]

    async def apply(self, event: Event) -> None:
        handlers = {
            "OrderCreated":     self._handle_created,
            "OrderItemAdded":   self._handle_item_added,
            "OrderItemRemoved": self._handle_item_removed,
            "OrderShipped":     self._handle_shipped,
            "OrderCompleted":   self._handle_completed,
            "OrderCancelled":   self._handle_cancelled,
        }
        handler = handlers.get(event.event_type)
        if handler:
            await handler(event)

    async def _handle_created(self, event: Event):
        async with self.pool.acquire() as conn:
            await conn.execute("""
                INSERT INTO order_summaries
                  (order_id, agency_id, status, total_amount, item_count,
                   currency, created_at)
                VALUES ($1, $2, 'pending', 0, 0, $3, $4)
                ON CONFLICT (order_id) DO NOTHING   -- idempotente
            """,
                event.data['order_id'], event.data['agency_id'],
                event.data.get('currency', 'EUR'), event.data['created_at']
            )

    async def _handle_item_added(self, event: Event):
        async with self.pool.acquire() as conn:
            await conn.execute("""
                UPDATE order_summaries SET
                    total_amount = total_amount + $2,
                    item_count   = item_count + 1,
                    updated_at   = NOW()
                WHERE order_id = $1
            """,
                event.data['order_id'],
                event.data['price'] * event.data['quantity']
            )
Proyección 3: Daily Sales (Idempotencia con ON CONFLICT)
projections/daily_sales.py
Idempotente
class DailySalesProjection(Projection):
    """
    Agrega ventas por día para el dashboard de CultivaMarket.
    Usa ON CONFLICT para garantizar idempotencia — seguro para re-procesar.
    """

    @property
    def name(self) -> str: return "daily_sales"

    def handles(self) -> List[str]:
        return ["OrderCompleted", "OrderRefunded"]

    async def _increment_sales(self, event: Event):
        date = event.data['completed_at'][:10]  # YYYY-MM-DD
        async with self.pool.acquire() as conn:
            await conn.execute("""
                INSERT INTO daily_sales
                    (date, total_orders, total_revenue, total_items, total_refunds)
                VALUES ($1, 1, $2, $3, 0)
                ON CONFLICT (date) DO UPDATE SET
                    total_orders  = daily_sales.total_orders  + 1,
                    total_revenue = daily_sales.total_revenue + $2,
                    total_items   = daily_sales.total_items   + $3,
                    updated_at    = NOW()
            """,
                date,
                event.data['total_amount'],
                event.data['item_count']
            )

    async def _decrement_sales(self, event: Event):
        # Reembolso: descuenta del día original de la venta
        date = event.data['original_completed_at'][:10]
        async with self.pool.acquire() as conn:
            await conn.execute("""
                UPDATE daily_sales SET
                    total_orders  = total_orders  - 1,
                    total_revenue = total_revenue - $2,
                    total_refunds = total_refunds + $2,
                    updated_at    = NOW()
                WHERE date = $1
            """,
                date, event.data['refund_amount']
            )
Schema SQL — Tablas del Read Model
migrations/001_read_model_schema.sql
PostgreSQL
-- ============================================================
-- CultivaMarket — Read Model Schema
-- Optimizado para proyecciones CQRS (sin FKs a write model)
-- ============================================================

-- Checkpoints del Projector
CREATE TABLE projection_checkpoints (
    projection_name  TEXT PRIMARY KEY,
    global_position  BIGINT NOT NULL DEFAULT 0,
    updated_at       TIMESTAMPTZ DEFAULT NOW()
);

-- Order Summary (desnormalizado)
CREATE TABLE order_summaries (
    order_id         UUID PRIMARY KEY,
    agency_id        UUID NOT NULL,
    status           TEXT NOT NULL DEFAULT 'pending',
    total_amount     NUMERIC(10,2) NOT NULL DEFAULT 0,
    item_count       INT NOT NULL DEFAULT 0,
    currency         CHAR(3) NOT NULL DEFAULT 'EUR',
    created_at       TIMESTAMPTZ NOT NULL,
    shipped_at       TIMESTAMPTZ,
    completed_at     TIMESTAMPTZ,
    cancelled_at     TIMESTAMPTZ,
    updated_at       TIMESTAMPTZ DEFAULT NOW()
);
CREATE INDEX idx_order_summaries_agency ON order_summaries(agency_id, status);
CREATE INDEX idx_order_summaries_status  ON order_summaries(status, created_at DESC);

-- Daily Sales Aggregation
CREATE TABLE daily_sales (
    date             DATE PRIMARY KEY,
    total_orders     INT NOT NULL DEFAULT 0,
    total_revenue    NUMERIC(12,2) NOT NULL DEFAULT 0,
    total_items      INT NOT NULL DEFAULT 0,
    total_refunds    NUMERIC(12,2) NOT NULL DEFAULT 0,
    updated_at       TIMESTAMPTZ DEFAULT NOW()
);

-- Customer Activity (multi-tabla, inline projection)
CREATE TABLE customers (
    customer_id      UUID PRIMARY KEY,
    email            TEXT NOT NULL,
    name             TEXT NOT NULL,
    tier             TEXT NOT NULL DEFAULT 'bronze',
    created_at       TIMESTAMPTZ NOT NULL,
    updated_at       TIMESTAMPTZ DEFAULT NOW()
);

CREATE TABLE customer_activity_summary (
    customer_id      UUID PRIMARY KEY REFERENCES customers,
    total_orders     INT NOT NULL DEFAULT 0,
    total_spent      NUMERIC(12,2) NOT NULL DEFAULT 0,
    total_reviews    INT NOT NULL DEFAULT 0,
    last_order_at    TIMESTAMPTZ,
    last_review_at   TIMESTAMPTZ
);
Buenas Practicas Aplicadas

Do's

+ Idempotencia con ON CONFLICT DO NOTHING / DO UPDATE — seguro para replay de eventos
+ Transacciones para actualizaciones multi-tabla (customer_activity)
+ Checkpoint persistente por proyección — resume tras restart sin perder posición
+ Manejo de errores aislado por proyección — un fallo no para las demás
+ Desnormalización deliberada para optimizar patrones de query

Don'ts Evitados

- No acoplamos proyecciones entre sí — cada una es independiente
- No ignoramos el orden de eventos — global_position garantiza secuencia
- No sobreNormalizamos el read model — las tablas son planas intencionalmente
- No saltamos errores silenciosamente — logging + alerta en cada excepción
- No guardamos lógica de dominio en proyecciones — solo transformar y persistir