Módulo 2: The Spark Execution Model

Evaluación perezosa: transformations vs actions

Descripción

Esta es la lección más importante de todo el módulo, y probablemente de toda esta guía. Spark divide cada operación que le pides en dos categorías completamente distintas: transformations, que no hacen ningún trabajo real, y actions, que sí lo hacen. Confundir las dos —o no saber, frente a un método nuevo, en cuál categoría cae— es la fuente más común de sorpresas al aprender Spark. Esta lección no te pide que memorices una lista: te da el criterio para reconocer la diferencia tú mismo, y la verifica con evidencia real —el conteo de jobs de Spark, antes y después de cada operación—, no con la palabra de nadie.

Conexión con el módulo. Las lecciones 2 a 4 de este módulo te dieron el "quién" (driver/executors) y el "con qué" (DataFrame API sobre RDDs). Esta lección te da el "cuándo": exactamente en qué momento, dentro de una cadena de código de Spark, ocurre trabajo real. Las lecciones 6 y 7 dividen esta misma idea en dos partes prácticas —construir sin ejecutar, y disparar la ejecución—, y la lección 8 la verifica de punta a punta en un solo proyecto.

Una analogía: la lista de compras y la caja registradora

Imagina que vas a hacer las compras de la semana. Escribes en tu teléfono: "leche". Eso no cuesta nada — es solo texto en una nota. Agregas "huevos". Tampoco cuesta nada. Agregas "pan", "café", "diez naranjas". Puedes seguir agregando artículos, tachar alguno, cambiar de opinión, durante horas si quieres, y tu tarjeta de crédito nunca se entera de nada de eso. La lista, por más larga que sea, es gratis mientras solo exista como lista.

El gasto real ocurre en un momento muy específico y muy distinto: cuando llegas a la caja registradora del supermercado, pones los artículos en la cinta, y la cajera cobra. Ahí —y solo ahí— el dinero sale de tu cuenta. Todo lo que pasó antes (escribir, tachar, reordenar la lista) fue planificación; lo que pasa en la caja es ejecución.

Spark separa el trabajo exactamente así. Las transformations.select(), .filter(), .withColumn(), .join(), .groupBy() (sin .agg() todavía)— son escribir la lista: describen qué quieres, sin costo, sin leer ni un byte de datos reales. Las actions.count(), .show(), .collect(), .write.parquet(...)— son la caja registradora: el momento exacto en que Spark, por fin, lee los datos, ejecuta cada transformation acumulada, y produce un resultado real.

Ejemplo trabajado: contando trabajos, no adivinando

La forma correcta de verificar esta idea no es "confiar en que no tarda nada" — eso es un cronómetro disfrazado, y esta guía no usa cronómetros para medir nada. La forma correcta es contar, con la propia API de Spark, cuántos jobs (trabajos reales, los mismos que aparecen en la pestaña "Jobs" de la Spark UI en localhost:4040) se han disparado en la sesión, antes y después de cada operación.

# lazy_jobcount.py
from pyspark.sql import SparkSession
from pyspark.sql.types import (
    StructType, StructField, StringType, IntegerType, DoubleType, TimestampType,
)
from pyspark.sql.functions import col

spark = SparkSession.builder.appName("kiosko-spark").master("local[*]").getOrCreate()

orders_schema = StructType([
    StructField("order_id", StringType(), False),
    StructField("store_id", StringType(), False),
    StructField("product_id", StringType(), False),
    StructField("quantity", IntegerType(), False),
    StructField("unit_price", DoubleType(), False),
    StructField("order_ts", TimestampType(), False),
])
orders_df = spark.read.csv(
    "orders_2026-08-*.csv", schema=orders_schema, header=True, enforceSchema=False,
)

status = spark.sparkContext.statusTracker()
print(f"Jobs antes de construir la cadena: {len(status.getJobIdsForGroup())}")

selected_df = orders_df.select("order_id", "store_id", "product_id", "quantity", "unit_price")
filtered_df = selected_df.filter(col("store_id") == "S01")
filtered_df = filtered_df.filter(col("quantity") > 1)

print(f"Jobs despues de select()+filter()+filter() (sin accion todavia): {len(status.getJobIdsForGroup())}")
print(f"type(filtered_df) = {type(filtered_df)}")

result_count = filtered_df.count()

print(f"Jobs despues de .count() (una accion): {len(status.getJobIdsForGroup())}")
print(f"filtered_df.count() = {result_count}")

spark.stop()

spark.sparkContext.statusTracker() es la puerta de entrada, desde Python, a la misma información que alimenta la pestaña "Jobs" de la Spark UI — según la documentación oficial de monitoreo de Spark, la Spark UI y su API REST muestran, entre otras cosas, /applications/[app-id]/jobs, "una lista de todos los jobs de una aplicación dada". getJobIdsForGroup() devuelve la lista de IDs de todos los jobs disparados hasta ese momento en la sesión — contar cuántos hay, antes y después de cada línea, es una forma de verificación tan directa como abrir la Spark UI en el navegador, solo que programática y reproducible en un script.

