El job de ventas de la lección anterior funcionaba, pero cada pregunta nueva costaba un job más: ventas por productor y mercado era uno, el ranking por mercado otro, unirlo con el catálogo un tercero, y entre cada uno la salida completa iba a HDFS y volvía. Apache Spark nació en Berkeley (2009, Matei Zaharia) precisamente de esa frustración: los algoritmos iterativos de aprendizaje automático y las consultas interactivas tardaban en MapReduce diez o cien veces más de lo que la CPU justificaba, porque el tiempo se iba en escribir y leer disco entre fases. Spark conserva lo esencial del modelo (datos en particiones, tareas independientes y reejecutables, un shuffle para agrupar por clave) pero expresa todo el cálculo como un DAG de operadores que un planificador convierte en fases y que se ejecuta manteniendo los datos intermedios en memoria. Sobre esa base construyó una API de colecciones (RDD), después una de tablas con optimizador (DataFrames y Spark SQL), y encima librerías de aprendizaje automático (MLlib), grafos (GraphX) y flujos (Structured Streaming, que veremos en 05-04). Hoy es el motor de lotes de referencia. En esta lección analitica reescribe el cálculo de ventas diarias como servicios/analitica/ventas_diarias.py, con DataFrames y con RDD para comparar, lo lanza en local y en un clúster de docker-compose.yml, lee el plan de ejecución, y entrena unas recomendaciones con ALS sobre los clics.
Contenido
- Por qué Spark: un DAG en memoria frente a fases en disco
- Arquitectura: driver, cluster manager, executors, tasks y stages
- RDD: colecciones inmutables, transformaciones perezosas y linaje
- DataFrames y Spark SQL: Catalyst y formatos columnares
- Operaciones estrechas y anchas, shuffle y stages
- Optimizaciones: cache, broadcast join, particionado y sesgo
- MapReduce frente a Spark
- MLlib: recomendaciones con ALS sobre los clics
- Práctica:
servicios/analitica/ventas_diarias.py - Errores Comunes y Consejos
- Ejercicios
- Conclusión
- Por qué Spark: un DAG en memoria frente a fases en disco
En MapReduce (05-02), una cadena de tres agrupaciones son tres jobs y seis pasos por disco: cada job escribe su salida en HDFS (replicada tres veces) y el siguiente la lee. Para un algoritmo iterativo que recorre los mismos datos cien veces, son cien lecturas completas desde disco de un conjunto que no ha cambiado. Spark cambia dos cosas:
- Todo el cálculo es un único programa, un DAG de operadores (leer, filtrar, explotar, agrupar, unir, escribir) que el motor conoce completo antes de ejecutar nada. Con esa visión global puede encadenar en una sola pasada los operadores que no necesitan redistribuir datos, y solo materializa datos intermedios en los puntos donde hay un shuffle. Es el patrón dataflow de 05-01.
- Los datos intermedios viven en memoria de los ejecutores, y se pueden cachear explícitamente para reutilizarlos en varias pasadas. Cien iteraciones sobre un conjunto cacheado leen disco una vez.
El precio de no escribir a disco es la tolerancia a fallos: si un nodo muere, sus datos intermedios en memoria desaparecen. La solución de Spark, y su idea original, es el linaje (apartado 3): en lugar de guardar los datos intermedios, guarda la receta para recalcularlos a partir de la entrada, y recomputa solo las particiones perdidas. Es la reejecución determinista de 05-01, aplicada a fragmentos de un cálculo en vez de a tareas enteras.
- Arquitectura: driver, cluster manager, executors, tasks y stages
flowchart TB
subgraph D[Driver: ventas_diarias.py]
SC[SparkSession / SparkContext<br/>DAG scheduler + task scheduler]
end
CM[Cluster manager<br/>standalone · YARN · Kubernetes]
subgraph W1[Nodo worker 1]
E1[Executor<br/>JVM con 4 núcleos y 8 GB]
T1a[task] --- E1
T1b[task] --- E1
end
subgraph W2[Nodo worker 2]
E2[Executor]
T2a[task] --- E2
T2b[task] --- E2
end
SC -- "pide executors" --> CM
CM -- "lanza" --> E1
CM -- "lanza" --> E2
SC -- "envía tasks, recibe resultados" --> E1
SC -- "envía tasks, recibe resultados" --> E2
E1 <-- "shuffle" --> E2
E1 --> HDFS[(HDFS / MinIO)]
E2 --> HDFS
- Driver. El proceso que ejecuta el programa principal (
ventas_diarias.py). Contiene elSparkSession(antesSparkContext), construye el DAG, lo divide en jobs, stages y tasks, las envía a los executors y recoge resultados. Es el maestro de 05-02, uno por aplicación: si el driver muere, la aplicación muere (en YARN en cluster mode, el driver es el ApplicationMaster y se puede reintentar). - Cluster manager. Reparte recursos entre aplicaciones. Spark trae el suyo (standalone, un maestro y workers: el que usaremos en
docker-compose.yml), y se integra con YARN (05-02: el driver pide contenedores al ResourceManager) y con Kubernetes (07-05: cada executor es un pod). Al driver le da igual cuál sea; la API es la misma. - Executor. Un proceso JVM en un nodo worker, con N núcleos y M GB asignados, que vive toda la aplicación. Ejecuta tasks en hilos (un task por núcleo a la vez) y guarda en memoria las particiones cacheadas y los datos de shuffle. A diferencia de MapReduce, no se arranca una JVM por tarea: la JVM se arranca una vez y ejecuta miles de tareas, lo que elimina el coste fijo por tarea.
- Task. La unidad de trabajo: ejecutar una cadena de operadores sobre una partición. Una fase con 200 particiones son 200 tasks.
- Stage. Un conjunto de tasks que se pueden ejecutar sin shuffle. El DAG se corta en stages en cada operación ancha (apartado 5); dentro de un stage, los operadores se encadenan (pipelining) y una fila pasa por todos ellos sin tocar disco.
- Job. Todo lo que desencadena una acción (apartado 3):
write,collect,count. Un programa puede tener varios jobs; cada job tiene uno o más stages.
Con PySpark, el driver es un proceso Python que se comunica con una JVM (vía Py4J), y en los executors el código Python de las funciones de RDD se ejecuta en procesos Python auxiliares con serialización de ida y vuelta; con DataFrames, en cambio, la mayoría de operaciones se traducen a código JVM y Python no interviene por fila. Es la razón principal por la que la API de DataFrames es la recomendada.
- RDD: colecciones inmutables, transformaciones perezosas y linaje
El RDD (Resilient Distributed Dataset) es la abstracción original de Spark: una colección de elementos particionada entre los executors, inmutable (no se modifica: se crea otro RDD a partir de él) y resiliente por linaje. Se opera con dos tipos de métodos:
| Tipo | Qué hace | Ejemplos | Ejecuta algo |
|---|---|---|---|
| Transformación | Devuelve un RDD nuevo definido a partir de otro | map, filter, flatMap, reduceByKey, groupByKey, join, distinct, repartition |
No: es perezosa, solo añade un nodo al DAG |
| Acción | Devuelve un resultado al driver o escribe en un sink | collect, count, take, reduce, saveAsTextFile, foreach |
Sí: lanza un job que ejecuta todas las transformaciones pendientes |
La pereza es la que permite optimizar: cuando el programa dice filter y luego map y luego count, Spark no ejecuta tres pasadas; al llegar a count conoce toda la cadena y la ejecuta en una sola pasada por partición. Y el linaje es el DAG de transformaciones que llevó a cada RDD: si se pierde una partición del RDD ventas porque murió su executor, Spark mira el linaje (ventas = lineas.reduceByKey(...), lineas = eventos.flatMap(...), eventos = sc.textFile(...)) y recalcula solo esa partición desde el bloque de HDFS correspondiente. No hay que guardar nada intermedio en disco; el coste es recomputar, que es asumible mientras el linaje sea corto (para linajes largos, iterativos, existe checkpoint(), que sí materializa en HDFS y corta el linaje).
El cálculo de ventas con RDD, que servirá de comparación en la práctica:
# Fragmento de servicios/analitica/ventas_diarias.py (versión RDD, ver apartado 9)
import json
from datetime import datetime, timezone
def dia_de(ev): # tiempo de evento (01-05), no de proceso
return datetime.fromtimestamp(ev["fecha_ms"] / 1000, tz=timezone.utc).strftime("%Y-%m-%d")
def ventas_rdd(sc, entrada):
eventos = sc.textFile(entrada) # RDD[str]: una partición por bloque HDFS
pedidos = eventos.map(json.loads).filter(lambda e: e["tipo"] == "pedido.creado")
lineas = pedidos.flatMap(lambda e: [ # una fila por línea de pedido
((dia_de(e), e["datos"]["mercado"], ln["productor"]), ln["cantidad"] * ln["precio"])
for ln in e["datos"]["lineas"]])
ventas = lineas.reduceByKey(lambda a, b: a + b) # ancha: shuffle por clave; combina localmente antes
return ventas # nada se ha ejecutado todavíaHasta aquí no se ha leído ni un byte: textFile, map, filter, flatMap y reduceByKey son transformaciones. ventas.collect() o ventas.saveAsTextFile(...) lanzarían el job. Dos detalles que distinguen a un usuario de Spark: reduceByKey hace la reducción local en cada partición antes del shuffle (el combiner de 05-02, automático), mientras que groupByKey seguido de una suma mueve todos los valores por la red; y las funciones lambda viajan serializadas a los executors, así que no pueden capturar objetos no serializables (una conexión a base de datos, por ejemplo).
- DataFrames y Spark SQL: Catalyst y formatos columnares
Un RDD es una colección de objetos opacos para Spark: no sabe que ln["productor"] es una columna, así que no puede optimizar nada más allá del encadenamiento. Un DataFrame es una tabla con esquema (columnas con nombre y tipo), distribuida en particiones como un RDD, y sobre la que se opera con una API relacional (select, filter, groupBy, join, agg) o directamente con SQL (spark.sql("SELECT ...")). Ambas producen el mismo plan lógico, que pasa por Catalyst, el optimizador:
- Análisis: resolver nombres de columnas y tipos contra el esquema.
- Optimización lógica: reglas como empujar los filtros hacia la lectura (predicate pushdown: filtrar
tipo = 'pedido.creado'al leer, no después), leer solo las columnas necesarias (column pruning), simplificar expresiones, reordenar joins. - Planificación física: elegir algoritmos: un join se hace por broadcast hash si un lado es pequeño, por sort-merge si no; una agregación en dos fases (parcial en cada partición, final tras el shuffle).
- Generación de código: Tungsten compila el plan a bytecode Java especializado (whole-stage codegen) y gestiona la memoria fuera del heap de la JVM en formato binario compacto. Es lo que hace que PySpark con DataFrames sea tan rápido como Scala: Python solo describe el plan.
Desde Spark 3, el Adaptive Query Execution (AQE) reoptimiza el plan durante la ejecución con estadísticas reales: reduce el número de particiones tras un shuffle pequeño, convierte un sort-merge join en broadcast si descubre que un lado cabe, y divide particiones sesgadas (apartado 6).
Los DataFrames se combinan con los formatos columnares. Un fichero Parquet guarda los datos por columnas en lugar de por filas: todos los productor juntos, todos los importe juntos, comprimidos con codificaciones adecuadas a cada tipo (diccionario para cadenas repetidas, run-length para valores consecutivos) y con estadísticas (mínimo, máximo, nulos) por bloque de filas. Una consulta que solo necesita productor e importe lee esas dos columnas y salta el resto del fichero, y un filtro dia = '2026-09-14' salta bloques enteros cuyo rango no lo contiene. Frente a JSONL, que obliga a leer y parsear cada byte, Parquet reduce en un orden de magnitud tanto los bytes leídos como el tamaño en disco. Por eso el lago de 04-02 empieza en JSONL (lo que produce Kafka) y el primer paso de la analítica lo convierte a Parquet, particionado por día.
- Operaciones estrechas y anchas, shuffle y stages
La división del DAG en stages depende de un único criterio: qué dependencia tiene cada partición de salida respecto a las de entrada.
| Dependencia estrecha (narrow) | Dependencia ancha (wide) | |
|---|---|---|
| Cada partición de salida depende de | Una partición de entrada (o unas pocas fijas) | Todas (o muchas) las particiones de entrada |
| Operaciones | map, filter, flatMap, select, withColumn, union, coalesce, join con broadcast, join con particionado idéntico |
groupBy/reduceByKey, distinct, join sort-merge, orderBy, repartition |
| Coste | Se encadena en la misma task, sin red | Shuffle: escribir, transferir, leer; corta el DAG en un nuevo stage |
| Recuperación tras fallo | Recomputar la partición perdida desde su única entrada | Recomputar puede exigir releer muchas particiones (Spark guarda los ficheros de shuffle para evitarlo) |
flowchart LR
subgraph S0[Stage 0: sin shuffle]
R[read json<br/>2 particiones] --> F[filter tipo] --> X[explode líneas] --> P[agg parcial<br/>por dia, mercado, productor]
end
P == "shuffle<br/>hashpartitioning(dia, mercado, productor)<br/>200 particiones" ==> S1
subgraph S1[Stage 1]
A[agg final] --> J[broadcast join<br/>con catálogo] --> W[write parquet]
end
C[read catálogo<br/>1 partición] -. broadcast a todos los executors .-> J
El DAG de ventas tiene exactamente un shuffle, el groupBy, y por tanto dos stages. Todo lo demás se encadena: una fila del JSON se filtra, se explota y se agrega parcialmente en la misma task sin que nadie la escriba. El join con el catálogo, que en MapReduce sería un tercer job, es estrecho gracias al broadcast (apartado 6). En la interfaz web de Spark (http://localhost:4040 durante la ejecución) cada job aparece con sus stages, cada stage con sus tasks, y para cada stage los bytes de shuffle write y shuffle read: son los contadores de 05-02, y el criterio de diagnóstico es el mismo.
El shuffle de Spark escribe también en disco local (los ficheros de shuffle, servidos por los executors o por un external shuffle service), así que "en memoria" no significa "sin disco": significa que entre operaciones estrechas no se escribe, y que los datos cacheados se sirven de memoria. Con 200 particiones de shuffle por defecto (spark.sql.shuffle.partitions), un job pequeño produce 200 tasks minúsculas en el segundo stage; AQE las fusiona, pero en versiones antiguas o con AQE desactivado conviene bajar el número.
- Optimizaciones: cache, broadcast join, particionado y sesgo
cache() y persist(). Marcan un DataFrame o RDD para que, la primera vez que se calcule, sus particiones se guarden (en memoria por defecto; persist(StorageLevel.MEMORY_AND_DISK) permite desbordar a disco, DISK_ONLY, o serializado para ahorrar memoria). Compensa cuando el mismo resultado intermedio se usa en varias acciones: el DataFrame lineas de la práctica alimenta las ventas por productor, un ranking y una comprobación de calidad; sin cache, cada acción vuelve a leer y parsear el JSON de HDFS. Con unpersist() se libera. Cachear algo que se usa una vez es puro coste.
Broadcast join. Unir las ventas (millones de filas, repartidas) con el catálogo (unas decenas de filas) no debería mover las ventas. Con F.broadcast(catalogo), Spark envía una copia del catálogo a cada executor y el join se resuelve localmente, sin shuffle del lado grande: es el map-side join de 05-02, automático. Spark lo hace solo si estima que el lado pequeño ocupa menos de spark.sql.autoBroadcastJoinThreshold (10 MB por defecto); forzarlo con broadcast() conviene cuando la estimación falla (un CSV sin estadísticas). Con dos lados grandes, el plan es un sort-merge join: shuffle de ambos por la clave y mezcla ordenada.
repartition(n) y coalesce(n). El número de particiones gobierna el paralelismo. repartition(n) hace un shuffle completo para obtener n particiones equilibradas (o repartition("dia") para colocar cada día en su partición, útil antes de escribir particionado); coalesce(n) reduce el número sin shuffle, fusionando particiones locales, y es la forma barata de no escribir 200 ficheros Parquet de 3 KB. Como orientación: particiones de 100–200 MB en memoria y entre 2 y 4 tasks por núcleo disponible.
Sesgo con salting. La Semana del Queso Artesano vuelve: al agrupar por productor, la partición de queseria-montblanc recibe la mitad de las filas. En DataFrames, la agregación parcial por partición (que Catalyst inserta siempre) alivia el caso de las sumas, igual que el combiner; el sesgo duele cuando la operación no se prerreduce (un join sort-merge, collect_list, ventanas). La técnica de 05-01, ahora con columnas:
from pyspark.sql import functions as F
N_SAL = 8
# 1) añadir sal determinista a la clave caliente (derivada del id de pedido: reejecutable)
con_sal = lineas.withColumn("sal", F.when(F.col("productor") == "queseria-montblanc",
F.pmod(F.hash("pedido_id"), F.lit(N_SAL))).otherwise(F.lit(0)))
# 2) agregar por (clave, sal): la clave caliente se reparte en 8 particiones
parcial = con_sal.groupBy("dia", "mercado", "productor", "sal").agg(F.sum("importe").alias("importe"))
# 3) segunda agregación, ya pequeña, quitando la sal
ventas = parcial.groupBy("dia", "mercado", "productor").agg(F.sum("importe").alias("importe"))Y desde Spark 3, AQE con spark.sql.adaptive.skewJoin.enabled=true detecta particiones de join sesgadas (por tamaño relativo a la mediana) y las divide automáticamente, sin sal manual. Para agregaciones sesgadas sin prerreducción posible sigue haciendo falta la sal.
- MapReduce frente a Spark
| MapReduce (Hadoop) | Spark | |
|---|---|---|
| Modelo | Map y reduce; cadenas de jobs | DAG de operadores en un programa |
| Datos intermedios | Disco local + HDFS entre jobs | Memoria (y disco local solo en el shuffle) |
| Tolerancia a fallos | Reejecución de tareas desde disco | Recomputación por linaje; ficheros de shuffle; checkpoint opcional |
| Coste fijo | JVM por tarea; decenas de segundos por job | Executors persistentes; tasks de milisegundos |
| API | Java (Streaming para otros lenguajes); solo map/reduce | Scala, Java, Python, R, SQL; decenas de operadores; DataFrames con optimizador |
| Joins, iteraciones | A mano, varios jobs | Nativos; cache para iterar |
| Iterativo (ML, grafos) | 10–100× más lento por relecturas | Diseñado para ello (MLlib, GraphX) |
| Interactivo | No | Sí (spark-shell, notebooks, Spark SQL) |
| Streaming | No | Structured Streaming (05-04) |
| Memoria necesaria | Modesta | Mayor: los executors necesitan RAM para cache y shuffle |
| Estado actual | Base histórica; YARN y HDFS siguen en uso | Motor de lotes de referencia |
La ventaja de Spark en el job de ventas de 05-02 se ve en los números: el mismo cálculo que en MapReduce tardaba 52 s en el clúster de pruebas tarda 6 s en Spark sobre YARN con los mismos recursos, y encadenar el ranking y el join con el catálogo no añade jobs, sino un stage. La contrapartida es la memoria: un executor mal dimensionado (demasiado poca memoria para el shuffle o la cache) falla con OutOfMemoryError donde MapReduce simplemente escribía a disco.
- MLlib: recomendaciones con ALS sobre los clics
El segundo caso de la analítica de Kilómetro Cero, "productos recomendados para Ana", es un problema iterativo: exactamente el tipo de problema que MapReduce resolvía mal. MLlib trae algoritmos distribuidos listos, y para recomendación el clásico es ALS (Alternating Least Squares), una factorización de la matriz clientes × productos que aprende un vector de factores por cliente y otro por producto a partir de interacciones (compras, clics) y predice el interés de cada pareja no observada. Con clics no hay "puntuación", sino señales implícitas (vio, añadió al carrito, compró), y ALS tiene un modo para eso.
Los clics del lago, /km0/clics/2026-09-14/hora=13/web-01.jsonl (04-02), tienen esta forma:
{"fecha_ms":1789390812000,"cliente":"ana","producto":"queso-curado","accion":"vio"}
{"fecha_ms":1789390834000,"cliente":"ana","producto":"vino-crianza","accion":"carrito"}
{"fecha_ms":1789390901000,"cliente":"marc","producto":"vino-crianza","accion":"compro"}
{"fecha_ms":1789391010000,"cliente":"lucia","producto":"queso-fresco","accion":"vio"}
{"fecha_ms":1789391044000,"cliente":"lucia","producto":"calabacin","accion":"compro"}# km0/servicios/analitica/recomendaciones_als.py
from pyspark.sql import SparkSession, functions as F
from pyspark.ml.feature import StringIndexer
from pyspark.ml.recommendation import ALS
spark = SparkSession.builder.appName("km0-recomendaciones").getOrCreate()
clics = spark.read.json("hdfs://namenode:8020/km0/clics/2026-09-*/") # 7 días, todas las horas
peso = F.when(F.col("accion") == "compro", 5).when(F.col("accion") == "carrito", 2).otherwise(1)
interes = clics.withColumn("peso", peso).groupBy("cliente", "producto").agg(F.sum("peso").alias("interes"))
# ALS necesita ids enteros: StringIndexer asigna uno por valor distinto
idx_cli = StringIndexer(inputCol="cliente", outputCol="cliente_id").fit(interes)
idx_pro = StringIndexer(inputCol="producto", outputCol="producto_id").fit(interes)
datos = idx_pro.transform(idx_cli.transform(interes))
als = ALS(userCol="cliente_id", itemCol="producto_id", ratingCol="interes",
implicitPrefs=True, # las señales son implícitas (clics), no notas
rank=10, maxIter=10, regParam=0.1, coldStartStrategy="drop", seed=42)
modelo = als.fit(datos) # 10 iteraciones: cada una alterna factores de clientes y productos
recomendaciones = modelo.recommendForAllUsers(3) # 3 productos por cliente
etiquetas = spark.createDataFrame(enumerate(idx_pro.labels), ["producto_id", "producto"])
(recomendaciones.select("cliente_id", F.explode("recommendations").alias("r"))
.join(etiquetas, F.col("r.producto_id") == etiquetas.producto_id)
.join(spark.createDataFrame(enumerate(idx_cli.labels), ["cliente_id", "cliente"]), "cliente_id")
.select("cliente", "producto", F.round("r.rating", 3).alias("afinidad"))
.write.mode("overwrite").parquet("hdfs://namenode:8020/km0/recomendaciones/2026-09-14"))Cada iteración de fit es un job con varios stages sobre los mismos datos, que MLlib cachea internamente; con 10 iteraciones en MapReduce serían 20 jobs y 20 lecturas de HDFS. El resultado (ana → vino-crianza, lucia → queso-curado...) lo carga el pipeline de 05-05 en la base de catalogo para que la web lo muestre. Lo que MLlib no decide es la calidad: elegir rank, regParam y los pesos de las acciones exige evaluar (RegressionEvaluator o métricas de ranking sobre un conjunto reservado), y eso pertenece a un curso de aprendizaje automático, no a este.
- Práctica:
servicios/analitica/ventas_diarias.py
servicios/analitica/ventas_diarias.py9.1 Entorno: Spark en docker-compose.yml
Al docker-compose.yml de 04-02 (HDFS) y 05-02 (YARN) se añade un clúster Spark standalone, que es más ligero que YARN para desarrollar. El driver correrá en nuestra máquina o en el contenedor del maestro:
# km0/docker-compose.yml (fragmento)
services:
spark-master:
image: bitnami/spark:3.5
environment:
- SPARK_MODE=master
ports:
- "8080:8080" # interfaz web del maestro standalone
- "7077:7077" # puerto al que se conectan drivers y workers
volumes:
- ./servicios/analitica:/app/analitica
- ./eventos:/app/eventos
spark-worker:
image: bitnami/spark:3.5
environment:
- SPARK_MODE=worker
- SPARK_MASTER_URL=spark://spark-master:7077
- SPARK_WORKER_CORES=2
- SPARK_WORKER_MEMORY=2G
deploy:
replicas: 2 # docker compose up --scale spark-worker=2
volumes:
- ./servicios/analitica:/app/analitica
- ./eventos:/app/eventosLos dos workers ofrecen 4 núcleos en total, y HDFS es alcanzable como hdfs://namenode:8020 desde la misma red de Compose. El catálogo es un CSV pequeño, servicios/analitica/catalogo.csv, que en producción vendría de una exportación diaria de km0_catalogo:
productor,nombre_productor,provincia
huerta-la-vega,Huerta La Vega,Girona
queseria-montblanc,Quesería Montblanc,Tarragona
bodega-roble-alto,Bodega Roble Alto,Lleida9.2 El programa con DataFrames
# km0/servicios/analitica/ventas_diarias.py
"""Ventas por productor, mercado y día a partir de los eventos pedido.creado del lago.
Uso: spark-submit ventas_diarias.py <entrada> <catalogo.csv> <salida> [--rdd]
entrada: hdfs://namenode:8020/km0/eventos/2026-09-14/pedidos.jsonl (o una ruta local, o un glob)
salida: hdfs://namenode:8020/km0/agregados/ventas_diarias (Parquet particionado por dia)
"""
import sys
from pyspark.sql import SparkSession, functions as F, types as T
ESQUEMA = T.StructType([ # declarar el esquema evita una pasada de inferencia
T.StructField("id_evento", T.StringType()),
T.StructField("tipo", T.StringType()),
T.StructField("version", T.IntegerType()),
T.StructField("fecha_ms", T.LongType()),
T.StructField("origen", T.StringType()),
T.StructField("datos", T.StructType([
T.StructField("pedido_id", T.StringType()),
T.StructField("cliente", T.StringType()),
T.StructField("mercado", T.StringType()),
T.StructField("lineas", T.ArrayType(T.StructType([
T.StructField("producto", T.StringType()),
T.StructField("productor", T.StringType()),
T.StructField("cantidad", T.IntegerType()),
T.StructField("precio", T.DoubleType()),
]))),
])),
])
def lineas_de_pedido(spark, entrada):
"""DataFrame con una fila por línea de pedido: dia, mercado, pedido_id, producto, productor, cantidad, importe."""
eventos = spark.read.schema(ESQUEMA).json(entrada)
pedidos = eventos.filter(F.col("tipo") == "pedido.creado") # Catalyst lo empuja a la lectura
return (pedidos
.select(
F.to_date(F.from_unixtime(F.col("fecha_ms") / 1000)).alias("dia"), # tiempo de evento
F.col("datos.mercado").alias("mercado"),
F.col("datos.pedido_id").alias("pedido_id"),
F.explode("datos.lineas").alias("ln")) # una fila por elemento del array
.select("dia", "mercado", "pedido_id",
F.col("ln.producto").alias("producto"),
F.col("ln.productor").alias("productor"),
F.col("ln.cantidad").alias("cantidad"),
(F.col("ln.cantidad") * F.col("ln.precio")).alias("importe")))
def ventas_dataframe(spark, entrada, ruta_catalogo):
lineas = lineas_de_pedido(spark, entrada).cache() # se usa en dos acciones (ventas y control)
catalogo = spark.read.option("header", True).csv(ruta_catalogo) # 3 filas: candidato a broadcast
ventas = (lineas
.groupBy("dia", "mercado", "productor") # UNA operación ancha: un shuffle
.agg(F.round(F.sum("importe"), 2).alias("importe"),
F.sum("cantidad").alias("unidades"),
F.countDistinct("pedido_id").alias("pedidos"))
.join(F.broadcast(catalogo), "productor", "left") # estrecha: el catálogo viaja a cada executor
.select("dia", "mercado", "productor", "nombre_productor", "provincia", "importe", "unidades", "pedidos"))
control = lineas.agg(F.countDistinct("pedido_id").alias("pedidos"), F.round(F.sum("importe"), 2).alias("total"))
print("Control:", control.first().asDict()) # acción 1: usa la cache
return ventas # la acción 2 será el write
if __name__ == "__main__":
entrada, ruta_catalogo, salida = sys.argv[1:4]
spark = (SparkSession.builder.appName("km0-ventas-diarias")
.config("spark.sql.shuffle.partitions", "8") # job pequeño: 200 sería absurdo
.config("spark.sql.sources.partitionOverwriteMode", "dynamic") # sobrescribir SOLO las particiones escritas
.getOrCreate())
if "--rdd" in sys.argv:
from ventas_rdd import ventas_rdd # versión RDD del apartado 3
for (dia, mercado, productor), importe in sorted(ventas_rdd(spark.sparkContext, entrada).collect()):
print(f"{dia} {mercado:10s} {productor:20s} {importe:12,.2f}")
else:
ventas = ventas_dataframe(spark, entrada, ruta_catalogo)
ventas.explain() # plan físico (apartado 9.3)
(ventas.coalesce(1) # un fichero por partición de salida
.write.mode("overwrite")
.partitionBy("dia") # salida/dia=2026-09-14/part-....parquet
.parquet(salida))
ventas.orderBy("dia", "mercado", "productor").show(truncate=False)
spark.stop()Puntos clave del programa:
- Esquema explícito. Sin él,
spark.read.jsonhace una pasada completa solo para inferir tipos. Con él, la lectura es una pasada y los tipos son los que queremos (fecha_mscomoLongType, nodouble). explodeconvierte el arraylineasen filas, que es lo que hacía el buclefor ln in ev["datos"]["lineas"]del mapper.- Un solo shuffle. El
groupByes la única operación ancha.filter,select,explodey el join broadcast se encadenan en el mismo stage que la lectura. Catalyst inserta una agregación parcial antes del shuffle (HashAggregateconpartial_sum), así que lo que cruza la red son parciales por(dia, mercado, productor)por partición: el combiner, sin escribirlo. cache()enlineasporque hay dos acciones (control.first()y elwrite); sin él, el JSON se leería y parsearía dos veces.partitionBy("dia")conpartitionOverwriteMode=dynamic. La salida es un directorio por día (dia=2026-09-14/), yoverwriteen modo dinámico reemplaza solo los días que este job escribe, dejando intactos los demás. Ejecutar dos veces el job del 14 de septiembre produce exactamente la misma salida: es la idempotencia por partición que el pipeline de 05-05 necesita para reprocesar y hacer backfill.coalesce(1)antes de escribir, porque el agregado son unas 12 filas por día y no queremos 8 ficheros Parquet de 2 KB. Con agregados grandes se quitaría o se ajustaría.
9.3 Lanzamiento y plan de ejecución
En local, con cuatro hilos como executors simulados (local[4]), sobre el fichero generado en 05-01:
$ pip install pyspark==3.5.1
$ spark-submit --master 'local[4]' servicios/analitica/ventas_diarias.py \
eventos/2026-09-14/pedidos.jsonl servicios/analitica/catalogo.csv salida/ventas_diarias
Control: {'pedidos': 400000, 'total': 3590252.0}
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- Project [dia, mercado, productor, nombre_productor, provincia, importe, unidades, pedidos]
+- BroadcastHashJoin [productor], [productor], LeftOuter, BuildRight
:- HashAggregate(keys=[dia, mercado, productor], functions=[sum(importe), sum(cantidad), count(distinct pedido_id)])
: +- Exchange hashpartitioning(dia, mercado, productor, 8)
: +- HashAggregate(keys=[dia, mercado, productor], functions=[partial_sum(importe), partial_sum(cantidad), ...])
: +- InMemoryTableScan [dia, mercado, pedido_id, productor, cantidad, importe]
: +- InMemoryRelation ... (lineas cacheado)
: +- Generate explode(datos.lineas) ...
: +- Filter (tipo = pedido.creado)
: +- FileScan json [tipo, fecha_ms, datos] PushedFilters: [IsNotNull(tipo), EqualTo(tipo,pedido.creado)]
+- BroadcastExchange HashedRelationBroadcastMode
+- FileScan csv [productor, nombre_productor, provincia]
+----------+---------+-------------------+------------------+---------+----------+--------+-------+
|dia |mercado |productor |nombre_productor |provincia|importe |unidades|pedidos|
+----------+---------+-------------------+------------------+---------+----------+--------+-------+
|2026-09-14|girona |bodega-roble-alto |Bodega Roble Alto |Lleida |246187.40 |25121 |23967 |
|2026-09-14|girona |huerta-la-vega |Huerta La Vega |Girona |178013.70 |58940 |31502 |
|2026-09-14|girona |queseria-montblanc |Quesería Montblanc|Tarragona|473402.10 |46330 |41218 |
...El plan se lee de abajo arriba, y cada línea confirma una decisión de los apartados anteriores: FileScan json con PushedFilters (el filtro empujado a la lectura), Generate explode, InMemoryRelation (la cache), el HashAggregate con partial_sum antes del Exchange hashpartitioning(..., 8) (agregación parcial, después el único shuffle con 8 particiones) y el HashAggregate final, y BroadcastHashJoin con BroadcastExchange del CSV (el catálogo viaja, las ventas no). Si en lugar de BroadcastHashJoin apareciera SortMergeJoin con dos Exchange, sabríamos que el broadcast no se aplicó y que estamos pagando un shuffle de más.
Contra el clúster standalone de Compose, con el fichero en HDFS:
$ docker compose exec spark-master spark-submit \
--master spark://spark-master:7077 \
--executor-memory 1G --executor-cores 2 --num-executors 2 \
/app/analitica/ventas_diarias.py \
hdfs://namenode:8020/km0/eventos/2026-09-14/pedidos.jsonl \
/app/analitica/catalogo.csv \
hdfs://namenode:8020/km0/agregados/ventas_diarias
$ docker compose exec namenode hdfs dfs -ls /km0/agregados/ventas_diarias/
drwxr-xr-x - spark supergroup 0 /km0/agregados/ventas_diarias/dia=2026-09-14
-rw-r--r-- 2 spark supergroup 0 /km0/agregados/ventas_diarias/_SUCCESSCon --master yarn y la configuración de Hadoop en HADOOP_CONF_DIR, el mismo fichero se lanzaría sobre el YARN de 05-02, con el driver como ApplicationMaster (--deploy-mode cluster), sin cambiar una línea de código. En la interfaz del maestro (localhost:8080) se ven los workers y las aplicaciones; en la del driver (localhost:4040, mientras corre) los jobs, stages y tasks, con sus bytes de shuffle.
9.4 La versión RDD, para comparar
ventas_rdd.py contiene la función del apartado 3. Lanzado con --rdd, produce las mismas cifras, y en la interfaz se ven dos diferencias: el stage de lectura es más lento (cada línea pasa por un proceso Python con json.loads, en lugar del parser JSON de la JVM) y no hay plan que leer: Spark ejecuta las lambdas tal cual, sin empujar filtros ni podar columnas, porque no sabe qué hacen. En el fichero de 130 MB la diferencia es de 9 s frente a 4 s en local[4]; en terabytes, de horas. El RDD sigue siendo la herramienta correcta cuando los datos no tienen esquema o la lógica no cabe en expresiones de columna, y es lo que hay debajo de todo DataFrame; pero para analítica sobre eventos con esquema, DataFrames es la API.
Errores Comunes y Consejos
collect()sobre un DataFrame grande. Trae todas las filas al driver, que tiene unos GB de memoria:OutOfMemoryErroren el driver. Para mirar,show()otake(n); para guardar,write.- Olvidar que las transformaciones son perezosas. Un
filtermal escrito no falla al escribirlo, sino en la primera acción, con una traza que apunta alwrite. Y unprintdentro de una lambda no aparece en el driver: se ejecuta en los executors (mira sus logs). groupByKey+ suma en RDD. Mueve todos los valores por la red.reduceByKeyoaggregateByKeycombinan localmente. En DataFrames,groupBy().agg()ya lo hace.- UDF de Python por fila. Una
udfen PySpark serializa cada fila hacia un proceso Python y de vuelta; anula Catalyst y Tungsten. Casi todo se puede expresar conF.*; si no, usarpandas_udf(vectorizada por lotes con Arrow). - Cachear sin usar o no liberar.
cache()sobre algo que se usa una vez es un coste sin beneficio; cachear muchas cosas sinunpersist()desborda la memoria de los executors y provoca desalojos y recomputaciones silenciosas. - 200 particiones de shuffle para 10 MB. El valor por defecto de
spark.sql.shuffle.partitionses para clústeres grandes. Ajústalo o confía en AQE; ycoalesceantes de escribir para no producir cientos de ficheros de un kilobyte que ahogarán el NameNode (04-02). - Confiar en el broadcast automático con un CSV. Sin estadísticas, Spark puede no estimar el tamaño y hacer un sort-merge join.
F.broadcast()explícito y comprobar enexplain(). - Inferir el esquema del JSON en producción. Es una pasada extra y un esquema que cambia con los datos (un día sin
versiony la columna desaparece). Esquema explícito, versionado con el contrato del evento (02-05). - Escribir con
overwritesinpartitionOverwriteMode=dynamic. Borra todo el directorio de salida, incluidos los días que no se estaban recalculando. Es el error más caro que se comete con un pipeline de backfill.
Ejercicios
Ejercicio 1: Ranking por mercado en un solo programa
Amplía ventas_dataframe para que, además de las ventas, produzca para cada (dia, mercado) el productor con más ventas y su cuota sobre el total del mercado (porcentaje), usando funciones de ventana (pyspark.sql.Window). ¿Cuántos shuffles añade tu solución? Compáralo con los dos jobs del ejercicio 1 de 05-02.
Ejercicio 2: Leer un plan
Este es el plan de una versión modificada del programa, escrita por un compañero. Identifica cuatro problemas de rendimiento a partir del plan y di cómo corregir cada uno.
== Physical Plan ==
+- SortMergeJoin [productor], [productor], LeftOuter
:- Sort [productor ASC]
: +- Exchange hashpartitioning(productor, 200)
: +- HashAggregate(keys=[dia, mercado, productor], functions=[sum(importe)])
: +- Exchange hashpartitioning(dia, mercado, productor, 200)
: +- HashAggregate(keys=[dia, mercado, productor], functions=[partial_sum(importe)])
: +- BatchEvalPython [calcular_importe(cantidad, precio)]
: +- Generate explode(datos.lineas)
: +- Filter (tipo = pedido.creado)
: +- FileScan json [id_evento, tipo, version, fecha_ms, origen, datos]
+- Sort [productor ASC]
+- Exchange hashpartitioning(productor, 200)
+- FileScan csv [productor, nombre_productor, provincia]Ejercicio 3: Reprocesar la Semana de la Vendimia
Los eventos del 8 al 14 de septiembre de 2026 (la Semana de la Vendimia) tenían un error en el precio de vino-crianza que pedidos ha corregido regenerando los ficheros pedidos.jsonl de esos siete días en HDFS. Escribe la invocación (o las invocaciones) de ventas_diarias.py que recalcule exactamente esos siete días sin tocar los demás, y explica qué combinación de opciones del programa garantiza que (a) no queden datos antiguos de esos días, (b) no se borren los otros días, y (c) ejecutarlo dos veces dé el mismo resultado. ¿Qué cambiarías si la entrada fuera un glob pedidos-*.jsonl con varios ficheros por día?
Soluciones
Ejercicio 1.
from pyspark.sql import Window
por_mercado = Window.partitionBy("dia", "mercado")
ranking = (ventas
.withColumn("total_mercado", F.sum("importe").over(por_mercado))
.withColumn("cuota", F.round(F.col("importe") / F.col("total_mercado") * 100, 1))
.withColumn("puesto", F.row_number().over(por_mercado.orderBy(F.desc("importe"))))
.filter(F.col("puesto") == 1)
.select("dia", "mercado", "productor", "importe", "cuota"))Las dos ventanas comparten la partición (dia, mercado), así que Spark añade un shuffle (Exchange hashpartitioning(dia, mercado)) y un Sort dentro de cada partición para row_number; con AQE y las mismas claves, a veces reutiliza el intercambio. Total: dos shuffles en el programa (el groupBy y la ventana), en un solo job con tres stages, frente a los dos jobs de MapReduce con sus cuatro pasos por HDFS. Y el join con el catálogo sigue sin costar shuffle porque ventas ya lo traía resuelto.
Ejercicio 2.
SortMergeJoincon dosExchangeporproductoren lugar deBroadcastHashJoin: el catálogo de tres filas está provocando un shuffle de las ventas y una ordenación. Corrección:F.broadcast(catalogo).BatchEvalPython [calcular_importe]: una UDF de Python por fila, que además impide que elpartial_sumse compute en la JVM sin salir a Python. Corrección:F.col("ln.cantidad") * F.col("ln.precio").Exchange ... 200: 200 particiones de shuffle para un agregado de decenas de filas; 200 tasks minúsculas por stage. Corrección:spark.sql.shuffle.partitionsa 8 (o AQE con coalescencia activada).FileScan json [id_evento, tipo, version, fecha_ms, origen, datos]sinPushedFiltersni poda: se leen todas las columnas, y probablemente sin esquema explícito (inferencia, una pasada extra). Corrección: esquema declarado yselectde las columnas necesarias justo tras la lectura; el filtrotipodebería aparecer comoPushedFilters(lo hace cuando la columna se compara con un literal y el esquema es conocido).
Un quinto detalle: no hay InMemoryRelation, así que si el programa hace más de una acción, releerá el JSON cada vez.
Ejercicio 3.
Una sola invocación con un glob para los siete días, o siete invocaciones (una por día, que es lo que hará Airflow en 05-05 con el backfill):
spark-submit --master spark://spark-master:7077 /app/analitica/ventas_diarias.py \
'hdfs://namenode:8020/km0/eventos/2026-09-{08,09,10,11,12,13,14}/pedidos.jsonl' \
/app/analitica/catalogo.csv hdfs://namenode:8020/km0/agregados/ventas_diarias(a) mode("overwrite") reemplaza el contenido de cada partición dia=2026-09-NN/ que el job escribe, de forma atómica por partición (escritura en temporal y commit). (b) partitionOverwriteMode=dynamic limita el borrado a las particiones presentes en la salida del job; sin él, overwrite vaciaría ventas_diarias/ entero, incluidos agosto y el resto de septiembre. (c) El cálculo es determinista (misma entrada, mismo agregado) y la escritura reemplaza la partición completa: dos ejecuciones dejan los mismos bytes (salvo nombres de fichero internos), sin duplicar ni acumular. Es importante que dia se derive de fecha_ms (tiempo de evento) y no del nombre del directorio: si un evento del 14 llegara tarde y pedidos lo hubiera dejado en el fichero del 15, el job del 15 lo escribiría en dia=2026-09-14/, sobrescribiendo la partición del 14 con solo ese evento. Para evitarlo hay dos opciones: filtrar en el job por el rango de días que se está procesando (F.col("dia").between(...)) y descartar o registrar los fuera de rango, o que la partición de entrada y salida coincidan por construcción. Con un glob pedidos-*.jsonl por día no cambia nada en el programa (spark.read.json acepta globs y directorios); solo conviene comprobar que ninguno de los ficheros esté aún en escritura, que es exactamente lo que el sensor de 05-05 vigilará con _SUCCESS o con un fichero de cierre.
Conclusión
Spark toma el modelo de MapReduce (particiones, tareas reejecutables, shuffle por clave) y le quita lo que lo hacía lento: el cálculo entero es un DAG de operadores que el driver conoce completo, corta en stages solo donde hay una dependencia ancha, encadena todo lo demás en la misma task, y mantiene los datos intermedios en la memoria de executors que viven toda la aplicación. La tolerancia a fallos viene del linaje, que recomputa la partición perdida en lugar de haberla escrito a disco. Sobre los RDD, inmutables y perezosos, los DataFrames añaden esquema y un optimizador, Catalyst, que empuja filtros, poda columnas, inserta agregaciones parciales y elige broadcast o sort-merge para cada join; Tungsten compila el plan, y Parquet hace que leer dos columnas de un año de eventos cueste lo que ocupan esas dos columnas. Hemos aprendido a leer un explain() para verificar que el plan hace lo que creemos, a usar cache cuando hay varias acciones, broadcast para el catálogo, coalesce antes de escribir, la sal para Quesería Montblanc, y partitionBy con sobreescritura dinámica para que el job de ventas diarias sea idempotente por día. servicios/analitica/ventas_diarias.py calcula ahora en seis segundos, y en un solo programa, lo que en 05-02 eran tres jobs y un minuto; y ALS en MLlib ha convertido los clics del lago en recomendaciones sin que las cien iteraciones supongan cien lecturas.
Todo lo que hemos hecho parte de un conjunto acotado: el fichero del 14 de septiembre, ya cerrado, o los clics de siete días. Pero pedidos.eventos no se cierra nunca: los eventos stock.actualizado llegan a razón de cientos por segundo, y las posiciones de los repartidores a 2,4 millones al día, y ni el panel de reparto ni la alerta de stock bajo pueden esperar al lote de la noche. La siguiente lección trata el procesamiento de flujos: qué cambia cuando el conjunto no termina, cómo se define "los últimos cinco minutos" cuando los eventos llegan desordenados (las marcas de agua), y cómo Flink y Spark Structured Streaming ejecutan el mismo DAG de esta lección sobre un flujo que no se detiene.
Curso de Arquitecturas Distribuidas
Módulo 1: Introducción a los Sistemas Distribuidos
- Conceptos Básicos de Sistemas Distribuidos
- Modelos de Sistemas Distribuidos
- Ventajas y Desafíos de los Sistemas Distribuidos
- Las Falacias de la Computación Distribuida
- Tiempo, Relojes y Ordenación de Eventos
- Del Monolito a la Plataforma Distribuida: el Caso Kilómetro Cero
Módulo 2: Comunicación en Sistemas Distribuidos
- Protocolos de Comunicación
- RPC y RMI
- gRPC y Serialización de Datos
- Mensajería y Colas de Mensajes
- Patrones de Comunicación Asíncrona
Módulo 3: Consistencia y Replicación
- Modelos de Consistencia
- El Teorema CAP y PACELC
- Algoritmos de Consenso
- Replicación de Datos
- Transacciones Distribuidas y Sagas
Módulo 4: Almacenamiento Distribuido
- Particionado de Datos y Hashing Consistente
- Sistemas de Archivos Distribuidos
- Almacenamiento de Objetos
- Bases de Datos Distribuidas
- Cachés Distribuidos
Módulo 5: Computación Distribuida
- Modelos de Computación Distribuida
- MapReduce y Hadoop
- Spark y Computación en Memoria
- Procesamiento de Flujos de Datos
- Planificación de Trabajos y Pipelines de Datos
Módulo 6: Seguridad en Sistemas Distribuidos
- Autenticación y Autorización
- Cifrado y Protección de Datos
- Gestión de Identidades
- Seguridad entre Servicios: mTLS y Gestión de Secretos
- Puertas de Enlace, Limitación de Tasa y Auditoría
Módulo 7: Monitoreo y Mantenimiento
- Monitoreo de Sistemas Distribuidos
- Logs Centralizados y Trazabilidad Distribuida
- Gestión de Fallos y Recuperación
- Patrones de Resiliencia: Timeouts, Reintentos y Circuit Breaker
- Automatización y Orquestación
- Pruebas en Sistemas Distribuidos e Ingeniería del Caos
