Apache Spark 3.5+ en 2025: Casos reales, métricas y arquitectura LLM-ready
En 2016 escribíamos que Spark se había «asentado como pilar tecnológico». En 2025, la frase se queda corta: Spark es el plano de control donde convergen data engineering, ML clásico e IA generativa. El reto ya no es «procesar big data», sino servir features a un agente LangGraph en <50ms, indexar 50TB de PDFs para RAG sin reventar el presupuesto GPU, y gobernar todo bajo Unity Catalog o OpenLineage.
En BAOSS hemos acompañado a tres clientes —retail europeo, banca tier-1, fabricante industrial— en esa transición. Aquí no hay teoría de conferencia: arquitecturas desplegadas, métricas de producción y lecciones aprendidas (anonymizadas, con números reales).
El problema 2025-2026: la pila data se ha fragmentado (y caro)
Las empresas llegan a nosotros con el mismo patrón:
- Data lake «moderno» (Iceberg/Delta/Hudi) pero consultas SQL lentas por falta de Z-Order / liquid clustering.
- Pipelines batch en Airflow + Spark que tardan 6h; el negocio pide frescura <15 min.
- Experimentos LLM en notebooks (LangChain, LlamaIndex, CrewAI) que no escalan a producción: OOM en driver, embeddings single-thread, sin versionado de prompts.
- Coste cloud descontrolado: clústeres fijos 24/7, instancias GPU ociosas, shuffle spill a disco.
- Governance rota: linaje roto entre tabla raw y feature store; PII en embeddings sin enmascarar.
La respuesta no es «comprar Databricks/Fabric/EMR Serverless y rezar». Es rediseñar la arquitectura sobre Spark 3.5+ aprovechando:
- Spark Connect: clientes ligeros (VS Code, PySpark remoto, JDBC/ODBC) sin driver local.
- Structured Streaming + stateful processing: ventanas de sesión, deduplicación exact-once, join con tablas Delta «as-of».
- Pandas API on Spark / PySpark 3.5: migración 1:1 de código pandas legacy a distribuido.
- English SDK / AI Functions: `df.select(ai(«resume esta queja en 3 bullets»))` — llama a GPT-4o / Claude 4 Sonnet vía batch inference nativo.
- GPU scheduling nativo (plugin NVIDIA RAPIDS + vLLM / Ollama en worker): inferencia embebida sin mover datos.
Caso 1 — Retail: Personalización tiempo real < 50 ms P99
Contexto: 120M eventos/día (clics, add-to-cart, compras) en Kafka. Modelo de ranking (LightGBM) + re-ranking con LLM (Llama 3.1 8B LoRA) para «razonar» ofertas. SLA: end-to-end < 50 ms desde evento a recomendación en app.
Arquitectura desplegada
- Ingest: Kafka → Spark Structured Streaming (watermark 10 min, stateRocksDB off-heap).
- Feature store online: Redis Cluster (TTL 7d) poblado por micro-batch cada 30s desde Delta Lake (CDF).
- Inferencia híbrida:
- LightGBM (CPU) → broadcast variable, scoring en mapPartitions (vectorizado Pandas UDF).
- LLM re-ranking → vLLM en workers GPU (A10G 24GB) invocado vía gRPC desde mapPartitions; batch dinámico 32 reqs.
- Serving: Spark Connect endpoint + gRPC gateway (Go) expuesto a Kubernetes (KServe).
Métricas duras (3 meses en prod)
| KPI | Antes (batch nightly + heurísticas) | Después (streaming + LLM) |
|---|---|---|
| Latencia P99 evento→reco | ~4 h (batch) | 42 ms |
| CTR uplift | baseline | +18 % |
| Coste infra / mes | €42k (EMR fijo + SageMaker endpoints) | €25k (EMR Serverless spot + GPU autoscaling 0→12) |
| Tiempo despliegue nuevo modelo | 2 semanas (CI/CD manual) | 4 h (MLflow 2.14 + GitOps ArgoCD) |
Clave técnica: stateful streaming + broadcast model + vLLM colocalizado evita saltos de red. El 40 % de reducción de coste vino de apagar clústeres fijos y usar EMR Serverless + Spot Fleet (max 20 % on-demand) con dynamic allocation agresivo (scale-down 60s idle).
Caso 2 — Banca: RAG corporativo a escala (50 TB, 12M docs)
Reto: Indexar histórico legal/contratos (PDF, scans OCR, emails) para asistentes internos (Claude 4 Opus + RAG). Pipeline legacy: Python + LlamaIndex single-thread → 14 días full re-index. Objetivo: < 8 h y gobernanza total (linaje, PII, versionado prompts).
Pipeline Spark-native
- Ingest: Autoloader (cloudFiles) → Delta Lake bronze (binary + metadata).
- Parsing distribuido:
pandas_udfcon marker/pdfplumber (CPU) + Tesseract GPU para scans. Output: markdown limpio + bounding boxes. - Chunking semántico: LangGraph agent (recursive splitter + LLM boundary detection) ejecutado como mapInPandas (batch 512 docs/worker).
- Embeddings: vLLM + BGE-M3 (1024-d) en workers GPU;
foreachPartition→ upsert batch 1k a Milvus/Zilliz (HNSW, SQ8). - Governance:
- Unity Catalog: tabla
silver.chunks(linaje automático desde bronze). - OpenLineage + Marquez: trazabilidad prompt→chunk→embedding.
- PII detection (Presidio) en silver; masking policy en Catálogo.
- Unity Catalog: tabla
Resultados
- Full re-index 50 TB: 6.3 h (vs 14 días) — 53x speedup.
- Coste indexado: €0.18 / GB (Spot GPU g5.2xlarge, 90 % utilización).
- Freshness incremental: CDC (Change Data Feed) → micro-batch 15 min → chunks nuevos en vector DB < 3 min.
- Adopción: 1.2k usuarios activos/día; latencia RAG P99 1.2 s (retrieval 300 ms + LLM 900 ms).
El salto no fue «usar Spark», sino modelar el pipeline RAG como ETL declarativo (Delta Live Tables / DLT) con expectations de calidad (chunk_size, embedding_norm, pii_free) que frenan deployment si fallan.
Caso 3 — Industrial: Mantenimiento predictivo + GenAI reports (ROI 3.2x en 6 meses)
Planta: 4.5k sensores (vibración, temperatura, presión) @ 10 Hz → 3.8 TB/día. Objetivo: detectar anomalías y generar informe técnico en lenguaje natural para operario de turno.
Arquitectura «Lakehouse + AI» unificada
- Edge → Kafka → Iceberg (MinIO on-prem): Spark Structured Streaming (exactly-once, checkpoint S3-compatible).
- Feature engineering: Ventanas temporales (1h, 24h), FFT, estadísticos rolling — todo en Pandas API on Spark (migración directa de notebooks data science).
- ML clásico: Isolation Forest + XGBoost (Spark MLlib) entrenado en cluster CPU (r6gd.16xlarge). Model registry MLflow + Model Serving via Spark UDF (batch scoring nocturno).
- GenAI report: Alerta → LangGraph agent (retrieve contexto histórico + manuales PDF indexados en Caso 2) → GPT-4o mini (batch inference via
ai()function) → PDF/HTML servido en portal interno. - Human-in-the-loop: Operario valida/corrige → feedback loop a Delta (tabla
gold.operator_feedback) → re-entrenamiento semanal automatizado (DLT pipeline).

