Arquitecturas Modernas de Datos
Las arquitecturas son los blueprints que definen cómo fluyen, se almacenan y se procesan los datos en una organización. Entender sus trade-offs es lo que diferencia a un Senior DE de un programador de pipelines.
Lambda Architecture
Lambda divide el procesamiento en tres capas paralelas para manejar tanto datos históricos (batch) como datos en tiempo real (stream), combinando sus resultados en la capa de servicio.
Events, Logs, APIs] end subgraph "Batch Layer (Alto Throughput)" B[🗄️ Immutable Master Dataset
HDFS / S3] C[⚡ Batch Processing
Spark / MapReduce] D[📊 Batch Views
Precomputed Results] end subgraph "Speed Layer (Baja Latencia)" E[🔄 Stream Processing
Kafka Streams / Flink] F[⚡ Real-time Views
Redis / Cassandra] end subgraph "Serving Layer" G[🔀 Query Merger
Combina Batch + Speed] H[📱 Applications & APIs] end A --> B A --> E B --> C C --> D E --> F D --> G F --> G G --> H style B fill:#1e3a5f,stroke:#3b82f6,color:#fff style E fill:#1a3a2a,stroke:#10b981,color:#fff style G fill:#3d1a6b,stroke:#8b5cf6,color:#fff
Las 3 Capas de Lambda
- Batch Layer: Procesa el dataset completo histórico periódicamente (horas/días). Resultados precisos y exactos. Herramientas: Spark, Hadoop. Storage: HDFS, S3.
- Speed Layer: Procesa solo los datos recientes desde el último batch. Resultados aproximados pero de baja latencia. Herramientas: Kafka Streams, Flink, Spark Streaming.
- Serving Layer: Combina las vistas de batch y speed para responder queries. Responde la pregunta "¿cuál es el estado del sistema ahora mismo?" fusionando resultados exactos (batch) con los recientes (speed).
La misma lógica de negocio debe implementarse DOS veces: una en batch (ej: PySpark) y otra en streaming (ej: Kafka Streams). Esto es el "lambda problem": los dos sistemas inevitablemente divergen, causando inconsistencias sutiles. Es la principal razón por la que la industria migró hacia Kappa Architecture.
Kappa Architecture
Kappa simplifica Lambda eliminando la batch layer. Todo es streaming. Para el reprocessing histórico, se replay los eventos desde el log (Kafka) como si fueran streaming reciente.
Event Log
Replayable] end subgraph "Processing (Una sola capa)" B --> C[⚡ Stream Processing
Flink / Kafka Streams
Spark Structured Streaming] end subgraph "Serving" C --> D[🗄️ Serving Store
Cassandra / Redis / DWH] D --> E[📱 Applications] end subgraph "Reprocessing" B -.->|Replay histórico| C style B fill:#1e3a5f,stroke:#3b82f6,color:#fff end style C fill:#1a3a2a,stroke:#10b981,color:#fff style D fill:#3d1a6b,stroke:#8b5cf6,color:#fff
Kappa en Acción: Reprocessing
# Kappa Architecture: Re-procesamiento desde Kafka
# Cuando cambias la lógica, simplemente replay desde el offset 0
from confluent_kafka import Consumer, Producer
import json
def reprocess_from_beginning(topic: str, new_processor):
"""Replay todos los eventos desde el inicio para re-calcular resultados."""
consumer = Consumer({
'bootstrap.servers': 'kafka:9092',
'group.id': f'reprocess_{topic}_{datetime.now().timestamp()}',
'auto.offset.reset': 'earliest', # 🔑 Desde el principio
})
consumer.subscribe([topic])
print(f"🔄 Starting reprocessing of topic: {topic}")
processed = 0
while True:
msg = consumer.poll(timeout=5.0)
if msg is None:
print(f"✅ Reprocessing complete: {processed:,} events")
break
if msg.error():
break
event = json.loads(msg.value().decode('utf-8'))
new_processor.process(event)
processed += 1
consumer.close()
# Ventaja clave: cambiar la lógica = solo replay el topic
# No hay batch layer que recalcular por separado
Medallion Architecture
La arquitectura Medallion organiza los datos en tres capas de calidad progresiva: Bronze (raw), Silver (cleaned), Gold (business-ready). Es la arquitectura más adoptada en el ecosistema cloud-native 2026.
Bronze Layer — Raw Data (Zona de aterrizaje)
Datos exactamente como llegan de la fuente. Sin transformaciones. Inmutable. Sirve como source of truth y permite replay. Formato: Parquet, Avro, Delta. Retención: larga (años). Schema: schema-on-read o schema del source.
Silver Layer — Cleaned & Enriched
Datos limpios, deduplicados, con schemas validados, datos nulos manejados, tipos correctos, uniones básicas aplicadas. Ya es útil para Data Scientists y análisis ad-hoc. PII puede anonimizarse aquí.
Gold Layer — Business-Ready Aggregates
Tablas de hechos y dimensiones, KPIs pre-calculados, modelos dbt finales, datos listos para BI tools. Alta calidad garantizada. Star schema, agregaciones, métricas de negocio. Optimizados para queries rápidas.
Implementación con Delta Lake y dbt
# ── BRONZE: Ingesta raw con Delta Lake ────────────────────────────────────
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("medallion").getOrCreate()
# Ingesta raw desde S3 → Bronze Delta Table
raw_df = spark.read.json("s3://bucket/raw/orders/2026/01/01/")
(raw_df
.withColumn("_ingested_at", current_timestamp())
.withColumn("_source_file", input_file_name())
.write.format("delta")
.mode("append")
.partitionBy("_ingested_date")
.save("s3://bucket/bronze/orders/")
)
# ── SILVER: Limpieza con PySpark ──────────────────────────────────────────
from pyspark.sql.functions import col, when, trim, lower, to_timestamp
bronze_df = spark.read.format("delta").load("s3://bucket/bronze/orders/")
silver_df = (bronze_df
# Eliminar duplicados
.dropDuplicates(["order_id"])
# Limpiar tipos
.withColumn("order_date", to_timestamp(col("order_date"), "yyyy-MM-dd'T'HH:mm:ss"))
.withColumn("amount", col("amount").cast("decimal(10,2)"))
.withColumn("customer_email", lower(trim(col("customer_email"))))
# Filtrar inválidos
.filter(col("order_id").isNotNull())
.filter(col("amount") >= 0)
# SCD2 merge
)
DeltaTable.forPath(spark, "s3://bucket/silver/orders/") \
.alias("target").merge(
silver_df.alias("source"),
"target.order_id = source.order_id"
).whenMatchedUpdateAll() \
.whenNotMatchedInsertAll() \
.execute()
# ── GOLD: Agregaciones con dbt ────────────────────────────────────────────
-- dbt model: gold/mart_sales_monthly.sql
-- Materialización como tabla (para mejor performance BI)
{{ config(materialized='table', schema='gold') }}
WITH silver_orders AS (
SELECT * FROM {{ ref('silver_orders') }}
WHERE order_status NOT IN ('cancelled', 'returned')
),
monthly_metrics AS (
SELECT
DATE_TRUNC('month', order_date) AS month,
customer_segment,
product_category,
COUNT(DISTINCT order_id) AS order_count,
COUNT(DISTINCT customer_id) AS unique_customers,
SUM(order_amount) AS gross_revenue,
SUM(discount_amount) AS total_discounts,
SUM(order_amount - discount_amount) AS net_revenue,
AVG(order_amount) AS avg_order_value,
SUM(order_amount) / COUNT(DISTINCT customer_id) AS revenue_per_customer
FROM silver_orders
GROUP BY 1, 2, 3
)
SELECT
*,
SUM(net_revenue) OVER (PARTITION BY customer_segment ORDER BY month) AS cumulative_revenue,
net_revenue / SUM(net_revenue) OVER (PARTITION BY month) AS market_share
FROM monthly_metrics
Data Mesh
Data Mesh no es una tecnología — es un paradigma sociotécnico. Propone descentralizar la propiedad de los datos, tratándolos como productos, con equipos de dominio responsables end-to-end de sus datos.
Los 4 Principios del Data Mesh
1️⃣ Domain Ownership
Los datos son propiedad del dominio que los genera. El equipo de "Orders" es responsable de los datos de órdenes, no un equipo central de datos.
2️⃣ Data as a Product
Cada dominio trata sus datos como un producto con SLAs, documentación, calidad garantizada y usuarios. Los datos son "first-class citizens".
3️⃣ Self-serve Platform
Una plataforma de datos centralizada provee las herramientas para que cualquier dominio pueda publicar y consumir datos sin fricción.
4️⃣ Federated Governance
Governance global (estándares, seguridad, compliance) con implementación descentralizada. "Federación" como en gobierno federal.
Storage, Compute,
Catalog, Governance] end subgraph "Dominio: Orders" O1[📋 Orders Data Product] O2[orders_v1 API] O3[orders_stream topic] end subgraph "Dominio: Customers" C1[👤 Customer Data Product] C2[customers_v2 API] C3[customer_events topic] end subgraph "Dominio: Finance" F1[💰 Finance Data Product] F2[revenue_metrics API] end subgraph "Consumidores" ML[🤖 ML Team] BI[📊 BI Team] EXT[🌐 External Teams] end P --- O1 P --- C1 P --- F1 O1 --> O2 O1 --> O3 C1 --> C2 C1 --> C3 O2 --> F1 C2 --> F1 F2 --> ML O2 --> BI C2 --> BI O3 --> EXT style P fill:#3d1a6b,stroke:#8b5cf6,color:#fff style O1 fill:#1e3a5f,stroke:#3b82f6,color:#fff style C1 fill:#065f46,stroke:#10b981,color:#fff style F1 fill:#7c3a00,stroke:#f59e0b,color:#fff
La plataforma central (self-serve) sigue siendo necesaria — es el "sistema operativo" sobre el que todos los dominios trabajan. Lo que cambia es que cada dominio tiene sus propios Data Engineers que entienden profundamente el negocio de su dominio. Data Mesh requiere equipos maduros y es difícil de implementar en organizaciones pequeñas.
Data Fabric
Data Fabric es una capa de integración de datos que conecta fuentes heterogéneas mediante metadatos inteligentes, ML y automatización, sin necesariamente mover los datos. Es tecnología-first vs Data Mesh que es organización-first.
Componentes del Data Fabric
- Knowledge Graph: Grafo de metadatos que conecta conceptos, datasets, transformaciones y usuarios
- Active Metadata: Metadatos que no solo describen sino que se usan para automatizar decisiones (qué pipeline correr, qué calidad validar)
- Unified Catalog: Vista única de todos los datos de la organización, sin importar dónde viven
- Intelligent Integration: ML que automatiza mapeos de datos, detección de anomalías y optimización de pipelines
- Virtual Data Access: Consultar datos en su ubicación original sin moverlos (data virtualization)
| Aspecto | Data Mesh | Data Fabric |
|---|---|---|
| Enfoque primario | Organizacional y cultural | Tecnológico e integración |
| Implementación | Bottom-up por dominios | Top-down con plataforma unificada |
| Ownership de datos | Descentralizado (dominios) | Puede ser centralizado |
| Movimiento de datos | Los datos viven en el dominio | Virtualización, puede ser in-place |
| ROI timeline | Largo (cultura + tech) | Más corto (es una capa de software) |
| Complementariedad | Pueden coexistir: Data Mesh como org model, Data Fabric como tech layer | |
Comparativa Global de Arquitecturas
| Arquitectura | Paradigma | Madurez Org. | Latencia | Complejidad | Cuándo usar |
|---|---|---|---|---|---|
| Lambda | Batch + Stream paralelo | Cualquiera | Baja (speed layer) | 🔴 Alta | Legacy. Evitar en nuevos proyectos |
| Kappa | Solo Streaming | Técnica media+ | Muy baja | 🟡 Media | Cuando todo es event-driven |
| Medallion | Capas de calidad | Cualquiera | Batch + near-RT | 🟢 Baja | El estándar 2026 para la mayoría |
| Data Mesh | Dominios descentralizados | Alta | Variable | 🔴 Alta (org) | Empresas grandes con múltiples dominios |
| Data Fabric | Integración virtual | Media+ | Variable | 🟡 Media | Múltiples silos, sin mover datos |
Para la mayoría de empresas: Medallion Architecture + dbt para transformaciones es la combinación ganadora. Simple, bien documentada, con tooling maduro y compatibilidad con todos los cloud providers. Si la empresa crece y tiene múltiples equipos, evolucionar hacia Data Mesh en paralelo. Data Fabric como capa adicional si hay muchos sistemas legacy heterogéneos.