Saltar al contenido

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 SparkSession
2from pyspark.sql.functions import col, sum as spark_sum, avg
3
4spark = 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")
6
7# NINGUNA de estas líneas ejecuta nada todavía
8# Spark solo está construyendo un PLAN
9ventas_madrid = ventas.filter(col("ciudad") == "Madrid") # transformación
10ventas_caras = ventas_madrid.filter(col("precio_unitario") > 100) # transformación
11resumen = ventas_caras.groupBy("producto").agg( # transformación
12 spark_sum("cantidad").alias("total_unidades"),
13 avg("precio_unitario").alias("precio_medio")
14)
15resultado = resumen.orderBy("total_unidades", ascending=False) # transformación
16
17print("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.")
20
21# AHORA sí — una ACCIÓN dispara la ejecución de todo el plan
22resultado.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:

  1. 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.
  2. 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.
  3. 03.Reordenar operaciones: Spark puede mover filtros antes de joins para reducir el volumen de datos que se mueven entre nodos.
  4. 04.Eliminar operaciones redundantes: si haces .filter(x > 5).filter(x > 10), Spark sabe que solo necesita x > 10.
  5. 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 operaciones
2
3# Tú escribes esto (join primero, luego filtro):
4resultado = (
5 ventas
6 .join(clientes, ventas.cliente_id == clientes.id) # JOIN de 1M x 50K filas
7 .filter(col("ciudad") == "Madrid") # filtro después
8 .select("nombre", "producto", "precio_unitario")
9)
10
11# 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 necesarias
15#
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.
18
19# 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).

El DAG se divide en stages. Cada stage se ejecuta en paralelo sin mover datos entre nodos. El shuffle es el corte.

### 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 completo
2resultado.explain(True)
3
4# == 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,...] csv
10
11# == 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!
16
17# == Physical Plan ==
18# AdaptiveSparkPlan
19# +- Sort [total_unidades DESC]
20# +- HashAggregate [producto], [sum, avg]
21# +- Exchange hashpartitioning(producto, 200) ← el SHUFFLE
22# +- 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 DataFrame
2df1 = ventas.filter(col("ciudad") == "Madrid") # df1 es NUEVO, ventas no cambia
3df2 = df1.filter(col("precio_unitario") > 100) # df2 es NUEVO, df1 no cambia
4df3 = df2.select("producto", "precio_unitario") # df3 es NUEVO, df2 no cambia
5
6# ventas sigue teniendo todas las ciudades y todas las columnas
7# Esto permite reutilizar DataFrames intermedios sin miedo
8ventas_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 pipeline
2top = (
3 ventas
4 .filter(col("ciudad") == "Barcelona")
5 .groupBy("producto")
6 .agg(spark_sum("cantidad").alias("total"))
7)
8
9top.count() # 1ª acción: lee el CSV, filtra, agrupa, cuenta
10top.show() # 2ª acción: vuelve a leer, filtrar y agrupar desde cero
11
12# Con cache: la primera acción calcula Y guarda en memoria
13top.cache() # marca para cachear (no ejecuta nada todavía)
14top.count() # 1ª acción: calcula Y guarda en memoria
15top.show() # 2ª acción: lee de memoria — mucho más rápido
16top.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

[01]

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
Cargando editor...
[02]

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.

Cargando editor...
[03]

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.

Cargando editor...

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