Qué esperar. Al correr python3 lazy_jobcount.py, la salida es exactamente esta (ejecutada en esta corrida):

Jobs antes de construir la cadena: 0
Jobs despues de select()+filter()+filter() (sin accion todavia): 0
type(filtered_df) = <class 'pyspark.sql.classic.dataframe.DataFrame'>
Jobs despues de .count() (una accion): 2
filtered_df.count() = 10

Esto es exactamente la evidencia que esta lección promete. Tres transformations encadenadas —un .select() y dos .filter()— dejan el conteo de jobs en cero, idéntico a como estaba antes de tocar orders_df siquiera. filtered_df es un objeto DataFrame real, completamente construido, con su plan de ejecución ya definido —pero ese plan nunca se ejecutó—. Solo cuando se llama a .count(), una acción, el conteo de jobs sube de 0 a 2, y filtered_df.count() por fin produce el número real: 10 órdenes de S01 con quantity > 1, el mismo resultado exacto que ya verificaste en las lecciones 2 y 4 de este módulo desde otros ángulos.

(¿Por qué 2 jobs y no 1, para una sola acción? Es una observación honesta que merece una respuesta honesta, no una simplificación: desde Spark 3.2, Adaptive Query Execution —AQE, activa por defecto— puede dividir incluso una operación aparentemente simple como .count() en más de un job, porque replanifica partes del trabajo con estadísticas reales de ejecución en vez de solo con estimaciones previas. El módulo 6 de esta guía explica AQE a fondo. Por ahora, lo único que importa —y lo que esta lección demuestra con evidencia real— es la distinción binaria: cero jobs antes de la acción, más de cero después. Exactamente cuántos jobs dispara una acción específica es un detalle de implementación que puede variar; que las transformations nunca disparan ninguno, no varía.)

Diagrama: el momento exacto donde el trabajo real empieza

flowchart LR
    A["orders_df = spark.read.csv(...)\nTRANSFORMATION -- 0 jobs"] --> B["selected_df = orders_df.select(...)\nTRANSFORMATION -- 0 jobs"]
    B --> C["filtered_df = selected_df.filter(...)\nTRANSFORMATION -- 0 jobs"]
    C --> D["filtered_df = filtered_df.filter(...)\nTRANSFORMATION -- 0 jobs"]
    D -.->|"hasta aqui: SOLO PLAN,\nningun byte leido"| E(("ACCION:\n.count()"))
    E -->|"AHORA SI:\nlee, filtra, cuenta de verdad"| F["result = 10\njobs > 0"]

    style E fill:#f96,stroke:#333,stroke-width:3px

Profundización: la cita oficial, y cómo reconocer una transformation de una action sin memorizar listas

La RDD Programming Guide —la misma fuente que ya citaste en la lección 3 para RDDs— describe la evaluación perezosa con una precisión que aplica igual de bien a DataFrames: "All transformations in Spark are lazy, in that they do not compute their results right away. Instead, they just remember the transformations applied to some base dataset (e.g. a file). The transformations are only computed when an action requires a result to be returned to the driver program" (todas las transformations en Spark son perezosas, en el sentido de que no calculan sus resultados de inmediato; en cambio, solo recuerdan las transformations aplicadas a algún conjunto de datos base —por ejemplo, un archivo—; las transformations solo se calculan cuando una action requiere que un resultado se devuelva al programa driver).

Esa última frase contiene el criterio que necesitas, sin memorizar ninguna lista: una operación es una action si su resultado tiene que volver al driver (un número, una lista de filas, la confirmación de que un archivo se escribió) — y es una transformation si su resultado sigue siendo, en sí mismo, otro DataFrame o RDD que describe más trabajo pendiente. .filter(...) devuelve un DataFrame — sigue siendo un plan, nada volvió al driver todavía. .count() devuelve un entero de Python — algo tuvo que volver al driver, así que algo tuvo que ejecutarse de verdad para producir ese entero.

Con ese criterio, puedes clasificar casi cualquier método nuevo que encuentres en el resto de esta guía sin tener que memorizarlo de antemano:

Método¿Qué devuelve?Transformation o Action
.select(...)otro DataFrameTransformation
.filter(...) / .where(...)otro DataFrameTransformation
.withColumn(...)otro DataFrameTransformation
.join(...)otro DataFrameTransformation
.groupBy(...) (sin .agg())un GroupedData (todavía plan, no resultado)Transformation
.orderBy(...)otro DataFrameTransformation
.count()un int de PythonAction
.show()None (imprime en pantalla, pero ejecuta)Action
.collect()una list de Row en el driverAction
.write.parquet(...)escribe archivos (efecto real)Action
.first() / .take(n)Row / list en el driverAction

