Modelos de lectura optimizados para el marketplace B2B de servicios de IA
Suscripción en tiempo real al stream de eventos. Procesa eventos a medida que se publican.
Procesa eventos históricos desde una posición determinada. Ideal para reconstruir modelos.
Almacena un checkpoint para continuar desde el último punto tras reinicios o fallos.
Ejecutada en la misma transacción que el write. Garantiza consistencia fuerte.
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.
order_summaries en PostgreSQL Read DBÍndice Elasticsearch de recursos digitales (prompts, plantillas, automatizaciones) con búsqueda full-text, filtros de categoría y rango de precios.
cultivamarket-products en Elasticsearch 8Agregació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.
daily_sales en PostgreSQL Read DBPerfil de actividad de agencias compradoras. Actualiza 3 tablas en una transacción: customers, customer_activity_summary, customer_order_history.
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.")
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'] )
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'] )
-- ============================================================ -- 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 );