lección 4
Transformaciones y acciones: por qué Spark es "perezoso" y es una ventaja
Lazy evaluation: las transformaciones no hacen nada hasta que una acción las ejecuta. Como escribir una lista de la compra antes de ir al supermercado.
⏱ 55 min
### La lista de la compra y el supermercado
Imagina que estás preparando una cena especial. Podrías ir al supermercado cada vez que necesitas un ingrediente: primero un viaje para los tomates, luego otro para la cebolla, luego otro para el aceite. Cada viaje tiene un coste fijo (conducir, aparcar, hacer cola). Es ineficiente. ¿Qué haces en la vida real? Escribes una lista de la compra COMPLETA y haces UN solo viaje al supermercado. Es exactamente lo que hace Spark.
Cuando escribes transformaciones en Spark (filter, select, groupBy, join...), Spark NO las ejecuta inmediatamente. Las va anotando en un plan, como tú anotas ingredientes en la lista. Solo cuando le pides un resultado concreto (show, count, write — llamadas "acciones"), Spark mira el plan completo, lo optimiza, y lo ejecuta todo de una vez de la forma más eficiente posible.
Este patrón se llama "lazy evaluation" (evaluación perezosa) y es una de las razones por las que Spark es tan rápido. Le permite optimizar el plan completo antes de ejecutar nada, en vez de ejecutar cada paso de forma aislada.
### Transformaciones: describir QUÉ hacer (sin hacerlo)
Las transformaciones son operaciones que DEFINEN una nueva forma de los datos pero no los mueven ni los calculan. Son instrucciones que Spark anota para ejecutar después. Las más comunes:
- select() — elegir columnas
- filter() / where() — filtrar filas
- withColumn() — crear o modificar una columna
- groupBy() — agrupar (necesita una agregación después)
- orderBy() / sort() — ordenar
- join() — combinar DataFrames
- distinct() — eliminar duplicados
- drop() — eliminar columnas
- withColumnRenamed() — renombrar columna
1from pyspark.sql import SparkSession2from pyspark.sql.functions import col, sum as spark_sum, avg34spark = SparkSession.builder.appName("lazy").master("spark://spark-master:7077").getOrCreate()5ventas = spark.read.option("header", "true").option("inferSchema", "true").csv("/home/jovyan/data/ventas.csv")67# NINGUNA de estas líneas ejecuta nada todavía8# Spark solo está construyendo un PLAN9ventas_madrid = ventas.filter(col("ciudad") == "Madrid") # transformación10ventas_caras = ventas_madrid.filter(col("precio_unitario") > 100) # transformación11resumen = ventas_caras.groupBy("producto").agg( # transformación12 spark_sum("cantidad").alias("total_unidades"),13 avg("precio_unitario").alias("precio_medio")14)15resultado = resumen.orderBy("total_unidades", ascending=False) # transformación1617print("Hasta aquí, Spark NO ha movido un solo byte de datos.")18print("Solo ha construido un plan de ejecución (DAG).")19print(f"El plan tiene {resultado.rdd.getNumPartitions()} particiones previstas.")2021# AHORA sí — una ACCIÓN dispara la ejecución de todo el plan22resultado.show(10) # ← ACCIÓN: Spark ejecuta todo el plan de golpe
Las 4 transformaciones se acumulan. Solo show() (la acción) dispara la ejecución real.
### Acciones: pedir el resultado (ejecutar el plan)
Las acciones son operaciones que requieren un resultado real — datos que devolver al driver o que escribir en disco. Cuando Spark encuentra una acción, ejecuta todas las transformaciones acumuladas:
- show(n) — muestra n filas en consola
- count() — cuenta el total de filas
- collect() — trae TODAS las filas al driver (¡cuidado con datasets grandes!)
- first() / head() — trae la primera fila
- take(n) — trae n filas como lista
- write.parquet() / write.csv() — escribe a disco
- toPandas() — convierte a DataFrame de Pandas (¡solo si cabe en memoria!)
NUNCA uses collect() con un DataFrame grande. Trae TODAS las filas al driver (una sola máquina), lo que puede causar un OutOfMemoryError. Es el equivalente a volver al problema original de Pandas. Usa show(), take() o write() para manejar resultados grandes.
### ¿Por qué la pereza es una ventaja?
La lazy evaluation permite a Spark optimizar el plan completo antes de ejecutarlo. El optimizador Catalyst analiza todas tus transformaciones juntas y puede:
- 01.Predicate pushdown: si filtras por ciudad="Madrid" y luego lees un Parquet particionado por ciudad, Spark solo lee la partición de Madrid. Sin lazy evaluation, leería todo y luego filtraría.
- 02.Projection pushdown: si seleccionas 3 columnas de una tabla con 50, Spark solo lee esas 3 del Parquet. No carga las 50 para luego descartar 47.
- 03.Reordenar operaciones: Spark puede mover filtros antes de joins para reducir el volumen de datos que se mueven entre nodos.
- 04.Eliminar operaciones redundantes: si haces .filter(x > 5).filter(x > 10), Spark sabe que solo necesita x > 10.
- 05.Elegir algoritmos de join: según el tamaño de las tablas, Spark elige broadcast join, sort-merge join o shuffle hash join.
1# Ejemplo: Spark optimiza automáticamente el orden de operaciones23# Tú escribes esto (join primero, luego filtro):4resultado = (5 ventas6 .join(clientes, ventas.cliente_id == clientes.id) # JOIN de 1M x 50K filas7 .filter(col("ciudad") == "Madrid") # filtro después8 .select("nombre", "producto", "precio_unitario")9)1011# Pero Spark EJECUTA esto (filtro primero, luego join):12# 1. Filtra ventas por Madrid (de 1M a ~100K filas)13# 2. Hace el JOIN con el dataset reducido (100K x 50K)14# 3. Selecciona solo las 3 columnas necesarias15#16# El resultado es idéntico, pero el rendimiento es radicalmente mejor.17# Esto es posible SOLO porque Spark ve el plan completo antes de ejecutar.1819# Para ver el plan optimizado:20resultado.explain(True)
Spark reordena operaciones para máxima eficiencia — gracias a la lazy evaluation
### El DAG: el plan de ejecución visual
Spark representa el plan de ejecución como un DAG (Directed Acyclic Graph): un grafo donde cada nodo es una operación y las flechas indican dependencias de datos. Cuando llamas a una acción, Spark divide el DAG en "stages" (etapas) separadas por shuffles, y cada stage en "tasks" (una por partición).
### explain(): ver el plan antes de ejecutar
Puedes pedir a Spark que te muestre su plan de ejecución SIN ejecutarlo. Es como pedir el presupuesto antes de aprobar la obra. Usa explain() con el parámetro True para ver el plan completo (parsed, analyzed, optimized y physical):
1# Ver el plan de ejecución completo2resultado.explain(True)34# == Parsed Logical Plan ==5# Sort [total_unidades DESC]6# +- Aggregate [producto], [producto, sum(cantidad) AS total_unidades, avg(precio_unitario) AS precio_medio]7# +- Filter (precio_unitario > 100)8# +- Filter (ciudad = Madrid)9# +- Relation [id,fecha,producto,...] csv1011# == Optimized Logical Plan ==12# Sort [total_unidades DESC]13# +- Aggregate [producto], [...]14# +- Filter ((ciudad = Madrid) AND (precio_unitario > 100)) ← ¡combinó los 2 filtros!15# +- Relation [producto,cantidad,precio_unitario,ciudad] csv ← ¡solo lee 4 columnas!1617# == Physical Plan ==18# AdaptiveSparkPlan19# +- Sort [total_unidades DESC]20# +- HashAggregate [producto], [sum, avg]21# +- Exchange hashpartitioning(producto, 200) ← el SHUFFLE22# +- HashAggregate [producto], [partial_sum, partial_avg]23# +- Filter ((ciudad = Madrid) AND (precio_unitario > 100))24# +- FileScan csv [producto,cantidad,precio_unitario,ciudad]
explain(True) revela cómo Spark optimizó tu consulta
Consejo de senior: SIEMPRE mira el explain() de tus queries complejas antes de ejecutarlas en producción. Es gratis (no ejecuta nada) y te dice si Spark está haciendo predicate pushdown, si está haciendo un broadcast join o un shuffle join costoso, y cuántas particiones va a crear. Es tu herramienta de diagnóstico número 1.
### Inmutabilidad: cada transformación crea un DataFrame nuevo
Otro concepto crucial: los DataFrames de Spark son inmutables. Cada transformación devuelve un DataFrame NUEVO sin modificar el original. Esto es diferente a Pandas, donde puedes modificar un DataFrame in-place. La inmutabilidad permite a Spark rastrear la lineage (linaje) de cada dato — de dónde vino y qué transformaciones pasó.
1# Cada transformación crea un nuevo DataFrame2df1 = ventas.filter(col("ciudad") == "Madrid") # df1 es NUEVO, ventas no cambia3df2 = df1.filter(col("precio_unitario") > 100) # df2 es NUEVO, df1 no cambia4df3 = df2.select("producto", "precio_unitario") # df3 es NUEVO, df2 no cambia56# ventas sigue teniendo todas las ciudades y todas las columnas7# Esto permite reutilizar DataFrames intermedios sin miedo8ventas_barcelona = ventas.filter(col("ciudad") == "Barcelona") # reutiliza 'ventas'
Inmutabilidad: cada transformación devuelve un nuevo DataFrame sin alterar el original
Esta inmutabilidad, combinada con la lazy evaluation, es lo que permite a Spark hacer optimizaciones agresivas y recuperarse de fallos. Si un nodo muere a mitad de un cómputo, Spark sabe exactamente qué transformaciones aplicar para recalcular las particiones perdidas, porque tiene el plan completo registrado.
### La otra cara de la pereza: cache()
Como el DataFrame no guarda los datos — solo la receta para calcularlos — cada acción vuelve a ejecutar la cadena entera desde el principio. Si llamas a .count() y luego a .show() sobre el mismo pipeline, Spark lo calcula dos veces. Con datos pequeños no lo notas; con un millón de filas y un groupBy de por medio, pagas el doble.
1# Sin cache: cada acción recalcula todo el pipeline2top = (3 ventas4 .filter(col("ciudad") == "Barcelona")5 .groupBy("producto")6 .agg(spark_sum("cantidad").alias("total"))7)89top.count() # 1ª acción: lee el CSV, filtra, agrupa, cuenta10top.show() # 2ª acción: vuelve a leer, filtrar y agrupar desde cero1112# Con cache: la primera acción calcula Y guarda en memoria13top.cache() # marca para cachear (no ejecuta nada todavía)14top.count() # 1ª acción: calcula Y guarda en memoria15top.show() # 2ª acción: lee de memoria — mucho más rápido16top.unpersist() # cuando ya no lo necesites, libéralo
cache() evita recalcular un pipeline cuando necesitas llamar varias acciones sobre el mismo resultado
Regla práctica: si vas a llamar dos o más acciones sobre el mismo DataFrame transformado, pon un .cache() antes de la primera. Y acuérdate de .unpersist() cuando acabes, porque la memoria del executor es limitada y lo que cacheas le quita sitio al shuffle.
## ejercicios
Clasificar operaciones: ¿transformación o acción?
Dada una lista de operaciones Spark, clasifica cada una como "transformación" (lazy, no ejecuta nada) o "acción" (dispara ejecución). Completa el diccionario.
💡 Resultado esperado
🔵 filter() → transformación 🔵 select() → transformación 🔵 withColumn() → transformación 🟠 count() → acción
Construir un pipeline lazy
El CEO quiere saber el top 5 de productos más vendidos (por cantidad) en Barcelona durante 2024. Encadena transformaciones sin ejecutar nada hasta el final.
Leer un plan de ejecución
Usa explain(True) para ver el plan optimizado del pipeline anterior. Identifica: ¿Spark combinó los dos filtros? ¿Cuántas columnas lee realmente del CSV? Escribe un comentario explicando lo que ves.
Regístrate para guardar tu progreso.
## comentarios
Reporta erratas, ayuda a otros o comparte tu opinión. Sé constructivo.
Inicia sesión para comentar y responder.
cargando comentarios...