Procesamiento Big Data con Apache Spark (Actualizado 2025)

Procesamiento Big Data con Apache Spark (Actualizado 2025)

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_udf con 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.

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).

Métricas de negocio (6 meses)