Event Sourcing con Datos Estáticos en 2025: Cómo Resolvimos la Expiración Temporal sin Schedulers Distribuidos
En agosto de 2020 publicábamos en este blog una aproximación a event sourcing desde datos estáticos usando Kafka. El problema era claro: los eventos no se «despiertan» solos. Una promoción con expiration_date no genera automáticamente un evento de desactivación al llegar la fecha. La solución entonces pasaba por KTable joins y foreign keys en Kafka 2.4 (KIP-213).
Cinco años después, el panorama ha cambiado radicalmente. En BAOSS hemos acompañado a 12 clientes enterprise (banca, retail, telco, logística) en la migración de arquitecturas event-driven legacy a plataformas nativas cloud con Kafka 3.7+, Flink SQL y Iceberg como capa de almacenamiento analítico. El problema de la expiración temporal no ha desaparecido —se ha complejizado con event versioning, schema evolution y la irrupción de agentes IA autónomos que consumen y producen eventos en tiempo real.
El Problema Real 2025-2026: Tiempo, Consistencia y Agentes IA
Hoy, una plataforma event-driven típica en nuestros clientes maneja:
- 50-200M eventos/día en tópicos core (órdenes, precios, inventario, pagos)
- 15-40 servicios consumiendo/produciendo con contratos Avro/Protobuf versionados (Schema Registry)
- 3-5 agentes IA autónomos (LangGraph, CrewAI, AutoGen) que ejecutan workflows multi-paso sobre streams de eventos
- Requisitos de consistencia eventual < 500ms y exactly-once en flujos críticos
- Cumplimiento DORA, NIS2, GDPR con trazabilidad completa de decisiones automatizadas
El caso clásico de la expiration_date en precios/promociones sigue existiendo, pero ahora convive con:
- SLA dinámicos calculados por modelos ML (GPT-4o / Claude 4 via function calling) que reevalúan prioridades cada minuto
- Políticas de retención regulatoria que exigen event replay selectivo por
correlation_idhasta 7 años - Agentes RAG + MCP que consultan knowledge graphs en Neo4j/GraphDB para enriquecer eventos con contexto histórico
- Feature flags distribuidos (LaunchDarkly, Unleash) que modifican comportamiento de consumidores sin redespiegue
Los schedulers distribuidos tradicionales (Quartz, Spring Batch, Airflow, incluso Big Ben de Walmart) introducen acoplamiento temporal, single points of failure y latencia impredecible. En 2025, la respuesta no es «más schedulers» —es stream processing nativo con semántica temporal.
Arquitectura de Referencia BAOSS 2025: Flink SQL + Iceberg + Kafka
Nuestra aproximación actual para event sourcing desde datos estáticos con atributos temporales combina tres capas:
- Capa de ingesta: Kafka 3.7+ (KRaft mode) con topic tiered storage a S3/GCS. Productores usan transactional.id para exactly-once end-to-end.
- Capa de procesamiento temporal: Flink SQL 1.19+ con watermarks, TTL state y temporal table joins. Aquí resolvemos la expiración sin schedulers.
- Capa analítica / replay: Apache Iceberg sobre object storage con time travel, partition evolution y row-level deletes para GDPR.
El Patrón: Temporal Table Join en Flink SQL
En lugar de KTable foreign key joins (limitados a claves primarias y sin semántica de tiempo), usamos temporal table joins de Flink. Una temporal table es una vista versionada en el tiempo: cada clave tiene múltiples versiones con valid_from / valid_to. El join se resuelve automáticamente según el watermark del evento entrante.
-- Definición de tabla temporal de precios (versionada por expiration_date)
CREATE TABLE prices_temporal (
price_id STRING,
product_id STRING,
amount DECIMAL(10,2),
currency STRING,
status STRING, -- 'ACTIVE', 'EXPIRED', 'SUPERSEDED'
valid_from TIMESTAMP(3),
valid_to TIMESTAMP(3),
WATERMARK FOR valid_from AS valid_from - INTERVAL '5' SECOND,
PRIMARY KEY (price_id, valid_from) NOT ENFORCED
) WITH (
'connector' = 'kafka',
'topic' = 'prices-cdc',
'properties.bootstrap.servers' = 'kafka:9092',
'scan.startup.mode' = 'earliest-offset',
'format' = 'avro-confluent',
'avro-confluent.schema-registry.url' = 'http://schema-registry:8081'
);
-- Stream de órdenes entrantes (eventos de negocio)
CREATE TABLE orders_stream (
order_id STRING,
product_id STRING,
quantity INT,
order_ts TIMESTAMP(3),
WATERMARK FOR order_ts AS order_ts - INTERVAL '10' SECOND
) WITH (
'connector' = 'kafka',
'topic' = 'orders',
'properties.bootstrap.servers' = 'kafka:9092',
'format' = 'avro-confluent'
);
-- Join temporal: cada orden obtiene el precio VIGENTE en su order_ts
SELECT
o.order_id,
o.product_id,
o.quantity,
p.amount AS unit_price,
p.currency,
(o.quantity * p.amount) AS total_amount,
o.order_ts
FROM orders_stream o
LEFT JOIN prices_temporal FOR SYSTEM_TIME AS OF o.order_ts p
ON o.product_id = p.product_id
WHERE p.status = 'ACTIVE';
Resultado: cero schedulers. La expiración es intrínseca al modelo de datos. Cuando el watermark supera valid_to, la versión expira automáticamente del state backend (RocksDB incremental checkpoints cada 30s). Flink emite changelog al tópico de salida —los consumidores downstream reciben la «desactivación» como un evento más.
Caso Real BAOSS (Anonimizado): Retail Multinacional — 47M Eventos/Día
Contexto: Cliente retail con 3.200 tiendas, e-commerce y marketplace. Plataforma legacy con Spring Batch nocturno para expirar promociones (ventana 4h, fallos frecuentes, replay manual). Migración a event-driven en 2024-2025.
Métricas Pre-Migración (Q4 2024)
- Latencia expiración promoción: 2-6 horas (batch nocturno)
- Incidencias mensuales por precios incorrectos: 12-18 (impacto €280K-420K/mes)
- Tiempo replay histórico (auditoría): 4-8 horas (dump + restore BD)
- Coste infra batch (EC2 r6g.4xlarge x 12): €18.500/mes
Implementación BAOSS (Q1-Q2 2025)
- Migración CDC (Debezium 2.5+) → Kafka → Flink SQL (Kubernetes, K8s Operator)
- Modelado temporal tables para 14 dominios (precios, promos, stock, SLA, fidelidad)
- Integración agentes IA: 3 agentes LangGraph consumiendo enriched events para dynamic pricing y stock allocation
- Capa Iceberg (Tabular/Trino) para time travel auditoría y feature store ML
- Observabilidad: OpenTelemetry + Tempo + Mimir + Grafana (latencia P99 < 200ms end-to-end)
Resultados Q3 2025 (6 meses en producción)
- Latencia expiración: < 500ms (watermark-driven, determinista)
- Incidencias precios incorrectos: 0 en 6 meses (vs 12-18/mes)
- Replay auditoría selectivo: < 3 min (Iceberg time travel + partition pruning)
- Coste infra streaming (EKS + MSK + S3): €11.200/mes (-39% vs batch)
- ROI 3.4x en 6 meses (ahorro directo + evitación sanciones + revenue recuperado)
- Despliegues nuevos dominios: 2 días vs 3 semanas (plantillas Flink SQL + CI/CD GitOps)
El equipo interno del cliente (12 ingenieros) ahora opera la plataforma de forma autónoma. BAOSS aportó arquitectura, coaching y transferencia —no «body shopping».
Segundo Caso: Banca — Cumplimiento DORA/NIS2 con Event Replay Certificado
Banco europeo (activos >€200B) necesitaba demostrar a regulador trazabilidad completa de decisiones automatizadas de crédito (modelos ML + reglas negocio) con capacidad de replay puntual por application_id hasta 7 años.
- Reto: Eventos en 38 tópicos, schemas Avro evolucionados (150+ versiones), PII en payloads.
- Solución: Flink SQL temporal joins para reconstruir state en cualquier
event_ts+ Iceberg row-level deletes para right to be forgotten (GDPR Art. 17) sin romper time travel. - Agentes IA: 2 agentes CrewAI (auditoría automática + generación narrativas regulatorio) consumen curated events via MCP (Model Context Protocol) desde Iceberg.
- Métrica clave: Auditoría DORA Q2 2025 completada en 4 días vs 6 semanas estimadas. Cero hallazgos críticos.