Vas a usar esta misma tabla, mentalmente, cada vez que veas un método nuevo el resto de esta guía —incluyendo .explain(), que verás en la lección 6 no está en ninguna de las dos categorías de forma estricta (no transforma el DataFrame ni lo ejecuta; solo lo describe), un caso especial que la siguiente lección aclara con evidencia.

Errores comunes

Esperar que una transformation "tarde" si el DataFrame es grande. Qué pasa: alguien construye una cadena larga de .filter()/.select()/.withColumn() sobre un DataFrame que, en teoría, representa millones de filas (como el kiosko_orders_at_scale del módulo 4), y espera que el script tarde en correr esa parte, aunque no haya ninguna acción todavía. Por qué pasa: la intuición de "más datos, más tiempo" es correcta para cualquier operación que de verdad procese datos — pero una transformation, sin importar cuántas filas represente el DataFrame, no procesa ni una sola. Cómo detectarlo: si tu script tarda segundos o minutos en una línea que solo tiene .filter()/.select(), sin ningún .count()/.show()/.collect() en el medio, algo más está pasando (quizás una acción oculta, como veremos en el próximo error). Cómo corregirlo: recuerda el criterio de esta lección — si el resultado de la línea sigue siendo un DataFrame, es una transformation, y una transformation nunca ejecuta nada sobre los datos reales, sin importar el volumen que el DataFrame represente.

No darse cuenta de que un print(df) normal de Python SÍ puede disparar trabajo. Qué pasa: alguien escribe print(orders_df) (en vez de orders_df.show()) esperando ver una vista previa de los datos, sin saber si eso cuenta como acción. Por qué pasa: print() sobre casi cualquier objeto de Python muestra información sobre su contenido, así que parece razonable esperar lo mismo de un DataFrame. Cómo detectarlo: prueba print(orders_df) tú mismo — la salida es algo como DataFrame[order_id: string, store_id: string, ...], solo el esquema, no los datos. print() sobre un DataFrame llama a su método __repr__, que solo describe la estructura —no dispara ningún job, exactamente como una transformation—. Cómo corregirlo: para ver datos reales, siempre necesitas una acción explícita: .show() para una vista formateada en la consola, .collect() para traer todo a una lista de Python (con el riesgo de memoria que eso implica sobre datos grandes), o .count() para solo el número de filas.

Asumir que dos acciones seguidas reutilizan el trabajo de la primera automáticamente. Qué pasa: alguien llama .count() y después .show() sobre el mismo DataFrame filtrado, y asume que la segunda acción es "gratis" porque el resultado del filtro "ya se calculó" en la primera. Por qué pasa: parece razonable pensar que, una vez que Spark leyó y filtró los datos para el .count(), ese trabajo queda disponible para la siguiente acción sin repetirse. Cómo detectarlo: el ejemplo trabajado de esta lección, y el proyecto de la lección 8, muestran con el mismo conteo de jobs que cada acción dispara su propio trabajo nuevo — el conteo de jobs sigue subiendo con cada acción adicional, no se detiene después de la primera. Cómo corregirlo: por defecto, Spark no recuerda resultados intermedios entre acciones — cada acción vuelve a ejecutar todo el plan desde el origen de los datos, a menos que uses .cache()/.persist() explícitamente para pedirle a Spark que sí guarde un resultado intermedio en memoria. El módulo 6 de esta guía (Catalyst, .explain() y caché) enseña exactamente cuándo esa decisión vale la pena y cuándo no.

Ejercicios

Ejercicio 1 — Clasifica cinco operaciones sin ejecutarlas. Sin correr ningún código, usando solo el criterio de esta lección ("¿el resultado vuelve al driver, o sigue siendo un DataFrame/RDD?"), clasifica estas cinco operaciones como transformation o action: .distinct(), .limit(5), .toPandas(), .printSchema(), .repartition(4).

Ver solución

.distinct()Transformation: devuelve otro DataFrame, sin duplicados, pero sigue siendo un plan pendiente. .limit(5)Transformation: devuelve otro DataFrame limitado a cinco filas, en un plan (aunque en la práctica Spark suele ejecutar algo internamente para optimizar .limit(), conceptualmente sigue devolviendo un DataFrame, no un valor al driver). .toPandas()Action: convierte el DataFrame completo en un pandas.DataFrame que vive en la memoria del driver — un resultado real tuvo que calcularse y traerse de vuelta. .printSchema() — ninguna de las dos, técnicamente: no ejecuta el plan de datos (no lee ningún byte), solo imprime el esquema ya conocido de antemano — es información puramente estructural, disponible sin acción, igual que .explain() en la lección 6. .repartition(4)Transformation: devuelve otro DataFrame con un número distinto de particiones en su plan, pero el repaquetado real solo ocurre cuando una acción posterior lo dispare.

Ejercicio 2 — Verifica tu clasificación con conteo de jobs real. Usando el patrón statusTracker() del ejemplo trabajado de esta lección, verifica con código real si .distinct() y .toPandas() disparan jobs o no, confirmando (o corrigiendo) tu respuesta del ejercicio 1.

Ver solución
status = spark.sparkContext.statusTracker()

print(f"Jobs antes: {len(status.getJobIdsForGroup())}")
distinct_df = orders_df.select("store_id").distinct()
print(f"Jobs despues de .distinct() (sin accion): {len(status.getJobIdsForGroup())}")

pdf = distinct_df.toPandas()
print(f"Jobs despues de .toPandas(): {len(status.getJobIdsForGroup())}")
print(pdf)

Salida esperada (ejecutada en esta corrida; .toPandas() requiere que pandas esté instalado, algo que el módulo 7 de esta guía retoma en profundidad con pandas_udf):

Jobs antes: 0
Jobs despues de .distinct() (sin accion): 0
Jobs despues de .toPandas(): 2
  store_id
0      S02
1      S01
2      S03

Confirmado: .distinct() no mueve el conteo de jobs (transformation), y .toPandas() sí lo hace (action) — exactamente lo que predijo el criterio del ejercicio 1, ahora verificado con evidencia real en vez de solo razonamiento. (El orden de las tres tiendas en el resultado —S02, S01, S03— es el orden interno de procesamiento de las particiones, no un orden garantizado; sin un .orderBy() explícito, nunca asumas ningún orden particular en el resultado de una acción — la lección 7 vuelve sobre este mismo punto con más detalle.)

Ejercicio 3 — Explica, sin código, la trampa de print(df). En 2-3 frases, explica por qué print(orders_df) no muestra los datos reales de Kiosko, aunque print() normalmente muestre el contenido de cualquier otro objeto de Python.

Ver solución

print() sobre cualquier objeto de Python llama internamente a su método __repr__, y el __repr__ de un DataFrame de Spark está diseñado para describir su estructura (el esquema, los nombres y tipos de columna) sin ejecutar el plan de ejecución subyacente — coherente con el hecho de que un DataFrame, hasta que se dispara una acción, no representa datos calculados, solo un plan. Por eso print(orders_df) produce algo como DataFrame[order_id: string, store_id: string, ...] en vez de filas reales: mostrar filas reales requeriría ejecutar el plan, lo cual solo ocurre con una acción explícita como .show().

Resumen y siguiente paso

En esta lección construiste el concepto más importante de este módulo: transformations (.select(), .filter(), .withColumn(), y todo lo que devuelve otro DataFrame) son perezosas y no ejecutan nada, mientras que actions (.count(), .show(), .collect(), y todo lo que devuelve un resultado real al driver) son las que disparan trabajo de verdad. Lo verificaste con el conteo real de jobs de Spark —la misma información que alimenta la pestaña "Jobs" de la Spark UI—, no con intuición: cero jobs después de tres transformations encadenadas, jobs reales solo después de la primera acción.

Antes de avanzar deberías poder: clasificar cualquier método de Spark como transformation o action usando el criterio de "¿el resultado vuelve al driver?"; explicar por qué print(df) no ejecuta nada; y explicar por qué dos acciones seguidas sobre el mismo DataFrame filtrado repiten el trabajo, en vez de reutilizar el resultado de la primera.

Las lecciones 6 y 7 dividen esta misma idea en dos pasos prácticos y más largos: la lección 6 construye una cadena de transformaciones más completa —incluyendo .explain() para inspeccionar el plan sin ejecutarlo— y confirma, de nuevo con evidencia, que nada se disparó; la lección 7 toma esa misma cadena y por fin la ejecuta con una acción real.

Recursos

  • Apache Spark — RDD Programming Guide (la cita exacta sobre evaluación perezosa: "All transformations in Spark are lazy... The transformations are only computed when an action requires a result to be returned to the driver program"). spark.apache.org/docs/latest/rdd-programming-guide.html. En inglés.
  • Apache Spark — Monitoring and Instrumentation (la Spark UI y su API REST, incluyendo /applications/[app-id]/jobs, la fuente de la información que statusTracker() expone desde Python). spark.apache.org/docs/latest/monitoring.html. En inglés.
  • Apache Spark — SQL Performance Tuning (Adaptive Query Execution, activa por defecto desde Spark 3.2, mencionada en esta lección como la razón por la que una sola acción puede disparar más de un job — se profundiza en el módulo 6). spark.apache.org/docs/latest/sql-performance-tuning.html. En inglés.