Sesión 4 — Spark como motor de procesamiento distribuido¶

Esta sesión es el siguiente paso después de S1, S2 y S3. En S1 se estudiaron HDFS, YARN y MapReduce; en S2 se exploraron los Parquet de TPC-DS SF1 como ficheros, sin catálogo; en S3 se comprobó que esos mismos Parquet pueden vivir en un almacenamiento compatible con S3. Ahora se incorpora un motor de procesamiento distribuido capaz de leer esos ficheros y ejecutar el trabajo sobre los DataNodes a través de YARN.

La pregunta central de esta sesión es:

¿Cómo pasa una aplicación Python de leer datos desde una ruta a construir y ejecutar un trabajo distribuido sobre los DataNodes?

Spark no se presentará todavía como un sistema de catálogo ni como una solución de lakehouse. En esta primera etapa es un motor que lee y escribe datos estructurados y que puede ejecutarse localmente o sobre YARN. El recorrido completo hasta este punto es:

S1: HDFS, YARN y MapReduce
  → S2: Parquet en /datalake/raw/tpcds, sin catálogo permanente
  → S3: Parquet sobre un almacenamiento S3-compatible
  → S4: Spark procesa los Parquet y utiliza YARN

Los datos de entrada siguen siendo los de TPC-DS SF1. Se conserva el materializado original de /datalake/raw/tpcds y se evita modificarlo durante los ejercicios: las transformaciones de Spark escribirán sus salidas en rutas de trabajo separadas, bajo /user/luser/s4.

Paquetes de Python de esta sesión. Esta sesión necesita pyspark. Respecto a S3 añade pyspark. Se instalan en el kernel de namenode en la sección «Instalar PySpark en este notebook»: una celda %%writefile requirements.txt crea el fichero dentro del contenedor —el requirements.txt de tu copia de s4/ está en tu equipo y namenode no lo ve— y otra ejecuta %pip install -r requirements.txt. La imagen de namenode ya trae estos paquetes en su entorno Python del curso, /opt/tcdm/venv, que es el intérprete del kernel y donde instala %pip; lo normal es que %pip sólo confirme que están; ejecuta esas celdas de todos modos: dejan explícito qué necesita la sesión y reparan el entorno si falta algo.

Diapositivas de la sesión¶

Las diapositivas se generan con el paquete jupyter-notebook-slide. El alias %%diapositiva permite mantener en español los tipos usados en la sesión y produce la misma salida HTML en Jupyter, Colab y la referencia HTML publicada.

En cada sesión encontrarás varios tipos de diapositivas que te ayudarán a seguir el hilo:

  • Sección (titulo): abre la sesión y resume qué vamos a ver y qué haremos.
  • A continuación (avance): anuncia el contenido de la parte siguiente, para saber en qué punto del programa estamos.
  • Recapitulación (resumen): condensa la parte anterior en puntos clave para recordarlos y revisarlos después.
  • Pregunta guía (pregunta): plantea cuestiones que orientan lo que viene a continuación; conviene intentar responderlas antes de seguir.
  • Evaluación de la sesión (evaluacion): recuerda cómo se evalúa el trabajo de esta sesión.
In [ ]:
%pip install -q "jupyter-notebook-slide @ git+https://github.com/dsevilla/jupyter-notebook-slide.git"
%load_ext notebook_slide

import notebook_slide as jnbs

# Colores de 26-27/teoria/tcdm.css, adaptados al tema de las sesiones.
jnbs.configure(
    background="#eaf2f8",
    foreground="#1f2933",
    border="#c9d6e1",
    heading="#0c304d",
    subheading="#0c304d",
    link="#174f7a",
    code_background="#eaf2f8",
    code_foreground="#0c304d",
    quote_background="#ffffff",
    font_family="Atkinson Hyperlegible, Inter, Aptos, Segoe UI, Helvetica, Arial, sans-serif",
    font_url="https://fonts.googleapis.com/css2?family=Atkinson+Hyperlegible:ital,wght@0,400;0,700;1,400;1,700&display=swap",
)
jnbs.register_slide_type("avance", "A continuación", "#174f7a")
jnbs.register_slide_type("resumen", "Recapitulación", "#a54467")
jnbs.register_slide_type("pregunta", "Pregunta guía", "#0c304d")
jnbs.register_slide_type("evaluacion", "Evaluación de la sesión", "#b3701a")
jnbs.register_slide_type(
    "titulo",
    "Sección",
    "#174f7a",
    layout="title",
    background="linear-gradient(135deg, #0c304d, #174f7a 62%, #a54467)",
    foreground="#ffffff",
    border="transparent",
    heading="#ffffff",
    subheading="#dceaf4",
)
jnbs.register_alias("diapositiva")

Sesión 4: Spark como motor de procesamiento distribuido

Cómo vamos a recorrer la sesión

  • Un motor nuevo sobre los mismos datos: los Parquet de TPC-DS SF1
  • DataFrames: transformaciones perezosas, acciones y planes de ejecución
  • Primero en local[2]; al final, de verdad sobre YARN con spark-submit
  • Las mismas preguntas de negocio de S2, ahora con otro motor

Capas que se deben distinguir¶

En esta sesión aparecen varias palabras que se reutilizarán más adelante y que conviene separar desde el principio:

Concepto Significado en esta sesión
HDFS o S3 Sistema de almacenamiento que contiene los datos
Parquet Formato columnar de los ficheros
Spark Motor que construye y ejecuta el procesamiento
YARN Gestor que asigna recursos a la aplicación Spark
Partición de Spark Parte lógica del trabajo que puede ejecutar una tarea
Partición Hive Organización física de datos por valores de columnas; se verá en S5
Catálogo Servicio que da nombres y metadatos a las tablas; se verá en S6 y S7

Una partición de ejecución de Spark no es un directorio de partición Hive. Tampoco es un bloque HDFS ni un fichero Parquet. Que compartan el nombre no significa que sean la misma unidad; a lo largo de esta sesión se insistirá varias veces en esa distinción.

Objetivos de la sesión¶

Al terminar esta sesión deberías poder:

  • crear una SparkSession y ejecutar una aplicación PySpark;
  • leer Parquet desde una ruta HDFS y examinar su esquema;
  • construir DataFrames con select, filter, withColumn, join y groupBy;
  • distinguir transformaciones de acciones;
  • explicar que Spark construye un plan antes de ejecutar el trabajo;
  • observar el plan lógico y físico con explain;
  • diferenciar el driver, los ejecutores y el gestor de recursos;
  • comparar los modos local y yarn;
  • comprender qué son las particiones de ejecución de Spark;
  • comprobar qué ejecutores remotos y qué DataNodes realizan trabajo;
  • escribir resultados Parquet sin confundir todavía una salida de Spark con una tabla registrada;
  • utilizar vistas temporales para introducir Spark SQL sin un catálogo permanente;
  • conservar consultas reproducibles y explicar la métrica que calcula cada ejemplo (líneas de venta frente a pedidos distintos, por ejemplo).
Evaluación de la sesión

Cómo se evalúa esta sesión

  • Reunión individual de unos 5 minutos con el profesor
  • Se muestra en vivo, en tu ordenador: la aplicación en YARN, el informe de

hosts que procesaron particiones y el resultado de negocio idéntico en

local y sobre YARN

  • Se entrega también una memoria breve (una o dos páginas)
  • No hace falta memorizar órdenes: sí explicar qué papel juega cada pieza

Antes de empezar¶

Dónde se ejecuta este notebook. Igual que en S2 y S3, Jupyter se ejecuta directamente dentro del contenedor namenode, como luser, con acceso de red directo a HDFS y YARN. Las celdas ejecutables hablan directamente con el clúster. Las órdenes que necesitan el Docker del host — arrancar o detener contenedores, make — se muestran como texto para ejecutarlas en una terminal de tu equipo, nunca como celdas de este notebook.

Esta sesión asume que el clúster de S1 ya está en marcha y que S2 generó los Parquet de TPC-DS SF1 en /datalake/raw/tpcds. Si no es así, desde una terminal de tu equipo situada en la raíz de la distribución de sesiones:

Animación: la orden make -C entorno warehouse-up se teclea en una terminal de tu equipo (host)

cd ~/tcdm-public
git pull
make -C entorno warehouse-up

Este notebook se ejecuta en el Jupyter de namenode, que no arranca solo. Como en S2, arráncalo desde otra terminal de tu equipo y déjala abierta mientras trabajes:

Animación: docker exec -it namenode bash se teclea en tu equipo y abre una terminal dentro del contenedor namenode, donde se ejecutan su - luser y jupyter lab

# en tu equipo: entra en namenode
docker exec -it namenode bash
# ya dentro de namenode, como luser: Jupyter con el token fijo «tcdm»
su - luser
jupyter lab --ip=0.0.0.0 --port=8888 --no-browser --IdentityProvider.token=tcdm

La primera orden es la única que se ejecuta en tu equipo; las otras dos, ya dentro del contenedor, arrancan Jupyter con el Python del curso (make -C entorno jupyter hace lo mismo en una sola orden). En Visual Studio Code, conecta este notebook con Select Kernel → Select Another Kernel → Existing Jupyter Server, usando la dirección http://127.0.0.1:8888/lab?token=tcdm y eligiendo el kernel Python 3 (ipykernel). El token es siempre tcdm.

Si algo no funciona —el notebook no conecta, una celda !hdfs no encuentra la orden, un servicio no responde—, consulta «Solución de problemas» en entorno/README.md.

Antes de instalar nada, comprueba que los datos de entrada siguen donde S2 los dejó:

In [ ]:
!hdfs dfs -ls -d /datalake/raw/tpcds

Si el directorio no existe, revisa S2 antes de continuar: esta sesión no regenera TPC-DS.

El progreso y los planes se observan con .show(), .explain() y la CLI de yarn. El ResourceManager (:8088) y el Timeline Server (:8188), ya conocidos desde S1, siguen disponibles para mirar el estado de las aplicaciones YARN en el navegador. A ellos se suma la interfaz web del propio Spark (:4040), que se presenta en cuanto exista la SparkSession.

Instalar PySpark en este notebook¶

El kernel de este notebook instalará PySpark directamente sobre sí mismo, igual que S2 y S3 instalaron PyArrow, Polars o DuckDB. Esta instalación es la que usaremos para todos los ejemplos en modo local del notebook: una única SparkSession construida más abajo y reutilizada en la mayoría de las secciones. La sección "Driver, ejecutores y YARN" no reutiliza esa SparkSession, porque no se puede cambiar el master de una sesión ya construida: lanza un proceso aparte con spark-submit, pero con este mismo PySpark del entorno /opt/tcdm/venv.

La celda siguiente crea requirements.txt con %%writefile dentro de namenode —el de tu copia de la distribución está en tu equipo y el kernel no lo ve— y la posterior lo instala con %pip en el intérprete del kernel. La imagen de namenode ya trae PySpark, así que lo normal es que %pip sólo confirme que está.

In [ ]:
%%writefile requirements.txt
# Dependencias del kernel interactivo de la sesión 4 (Spark en modo local).
# La serie de PySpark aceptada por el curso es >4,<4.2.
# El envío a YARN de la sección "Driver, ejecutores y YARN" usa este mismo
# PySpark (spark-submit de /opt/tcdm/venv).
pyspark>4,<4.2
In [ ]:
%pip install -q -r requirements.txt
In [ ]:
import sys

import pyspark

print(f"PySpark {pyspark.__version__} con Python {sys.version.split()[0]}.")
A continuación

El motor Spark

  • Por qué Spark aparece como respuesta a los encadenamientos MapReduce
  • Driver, ejecutores y gestor de recursos: tres papeles en cualquier modo
  • Job, stage, task y partición: las unidades que verás en los planes

El motor Spark¶

Apache Spark nació en 2009 en el AMPLab de UC Berkeley, motivado por lo costoso que resultaba encadenar trabajos iterativos con el modelo MapReduce que ya conoces de S1: cada paso de un pipeline MapReduce vuelve a leer y escribir HDFS, aunque los datos quepan en memoria. Se hizo open source en 2010 y pasó a la Apache Software Foundation en 2013. Spark generaliza el modelo de MapReduce: sigue dividiendo el trabajo en tareas sobre particiones de datos, pero encadena esas tareas en un grafo de ejecución y puede mantener resultados intermedios en memoria entre una etapa y la siguiente, en lugar de pasar obligatoriamente por HDFS entre cada paso.

Spark ofrece varios módulos sobre un mismo motor central: SQL/DataFrames, streaming estructurado, una biblioteca de aprendizaje automático (MLlib) y procesamiento de grafos. Esta sesión se queda únicamente en SQL/DataFrames, que es el módulo relevante para un data lake de ficheros Parquet.

API estructurada frente a RDDs¶

Spark ofrece dos niveles de API:

  • la API estructurada (DataFrames y, en Scala/Java, Datasets tipados), que describe qué se quiere calcular sobre columnas con nombre y deja que el optimizador Catalyst y el motor de ejecución Tungsten decidan cómo ejecutarlo;
  • la API de bajo nivel (RDDs, Resilient Distributed Datasets), la interfaz original de Spark 1.x: colecciones distribuidas de objetos Python arbitrarios, sin optimizador de por medio.

Los RDDs siguen existiendo y en algún punto muy concreto de esta sesión —al final, en la sección sobre YARN— se usará uno para hacer visible en qué host se ha ejecutado cada tarea, porque es la única forma directa de leer socket.gethostname() dentro de cada partición. Fuera de ese caso muy puntual, todo el notebook usa exclusivamente DataFrames: es la API recomendada, tanto por rendimiento (Catalyst puede reordenar y podar operaciones que RDDs ejecutaría literalmente) como por expresividad (select, filter, groupBy, join, SQL).

Driver, ejecutores y gestor de recursos¶

Tres papeles se repiten en cualquier ejecución de Spark, en modo local o sobre YARN:

  • el driver crea la SparkSession, construye el plan lógico y físico de cada consulta y coordina la ejecución;
  • los ejecutores son los procesos que ejecutan las tareas sobre las particiones de datos y pueden mantener resultados en memoria;
  • el gestor de recursos decide dónde se ejecutan el driver y los ejecutores. En modo local[N] ese gestor no existe realmente: todo vive en un único proceso JVM con N hilos simulando ejecutores. Sobre YARN, el gestor de recursos es el mismo ResourceManager que ya conoces de S1: el driver pide un ApplicationMaster y unos contenedores de ejecutor, y los NodeManagers de los DataNodes los conceden.

Términos que reaparecerán en los planes de ejecución¶

  • Job: conjunto de stages que dispara una acción.
  • Stage: grupo de tareas que se pueden ejecutar sin necesidad de redistribuir datos entre ellas (sin shuffle).
  • Task: unidad de trabajo sobre una única partición.
  • Partición: la porción de un DataFrame que procesa una tarea.
  • Dependencia estrecha (narrow): cada partición de salida depende de una única partición de entrada (select, filter, withColumn).
  • Dependencia ancha (wide): una partición de salida depende de varias particiones de entrada, lo que obliga a un shuffle (groupBy, orderBy, join sin broadcast, distinct).

Crear una SparkSession local¶

Construimos una única SparkSession en modo local[2]: dos hilos que simulan dos ejecutores dentro del mismo proceso Python/JVM del kernel. Esta sesión es la que reutilizaremos en casi todas las secciones siguientes.

La memoria del driver se deja deliberadamente modesta. Esta JVM comparte el contenedor namenode con el propio NameNode, el ResourceManager, el Timeline Server y el kernel de Jupyter (recuerda la tabla de límites de S1: 2 vCPU y 4096 MiB para todo el contenedor). Pedir varios gigabytes de memoria del driver aquí dejaría sin margen a los demonios de Hadoop; para los volúmenes de SF1 que leeremos con select/filter/agregaciones no hace falta.

In [ ]:
from pyspark import SparkContext
from pyspark.sql import SparkSession

spark: SparkSession = (
    SparkSession.builder.appName("TCDM_S4_Local")
    .master("local[2]")
    .config("spark.driver.memory", "1g")
    .config("spark.sql.shuffle.partitions", "8")
    .getOrCreate()
)
sc: SparkContext = spark.sparkContext

print(f"Spark version: {spark.version} | Master: {sc.master} | App: {sc.appName}")

spark.sql.shuffle.partitions controla cuántas particiones produce Spark tras una operación wide (por ejemplo, un groupBy). El valor por defecto es 200, pensado para clústeres grandes; lo bajamos a 8 porque en local[2], con los volúmenes de SF1, 200 particiones diminutas sólo añadirían coordinación sin paralelismo real. Volveremos sobre este parámetro más adelante, en la sección dedicada a la configuración del shuffle.

La interfaz web de Spark¶

Mientras exista una SparkSession, su driver sirve una interfaz web en el puerto 4040. El driver de esta sesión vive en namenode, y compose-hadoop-cluster.yml publica ese puerto en tu equipo: ábrela en http://localhost:4040 y déjala en una pestaña del navegador durante toda la sesión. Se actualiza con cada acción que ejecutes y sigue disponible, con todo lo ejecutado hasta entonces, hasta que se llame a spark.stop().

Pestaña Qué muestra Para qué sirve en esta sesión
Jobs Un job por cada acción (show, count, una escritura…), con su estado, su duración y una línea de tiempo Comprobar que una transformación no lanza nada y que una acción sí
Stages Las etapas de cada job y sus tareas: cuántas son, qué ejecutor hizo cada una, cuánto tardó y cuántos datos leyó o intercambió en el shuffle. «DAG Visualization» dibuja el grafo de dependencias entre las operaciones de la etapa Ver dónde hay un shuffle (dependencia ancha) y cuántas tareas genera cada número de particiones
Storage Los DataFrames persistidos con cache() o persist(), qué fracción está materializada y cuánto ocupa en memoria y en disco Sección de persistencia
Environment La configuración efectiva: propiedades de Spark, versiones de Java y Scala, rutas de bibliotecas Confirmar valores como spark.sql.shuffle.partitions
Executors Los ejecutores registrados, con sus núcleos, su memoria, las tareas completadas y el tiempo de recolección de basura. En local[2] sólo aparece driver Sobre YARN, ver un ejecutor por contenedor y en qué DataNode está cada uno
SQL / DataFrame Una entrada por consulta ejecutada, con el plan físico dibujado como grafo y las métricas reales de cada operador (filas producidas, ficheros y bytes leídos, registros de shuffle). «Details» muestra los planes lógico y físico en texto Es la versión gráfica, y con medidas, de explain("formatted")

explain() muestra el plan antes de ejecutar; la pestaña SQL / DataFrame muestra ese mismo plan después, con lo que ocurrió de verdad en cada operador. Las dos vistas se complementan durante toda la sesión.

Si tienes abierto otro notebook con su propia SparkSession, la segunda interfaz se sirve en el puerto 4041, que no está publicado: cierra antes la otra sesión con spark.stop().

Primera lectura de Parquet¶

Leemos primero una tabla pequeña, date_dim, y después una tabla de hechos, web_sales. HDFS es el sistema de ficheros por defecto de este clúster (igual que en S1 y S2), así que no hace falta ningún prefijo hdfs://: basta la ruta absoluta /datalake/raw/tpcds/<tabla>.

In [ ]:
from pyspark.sql.dataframe import DataFrame

date_dim_path: str = "/datalake/raw/tpcds/date_dim"
date_dim: DataFrame = spark.read.parquet(date_dim_path)

print(f"Ruta física leída: {date_dim_path}")
date_dim.printSchema()

spark.read.parquet(...) no ha leído todavía ninguna fila: sólo ha abierto los metadatos Parquet para conocer el esquema. Es una transformación, perezosa como todas. Las siguientes dos llamadas sí son acciones: obligan a Spark a ejecutar el plan sobre los datos.

In [ ]:
date_dim.show(5)

n_date_dim: int = date_dim.count()
print(f"Filas en date_dim: {n_date_dim}")
print(f"Particiones de ejecución: {date_dim.rdd.getNumPartitions()}")

assert n_date_dim == 73_049

Ahora repetimos la misma secuencia con web_sales, la tabla de hechos que usaremos en la mayor parte de la sesión. Fíjate en que el número de particiones de ejecución no tiene por qué coincidir con el número de ficheros Parquet que escribió S2: Spark decide la partición de lectura a partir del tamaño de los ficheros y de spark.sql.files.maxPartitionBytes, no copiando uno a uno los ficheros físicos.

In [ ]:
web_sales_path: str = "/datalake/raw/tpcds/web_sales"
web_sales: DataFrame = spark.read.parquet(web_sales_path)

print(f"Ruta física leída: {web_sales_path}")
web_sales.printSchema()
In [ ]:
web_sales.show(5)

n_web_sales: int = web_sales.count()
print(f"Filas en web_sales: {n_web_sales}")
print(f"Particiones de ejecución: {web_sales.rdd.getNumPartitions()}")

assert n_web_sales == 719_384

El objetivo de esta lectura no es traer SF1 completo a la memoria del driver: es distribuir la lectura y el cómputo entre particiones que Spark puede procesar en paralelo. web_sales es sólo una parte de los 261 MiB que ocupa todo SF1 en Parquet (su tamaño exacto lo da hdfs dfs -du -h /datalake/raw/tpcds/web_sales); TPC-DS a escalas mayores (SF100, SF1000...) seguiría funcionando con el mismo código precisamente porque el driver nunca materializa la tabla entera, sólo coordina el trabajo sobre sus particiones.

Mira ahora la pestaña Jobs de http://localhost:4040: deben aparecer los jobs que han lanzado los show() y los count() de estas celdas, y ninguno por las dos llamadas a spark.read.parquet(...).

Recapitulación

Spark también lee los Parquet de S2

  • spark.read.parquet(...) sólo abre metadatos: el esquema no es una acción
  • show() y count() son acciones: ahí se ejecuta el plan
  • El número de particiones de ejecución no coincide con el de ficheros:

Spark lo decide por tamaño, no por directorio

Transformaciones, acciones y planificación¶

Encadenamos varias transformaciones sobre web_sales: select, filter y withColumn. Ninguna ejecuta nada todavía.

In [ ]:
from pyspark.sql.functions import col

web_sales_profitable: DataFrame = (
    web_sales.select(
        "ws_order_number",
        "ws_bill_customer_sk",
        "ws_net_profit",
        "ws_ext_sales_price",
    )
    .filter(col("ws_net_profit") > 0)
    .withColumn("margen_sobre_venta", col("ws_net_profit") / col("ws_ext_sales_price"))
)

print(f"Tipo del objeto: {type(web_sales_profitable)}")

web_sales_profitable es un DataFrame con un plan pendiente, no un resultado. La siguiente celda dispara una acción (show) y es el momento en el que Spark realmente lee, filtra y calcula.

In [ ]:
web_sales_profitable.show(5)

Seleccionar pocas columnas frente a select *¶

Parquet es un formato columnar: leer sólo las columnas que hacen falta reduce los bytes que hay que transportar desde el fichero. Comparamos el plan físico de una consulta que sólo toca dos columnas con el de otra que lee todas.

In [ ]:
narrow_query: DataFrame = web_sales.select("ws_order_number", "ws_net_profit").filter(
    col("ws_net_profit") > 0
)
wide_query: DataFrame = web_sales.select("*").filter(col("ws_net_profit") > 0)

print("Plan de la consulta con columnas seleccionadas:")
narrow_query.explain("formatted")
In [ ]:
print("Plan de la consulta con select(\"*\"):")
wide_query.explain("formatted")

Compara la sección ReadSchema de ambos planes: en la consulta estrecha sólo aparecen ws_order_number y ws_net_profit; en la de select("*") aparece el esquema completo de web_sales. El filtro ws_net_profit > 0 también aparece empujado hacia el PushedFilters del Scan parquet en los dos casos: eso es predicate/projection pushdown. El efecto hay que comprobarlo en el plan, no darlo por hecho sólo porque la función se llame filter.

Pregunta guía
  • Las dos consultas devuelven las mismas filas filtradas: ¿qué cambia

entonces entre el plan «estrecho» y el de select("*")?

  • El filtro aparece en PushedFilters: ¿por qué hay que comprobarlo en el

plan en lugar de darlo por supuesto porque la función se llame filter?

Filtrado, ordenación y agrupamiento¶

Contamos, para cada cliente, cuántas líneas de venta con beneficio positivo tiene en web_sales, y mostramos los diez clientes con más líneas.

Antes de agrupar descartamos las líneas sin cliente. TPC-DS deja a propósito algunas claves nulas en las tablas de hechos, para simular datos de origen imperfectos, y groupBy reúne todos los nulos en un único grupo: sin el filtro, ese grupo NULL encabezaría la lista con más líneas que cualquier cliente real, como si fuera uno más. S2 hizo lo mismo con is_valid al contar pedidos por cliente con PyArrow.

Aquí seguimos contando líneas, no pedidos: cada fila de web_sales es una línea de un pedido, y un mismo pedido puede tener varias líneas.

In [ ]:
lineas_con_beneficio: DataFrame = (
    web_sales.filter((col("ws_net_profit") > 0) & col("ws_bill_customer_sk").isNotNull())
    .groupBy("ws_bill_customer_sk")
    .count()
    .orderBy(col("count").desc())
)
lineas_con_beneficio.show(10)

rollup y cube¶

Unimos web_sales con date_dim (por ws_sold_date_sk = d_date_sk) y con item (por ws_item_sk = i_item_sk) para agregar el beneficio neto por año y por categoría de producto.

In [ ]:
item: DataFrame = spark.read.parquet("/datalake/raw/tpcds/item")

web_sales_dated_categorized: DataFrame = (
    web_sales.join(date_dim, web_sales.ws_sold_date_sk == date_dim.d_date_sk)
    .join(item, web_sales.ws_item_sk == item.i_item_sk)
    .select("d_year", "i_category", "ws_net_profit")
)

n_item: int = item.count()
print(f"Filas en item: {n_item}")
assert n_item == 18_000
In [ ]:
from pyspark.sql.functions import sum as spark_sum

rollup_year_category: DataFrame = (
    web_sales_dated_categorized.rollup("d_year", "i_category")
    .agg(spark_sum("ws_net_profit").alias("beneficio_neto"))
    .orderBy("d_year", "i_category")
)
rollup_year_category.show(100)

En el resultado del rollup, un null en i_category con d_year relleno es el subtotal de ese año en todas las categorías; un null en ambas columnas es el total general. rollup añade subtotales siguiendo el orden de las columnas indicadas (primero por d_year, luego por la combinación de ambas), pero no añade el subtotal de cada categoría sumando todos los años.

cube sí genera todas las combinaciones de subtotales, incluida esa:

In [ ]:
cube_year_category: DataFrame = (
    web_sales_dated_categorized.cube("d_year", "i_category")
    .agg(spark_sum("ws_net_profit").alias("beneficio_neto"))
    .orderBy("d_year", "i_category")
)
cube_year_category.show(100)

Comprueba en el resultado del cube que ahora también aparece una fila con d_year a null e i_category con un valor concreto: el beneficio de esa categoría sumando todos los años, algo que rollup no calculaba.

A continuación

El API de DataFrames en acción

  • Joins con las dimensiones: tipos, huérfanos y broadcast join
  • Tipos complejos: un historial de pedidos como array de structs
  • UDFs frente a funciones integradas: qué ve y qué no ve el optimizador
  • Ventanas: máximo por cliente, lag/lead, medias móviles y acumulados

Joins¶

Cargamos customer y customer_address, y preparamos una selección de web_sales con las claves de unión que necesitaremos.

In [ ]:
import pyspark.sql.functions as F

customer: DataFrame = spark.read.parquet("/datalake/raw/tpcds/customer")
customer_address: DataFrame = spark.read.parquet("/datalake/raw/tpcds/customer_address")

n_customer: int = customer.count()
n_customer_address: int = customer_address.count()
print(f"Filas en customer: {n_customer}")
print(f"Filas en customer_address: {n_customer_address}")
assert n_customer == 100_000
assert n_customer_address == 50_000

# Mismo filtro que usó S2 para reconocer las direcciones ya cruzadas con el
# INE: país fijado a España, municipio y provincia no nulos, y código de
# provincia de dos dígitos.
localized_addresses: DataFrame = customer_address.filter(
    (F.trim(F.col("ca_country")) == "España")
    & F.col("ca_city").isNotNull()
    & F.col("ca_county").isNotNull()
    & F.trim(F.col("ca_state")).rlike(r"^[0-9]{2}$")
)
assert localized_addresses.count() == n_customer_address

web_sales_slim: DataFrame = web_sales.select(
    "ws_order_number", "ws_bill_customer_sk", "ws_bill_addr_sk", "ws_net_profit"
)

Tipos de join¶

Probamos varios tipos de join entre web_sales_slim y customer_address por ws_bill_addr_sk = ca_address_sk. No todos los domicilios facturados en web_sales tienen por qué existir en customer_address con integridad referencial estricta —Parquet no la impone, y TPC-DS la deja opcional—, así que estos ejemplos también sirven para detectar huérfanos.

In [ ]:
from pyspark.sql.column import Column

join_expr: Column = web_sales_slim.ws_bill_addr_sk == customer_address.ca_address_sk

inner: DataFrame = web_sales_slim.join(customer_address, join_expr, "inner")
print(f"INNER JOIN: {inner.count()} filas.")
In [ ]:
left_outer: DataFrame = web_sales_slim.join(customer_address, join_expr, "left_outer")
print(f"LEFT OUTER JOIN: {left_outer.count()} filas (incluye líneas sin domicilio conocido).")
In [ ]:
right_outer: DataFrame = web_sales_slim.join(customer_address, join_expr, "right_outer")
print(f"RIGHT OUTER JOIN: {right_outer.count()} filas (incluye domicilios sin ventas).")
In [ ]:
left_semi: DataFrame = web_sales_slim.join(localized_addresses, join_expr, "left_semi")
print(
    f"LEFT SEMI JOIN: {left_semi.count()} filas de web_sales cuyo domicilio existe"
    " y está localizado (municipio del INE) en customer_address."
)
In [ ]:
left_anti: DataFrame = web_sales_slim.join(customer_address, join_expr, "left_anti")
print(
    f"LEFT ANTI JOIN: {left_anti.count()} filas de web_sales cuyo domicilio NO existe en customer_address."
)
# Producto cartesiano: cada fila de web_sales_slim con TODAS las filas de
# customer_address. NO debe ejecutarse: el resultado tendría
# aproximadamente 719.384 x 50.000 filas.
# cross = web_sales_slim.crossJoin(customer_address)

Broadcast join¶

customer_address tiene 50.000 filas: es pequeña comparada con las 719.384 de web_sales. Podemos pedirle a Spark que la distribuya entera a cada ejecutor con broadcast(), evitando el shuffle de la tabla grande.

In [ ]:
from pyspark.sql.functions import broadcast

sales_with_state: DataFrame = web_sales_slim.join(broadcast(customer_address), join_expr, "inner")

print("Plan físico del broadcast join:")
sales_with_state.explain("formatted")

En el plan busca BroadcastHashJoin y BroadcastExchange: la tabla grande (web_sales_slim) no aparece envuelta en un Exchange de shuffle, sólo la pequeña se distribuye entera. Guardamos sales_with_state porque la reutilizaremos más adelante, en la sección de persistencia.

Funciones escalares y agregadas¶

Calculamos estadísticos simples sobre ws_ext_sales_price, el importe extendido de venta antes de descuentos e impuestos.

In [ ]:
web_sales.select(
    F.avg("ws_ext_sales_price").alias("precio_medio"),
    F.max("ws_ext_sales_price").alias("precio_maximo"),
    F.min("ws_ext_sales_price").alias("precio_minimo"),
    F.count("ws_ext_sales_price").alias("lineas_con_precio"),
).show()
In [ ]:
web_sales.select("ws_ext_sales_price", "ws_net_profit").describe().show()

describe() añade además la desviación típica y los valores mínimo/máximo como cadenas de texto (por eso todas las columnas del resultado aparecen como string); es una forma rápida de comprobar rangos y detectar valores sospechosos antes de construir una consulta de negocio más elaborada.

Tipos complejos: arrays y structs¶

Spark puede representar columnas que contienen, a su vez, estructuras o listas. Vamos a construir, para cada cliente, su historial de pedidos web como un array de struct. Primero agregamos web_sales a nivel de pedido (no de línea): un mismo ws_order_number puede tener varias líneas, así que sumamos su importe neto pagado antes de construir el historial.

In [ ]:
orders_base: DataFrame = web_sales.groupBy("ws_bill_customer_sk", "ws_order_number").agg(
    F.sum("ws_net_paid").alias("importe_pedido"),
    F.first("ws_sold_date_sk").alias("ws_sold_date_sk"),
)
print(f"Pedidos distintos agregados: {orders_base.count()}")
In [ ]:
customer_order_history: DataFrame = orders_base.groupBy("ws_bill_customer_sk").agg(
    F.collect_list(
        F.struct(
            col("ws_order_number"),
            col("importe_pedido"),
            col("ws_sold_date_sk"),
        )
    ).alias("historial_pedidos")
)
customer_order_history.printSchema()
customer_order_history.show(5, truncate=False)

Cada elemento del array historial_pedidos es ya un pedido agregado (gracias al groupBy anterior por ws_order_number), no una línea suelta. Podemos acceder al primer pedido del array con getItem, contar cuántos pedidos tiene cada cliente con size, y filtrar clientes con muchos pedidos.

In [ ]:
customer_order_history.select(
    "ws_bill_customer_sk",
    F.col("historial_pedidos").getItem(0).alias("primer_pedido_del_array"),
    F.size("historial_pedidos").alias("numero_de_pedidos"),
).orderBy(F.col("numero_de_pedidos").desc()).show(10, truncate=False)
In [ ]:
clientes_con_muchos_pedidos: DataFrame = customer_order_history.filter(
    F.size("historial_pedidos") > 5
)
print(f"Clientes con más de 5 pedidos: {clientes_con_muchos_pedidos.count()}")

explode hace el camino inverso: convierte cada elemento del array en una fila independiente, recuperando una fila por pedido en vez de una fila por cliente.

In [ ]:
exploded_orders: DataFrame = customer_order_history.select(
    "ws_bill_customer_sk", F.explode("historial_pedidos").alias("pedido")
)
exploded_orders.select(
    "ws_bill_customer_sk",
    col("pedido.ws_order_number").alias("pedido_numero"),
    col("pedido.importe_pedido").alias("importe"),
).orderBy("ws_bill_customer_sk", "pedido_numero").show(10)

Funciones definidas por el usuario (UDFs)¶

Cuando una transformación no se puede expresar con las funciones integradas de Spark, se puede registrar una función Python como UDF. Vamos a clasificar cada línea de venta según su beneficio neto en tres etiquetas.

In [ ]:
from decimal import Decimal

from pyspark.sql.types import StringType


@F.udf(returnType=StringType())
def clasificar_beneficio(net_profit: Decimal | None) -> str:
    # Clasifica una línea de venta según su beneficio neto (ws_net_profit).
    if net_profit is None:
        return "desconocido"
    if net_profit < 0:
        return "perdida"
    if net_profit < 50:
        return "margen_bajo"
    return "margen_alto"


web_sales_classified: DataFrame = web_sales.withColumn(
    "clasificacion_beneficio", clasificar_beneficio(col("ws_net_profit"))
)
web_sales_classified.groupBy("clasificacion_beneficio").count().orderBy(
    F.col("count").desc()
).show()

Aviso de rendimiento. Esta UDF funciona, pero Spark no puede mirar dentro del código Python que hemos escrito: no sabe qué hace clasificar_beneficio, así que no puede aplicarle las optimizaciones de Catalyst (reordenar, podar columnas dentro de la función) ni ejecutarla en Tungsten como haría con una expresión nativa. Además, cada fila viaja de la JVM al proceso Python y vuelve, lo que añade serialización. La misma clasificación con funciones integradas —F.when(...).otherwise(...)— se ejecuta enteramente dentro del motor y suele ser bastante más rápida:

In [ ]:
web_sales_classified_builtin: DataFrame = web_sales.withColumn(
    "clasificacion_beneficio",
    F.when(col("ws_net_profit").isNull(), "desconocido")
    .when(col("ws_net_profit") < 0, "perdida")
    .when(col("ws_net_profit") < 50, "margen_bajo")
    .otherwise("margen_alto"),
)
web_sales_classified_builtin.groupBy("clasificacion_beneficio").count().orderBy(
    F.col("count").desc()
).show()

Regla de la sesión: usa una UDF sólo cuando no exista una combinación de funciones integradas que exprese la misma lógica.

Recapitulación

Joins, tipos complejos y UDFs

  • El tipo de join decide qué filas sin pareja sobreviven; un broadcast join evita el shuffle de la tabla grande
  • Un array de struct guarda el historial de pedidos de cada cliente; explode lo devuelve a filas
  • Una UDF de Python es opaca para el optimizador: sólo cuando no exista una función integrada equivalente

Funciones de ventana¶

Reutilizamos orders_base (un pedido por fila, con su importe neto pagado) y le añadimos la fecha de negocio uniendo con date_dim. A partir de ahí construimos cuatro ejemplos clásicos de funciones de ventana.

In [ ]:
from pyspark.sql.window import Window, WindowSpec

orders_dated: DataFrame = orders_base.join(
    date_dim, orders_base.ws_sold_date_sk == date_dim.d_date_sk
).select(
    col("ws_bill_customer_sk").alias("cliente"),
    col("ws_order_number").alias("pedido"),
    col("d_date").alias("fecha_pedido"),
    col("importe_pedido"),
)
orders_dated.orderBy("cliente", "fecha_pedido").show(10)

Máximo por partición: el pedido mayor de cada cliente¶

Window.partitionBy no reduce filas —a diferencia de groupBy—: añade una columna calculada a cada fila original, agregada dentro de su partición.

In [ ]:
from pyspark.sql.column import Column

ventana_por_cliente: WindowSpec = Window.partitionBy("cliente")
importe_maximo_cliente: Column = F.max("importe_pedido").over(ventana_por_cliente)

orders_dated.select(
    "cliente",
    "pedido",
    "importe_pedido",
    importe_maximo_cliente.alias("pedido_mayor_del_cliente"),
).withColumn(
    "es_su_pedido_mayor", col("importe_pedido") == col("pedido_mayor_del_cliente")
).orderBy(
    "cliente", "fecha_pedido"
).show(
    15
)

lag/lead: días entre pedidos consecutivos de un mismo cliente¶

Aquí sí importa el orden dentro de la ventana: la particionamos por cliente y la ordenamos por fecha de pedido.

In [ ]:
ventana_ordenada: WindowSpec = Window.partitionBy("cliente").orderBy("fecha_pedido")

fecha_pedido_anterior: Column = F.lag("fecha_pedido", 1).over(ventana_ordenada)
fecha_pedido_siguiente: Column = F.lead("fecha_pedido", 1).over(ventana_ordenada)

orders_dated.select(
    "cliente",
    "pedido",
    "fecha_pedido",
    F.datediff("fecha_pedido", fecha_pedido_anterior).alias("dias_desde_pedido_anterior"),
    F.datediff(fecha_pedido_siguiente, "fecha_pedido").alias("dias_hasta_pedido_siguiente"),
).orderBy("cliente", "fecha_pedido").show(15)

Media móvil de los últimos 3 pedidos por cliente¶

rowsBetween(-2, 0) define un marco que incluye la fila actual y las dos anteriores dentro de la misma partición ordenada.

In [ ]:
ventana_movil: WindowSpec = ventana_ordenada.rowsBetween(-2, 0)
media_movil_3_pedidos: Column = F.avg("importe_pedido").over(ventana_movil)

orders_dated.select(
    "cliente",
    "pedido",
    "fecha_pedido",
    "importe_pedido",
    media_movil_3_pedidos.alias("media_movil_3_pedidos"),
).orderBy("cliente", "fecha_pedido").show(15)

Suma acumulada del gasto por cliente¶

Window.unboundedPreceding hasta la fila actual acumula todas las filas anteriores de la partición: es el patrón habitual para un total corriente.

In [ ]:
ventana_acumulada: WindowSpec = ventana_ordenada.rowsBetween(
    Window.unboundedPreceding, Window.currentRow
)
gasto_acumulado: Column = F.sum("importe_pedido").over(ventana_acumulada)

orders_dated.select(
    "cliente",
    "pedido",
    "fecha_pedido",
    "importe_pedido",
    gasto_acumulado.alias("gasto_acumulado_cliente"),
).orderBy("cliente", "fecha_pedido").show(15)
A continuación

Las mismas preguntas, ahora en SQL

  • createOrReplaceTempView da un nombre SQL a un DataFrame sólo mientras dura la sesión
  • No crea un catálogo, ni un esquema, ni una tabla Hive
  • Repetimos las tres preguntas de negocio de S2: cliente con más pedidos, provincia con más compras y distribución horaria

Spark SQL con vistas temporales¶

Un DataFrame se puede registrar como vista temporal y consultarse con SQL:

web_sales.createOrReplaceTempView("web_sales")
spark.sql("SELECT ... FROM web_sales")

La vista desaparece al terminar la sesión: no crea ningún nombre persistente, ni un esquema de catálogo, ni una tabla Hive. Es sólo una forma de escribir la misma consulta con SQL en lugar de con el API de DataFrames. Reproducimos aquí, en Spark SQL, las tres preguntas de negocio que S2 ya respondió con PyArrow, Polars y DuckDB sobre los mismos ficheros, para comprobar que el motor cambia pero el significado de los datos no.

In [ ]:
web_sales.createOrReplaceTempView("web_sales")
customer.createOrReplaceTempView("customer")
customer_address.createOrReplaceTempView("customer_address")
date_dim.createOrReplaceTempView("date_dim")

time_dim: DataFrame = spark.read.parquet("/datalake/raw/tpcds/time_dim")
n_time_dim: int = time_dim.count()
print(f"Filas en time_dim: {n_time_dim}")
assert n_time_dim == 86_400
time_dim.createOrReplaceTempView("time_dim")

Cliente con más pedidos¶

Como en S2, contamos pedidos distintos con count(DISTINCT ws_order_number): contar * contaría líneas, no pedidos.

In [ ]:
spark.sql("""
    SELECT
        c.c_customer_id,
        count(DISTINCT s.ws_order_number) AS pedidos,
        count(*) AS lineas_de_venta
    FROM web_sales AS s
    JOIN customer AS c
        ON s.ws_bill_customer_sk = c.c_customer_sk
    GROUP BY c.c_customer_id
    ORDER BY pedidos DESC, c.c_customer_id
    LIMIT 10
""").show()

Provincia con más compras¶

S2 cruzó las direcciones sintéticas con la relación municipal del INE. ca_county contiene ahora la provincia y ca_state su código INE de dos caracteres; las claves que enlazan las ventas no cambiaron.

In [ ]:
spark.sql("""
    SELECT
        a.ca_state AS codigo_provincia,
        a.ca_county AS provincia,
        count(DISTINCT s.ws_order_number) AS pedidos,
        sum(s.ws_ext_sales_price) AS ventas
    FROM web_sales AS s
    JOIN customer_address AS a
        ON s.ws_bill_addr_sk = a.ca_address_sk
    GROUP BY a.ca_state, a.ca_county
    ORDER BY pedidos DESC, a.ca_state
    LIMIT 10
""").show()

Distribución horaria de las compras¶

Unimos por ws_sold_time_sk = t_time_sk, no por t_time: esa segunda columna es la representación numérica completa de la hora, no la clave que usa web_sales.

In [ ]:
spark.sql("""
    SELECT
        t.t_hour,
        count(DISTINCT s.ws_order_number) AS pedidos,
        sum(s.ws_ext_sales_price) AS ventas
    FROM web_sales AS s
    JOIN time_dim AS t
        ON s.ws_sold_time_sk = t.t_time_sk
    GROUP BY t.t_hour
    ORDER BY t.t_hour
""").show(24)

Las tres consultas anteriores son equivalentes, columna a columna y clave a clave, a las que S2 resolvió con PyArrow, Polars y DuckDB. Lo único que ha cambiado es el motor que las ejecuta y el hecho de que ahora puede distribuirse sobre YARN, como veremos más adelante.

Recapitulación

Las preguntas de S2, ahora con Spark SQL

  • Vistas temporales: un nombre para la sesión, sin catálogo ni tabla Hive
  • Mismo dataset, mismas claves, mismas métricas: cambia el motor, no el

significado

  • count(DISTINCT ws_order_number) sigue contando pedidos; count(*),

líneas

Beneficio neto por año¶

Esta consulta no viene de S2: es la misma agregación que ejecutará más adelante el trabajo que se envía a YARN (spark_sum(ws_net_profit), count(DISTINCT ws_order_number) y count(*), agrupando por d_year). Guardamos el resultado en local para compararlo con el que produzca YARN sobre los mismos datos.

In [ ]:
from pyspark.sql import Row

sales_by_year_local: DataFrame = spark.sql("""
    SELECT
        d.d_year,
        sum(s.ws_net_profit) AS beneficio_neto_total,
        count(DISTINCT s.ws_order_number) AS pedidos_distintos,
        count(*) AS lineas_de_venta
    FROM web_sales AS s
    JOIN date_dim AS d
        ON s.ws_sold_date_sk = d.d_date_sk
    GROUP BY d.d_year
    ORDER BY d.d_year
""")
sales_by_year_local.show()

# Se guarda como filas de Python: la SparkSession local se cierra antes de
# lanzar el trabajo sobre YARN y el DataFrame deja de poder usarse.
sales_by_year_local_rows: list[Row] = sales_by_year_local.collect()
A continuación

Escribir y paralelizar con criterio

  • Un fichero por partición; las vacías no llegan a escribir ninguno
  • repartition y coalesce: mismo verbo, costes muy distintos
  • spark.sql.shuffle.partitions y qué deja decidir a AQE

Guardar resultados¶

Escribimos la agregación de beneficio por año y categoría en HDFS, bajo una ruta de trabajo separada de /datalake. Es un groupBy normal sobre web_sales_dated_categorized, no el rollup de la sección anterior: aquí no hace falta ninguna fila de subtotal, sólo el detalle por año y categoría. Primero preparamos el espacio de trabajo de esta sesión.

In [ ]:
!hdfs dfs -mkdir -p /user/luser/s4
In [ ]:
sales_by_year_category: DataFrame = web_sales_dated_categorized.groupBy("d_year", "i_category").agg(
    spark_sum("ws_net_profit").alias("beneficio_neto")
)
print(f"Filas a escribir: {sales_by_year_category.count()}")

El problema de los ficheros pequeños¶

Si repartimos el resultado en muchas particiones antes de escribir, Spark crea hasta un fichero Parquet por partición (las particiones vacías no llegan a escribir fichero), y la mayoría de las que sí escriben quedan casi vacías: es el problema de los "ficheros pequeños" que tanto penaliza a HDFS y a los motores que después lean esa tabla.

Para listar el resultado usamos una orden de shell con una variable de Python entre llaves: en las celdas !orden, IPython sustituye {variable} por su valor antes de ejecutar la orden.

In [ ]:
many_files_path: str = "/user/luser/s4/sales-by-year-category-muchos-ficheros"
sales_by_year_category.repartition(50).write.mode("overwrite").parquet(many_files_path)
In [ ]:
!hdfs dfs -ls {many_files_path}

Consolidar antes de escribir con coalesce¶

coalesce(n) reduce el número de particiones sin forzar un shuffle completo (a diferencia de repartition(n)), agrupando particiones existentes. Es la forma habitual de evitar el problema anterior sin pagar el coste de una redistribución completa.

In [ ]:
few_files_path: str = "/user/luser/s4/sales-by-year-category-pocos-ficheros"
sales_by_year_category.coalesce(4).write.mode("overwrite").parquet(few_files_path)
In [ ]:
!hdfs dfs -ls {few_files_path}

En esta sección coalesce controla cuántos ficheros produce la escritura. S5 retomará esta idea al organizar físicamente los datos por valores de columna.

Particiones de ejecución¶

repartition(n) y coalesce(n) cambian cuántas particiones tiene un DataFrame, pero tienen costes distintos: repartition siempre hace un shuffle completo (puede aumentar o reducir el número de particiones, redistribuyendo todas las filas); coalesce sólo puede reducir particiones y evita el shuffle completo agrupando particiones vecinas, aunque a costa de un reparto potencialmente menos equilibrado.

Reparticionar antes de un join¶

Si dos DataFrames que se van a unir se reparticionan por la misma clave de unión y con el mismo número de particiones, Spark puede aprovechar esa coincidencia para reducir el shuffle del join. No es automático ni garantizado: hay que comprobarlo en el plan, no darlo por supuesto por haber llamado a repartition.

In [ ]:
print(f"web_sales_slim particiones (antes): {web_sales_slim.rdd.getNumPartitions()}")
print(f"customer_address particiones (antes): {customer_address.rdd.getNumPartitions()}")

web_sales_by_addr: DataFrame = web_sales_slim.repartition(8, "ws_bill_addr_sk")
customer_address_by_key: DataFrame = customer_address.repartition(8, "ca_address_sk")

print(f"web_sales_by_addr particiones (después): {web_sales_by_addr.rdd.getNumPartitions()}")
print(
    f"customer_address_by_key particiones (después): {customer_address_by_key.rdd.getNumPartitions()}"
)
In [ ]:
join_expr_repartitioned: Column = (
    web_sales_by_addr.ws_bill_addr_sk == customer_address_by_key.ca_address_sk
)
sales_with_state_repartitioned: DataFrame = web_sales_by_addr.join(
    customer_address_by_key, join_expr_repartitioned, "inner"
)

print("Plan físico tras reparticionar por la clave de unión:")
sales_with_state_repartitioned.explain("formatted")

Busca en el plan si sigue apareciendo un Exchange justo antes del Join para cada lado. Reparticionar por la misma clave y el mismo número de particiones facilita que Spark reconozca la coincidencia, pero el optimizador decide caso a caso; por eso hay que medir el plan en lugar de suponer el efecto sólo por el nombre de la función.

Configuración del shuffle¶

spark.sql.shuffle.partitions fija cuántas particiones produce cualquier operación wide. Lo comprobamos con una agregación de store_sales —2.880.404 filas, la mayor tabla de hechos que usaremos en esta sesión— por año de venta. Desactivamos primero Adaptive Query Execution (activo por defecto en Spark 4): si no, AQE podría fusionar particiones de shuffle pequeñas él mismo y el número de particiones que se imprime a continuación dejaría de depender solo de shuffle.partitions. La reactivamos al terminar esta demostración, antes de la sección dedicada a AQE.

In [ ]:
spark.conf.set("spark.sql.adaptive.enabled", "false")

store_sales: DataFrame = spark.read.parquet("/datalake/raw/tpcds/store_sales")
n_store_sales: int = store_sales.count()
print(f"Filas en store_sales: {n_store_sales}")
assert n_store_sales == 2_880_404

print(f"spark.sql.shuffle.partitions actual: {spark.conf.get('spark.sql.shuffle.partitions')}")
In [ ]:
store_sales_by_year: DataFrame = (
    store_sales.join(date_dim, store_sales.ss_sold_date_sk == date_dim.d_date_sk)
    .groupBy("d_year")
    .agg(spark_sum("ss_ext_sales_price").alias("ventas_tienda"))
)
print(
    f"Particiones tras el groupBy (shuffle.partitions=8): {store_sales_by_year.rdd.getNumPartitions()}"
)
In [ ]:
spark.conf.set("spark.sql.shuffle.partitions", 3)

store_sales_by_year_3: DataFrame = (
    store_sales.join(date_dim, store_sales.ss_sold_date_sk == date_dim.d_date_sk)
    .groupBy("d_year")
    .agg(spark_sum("ss_ext_sales_price").alias("ventas_tienda"))
)
print(
    f"Particiones tras el groupBy (shuffle.partitions=3): {store_sales_by_year_3.rdd.getNumPartitions()}"
)

# Volvemos a los valores usados en el resto del notebook
spark.conf.set("spark.sql.shuffle.partitions", 8)
spark.conf.set("spark.sql.adaptive.enabled", "true")

El número de filas del resultado no cambia (sigue habiendo un puñado de años en SF1); lo que cambia es en cuántas particiones queda repartido ese resultado, y por tanto cuántas tareas paralelas puede lanzar Spark sobre él en el siguiente paso del pipeline.

Adaptive Query Execution (AQE)¶

AQE deja que Spark reajuste el plan durante la ejecución, con estadísticas reales en lugar de sólo estimaciones previas. Una de sus funciones más visibles es fusionar particiones de shuffle que han quedado demasiado pequeñas tras un filtro selectivo.

Para verlo con datos reales, primero calculamos el rango de años presente en date_dim, en lugar de suponer un año concreto.

In [ ]:
from pyspark.sql.types import Row

year_bounds: Row | None = date_dim.agg(
    F.min("d_year").alias("min_year"), F.max("d_year").alias("max_year")
).first()
sample_year: int = (year_bounds["min_year"] + year_bounds["max_year"]) // 2
print(
    f"Años en date_dim: {year_bounds['min_year']}-{year_bounds['max_year']}. Usamos {sample_year}."
)
In [ ]:
spark.conf.set("spark.sql.shuffle.partitions", 200)
spark.conf.set("spark.sql.adaptive.enabled", "false")

sin_aqe: DataFrame = (
    store_sales.join(date_dim, store_sales.ss_sold_date_sk == date_dim.d_date_sk)
    .filter(col("d_year") == sample_year)
    .groupBy("d_year")
    .agg(spark_sum("ss_ext_sales_price").alias("ventas_tienda"))
)
sin_aqe.collect()
print(f"Particiones sin AQE (shuffle.partitions=200): {sin_aqe.rdd.getNumPartitions()}")
In [ ]:
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")

con_aqe: DataFrame = (
    store_sales.join(date_dim, store_sales.ss_sold_date_sk == date_dim.d_date_sk)
    .filter(col("d_year") == sample_year)
    .groupBy("d_year")
    .agg(spark_sum("ss_ext_sales_price").alias("ventas_tienda"))
)
con_aqe.collect()
print(f"Particiones con AQE (coalescePartitions habilitado): {con_aqe.rdd.getNumPartitions()}")

# Volvemos al valor usado en el resto del notebook
spark.conf.set("spark.sql.shuffle.partitions", 8)

El filtro d_year == sample_year selecciona un único año de los varios que hay en date_dim; el groupBy("d_year") posterior agrupa todas esas filas bajo una sola clave, así que el shuffle concentra casi todo en una partición y deja las otras 199 vacías. Sin AQE, esas 200 particiones de shuffle casi vacías se mantienen tal cual. Con coalescePartitions habilitado, Spark las fusiona en menos particiones después de ver cuántas filas sobrevivieron al filtro y a cuántas claves de agrupación se redujeron, no antes.

Recapitulación

Particiones y shuffle

  • Partición de ejecución ≠ partición Hive ≠ bloque HDFS ≠ fichero Parquet
  • shuffle.partitions fija el paralelismo tras una operación ancha
  • AQE corrige el plan con estadísticas reales: fusiona particiones después

del filtro, no antes

A continuación

Persistencia

  • Spark recalcula un DataFrame desde sus fuentes cada vez que una acción lo necesita
  • cache() y persist() guardan el resultado la primera vez que se materializa
  • MEMORY_AND_DISK lleva a disco las particiones que no caben en memoria
  • Al terminar cerramos la sesión local: la parte de YARN necesita un proceso nuevo

Persistencia¶

sales_with_state, el resultado del broadcast join de la sección de joins, es un buen candidato para cachear: lo vamos a consultar varias veces seguidas. Por defecto, Spark recalcula un DataFrame desde sus fuentes cada vez que una acción lo necesita; cache()/persist() evitan ese recálculo guardando el resultado la primera vez que se materializa.

In [ ]:
print(f"¿Está cacheado antes de cache()?: {sales_with_state.is_cached}")

sales_with_state.cache()
sales_with_state.count()  # acción que materializa la cache

print(f"¿Está cacheado después de cache()?: {sales_with_state.is_cached}")
print(f"Nivel de almacenamiento: {sales_with_state.storageLevel}")
In [ ]:
sales_with_state.groupBy("ca_state", "ca_county").count().orderBy(F.col("count").desc()).show(5)
sales_with_state.agg(F.avg("ws_net_profit")).show()

Ambas acciones anteriores han reutilizado la versión cacheada, sin volver a ejecutar el broadcast join. Pero la persistencia no se hereda: un DataFrame construido a partir de uno cacheado, mediante una nueva transformación, no queda cacheado automáticamente. Mientras sales_with_state está cacheado, la pestaña Storage de la interfaz de Spark lo lista con el espacio que ocupa.

In [ ]:
sales_with_state_profitable: DataFrame = sales_with_state.filter(col("ws_net_profit") > 0)
print(f"¿Está cacheado el DataFrame derivado?: {sales_with_state_profitable.is_cached}")
In [ ]:
from pyspark import StorageLevel

sales_with_state.unpersist()
print(f"¿Está cacheado tras unpersist()?: {sales_with_state.is_cached}")

sales_with_state.persist(StorageLevel.MEMORY_AND_DISK)
sales_with_state.count()
print(f"Nuevo nivel de almacenamiento: {sales_with_state.storageLevel}")
sales_with_state.unpersist()

MEMORY_AND_DISK guarda en disco las particiones que no quepan en memoria, en lugar de recalcularlas cada vez; es un compromiso razonable cuando el DataFrame es más grande que la memoria disponible del driver en modo local. Dejamos sales_with_state sin persistir al terminar esta sección, para no competir por memoria con las siguientes.

In [ ]:
spark.stop()

Cerramos aquí la sesión local: la siguiente sección no puede reutilizarla, porque Spark no permite cambiar el master de una SparkSession ya construida. El envío a YARN se hace como un proceso completamente aparte, mediante spark-submit.

A continuación

De local[2] a los DataNodes

  • Driver, ApplicationMaster y ejecutores como procesos reales
  • Un script aparte lanzado con spark-submit --master yarn: el master

no se puede cambiar en una sesión ya creada

  • El driver usa el PySpark de /opt/tcdm/venv; los ejecutores, sólo un

python3 en cada DataNode

Driver, ejecutores y YARN¶

Hasta ahora, "driver" y "ejecutores" eran hilos dentro del mismo proceso local[2]. Sobre YARN son procesos reales, potencialmente en máquinas distintas:

  • el driver sigue siendo el proceso que construye el plan y coordina la aplicación; en modo client (el que usaremos) vive en namenode, junto al proceso spark-submit que lo lanza;
  • el ApplicationMaster es el interlocutor de la aplicación ante el ResourceManager: pide y libera contenedores;
  • los ejecutores son contenedores YARN que se ejecutan en los NodeManagers de los DataNodes —los mismos datanode1, datanode2 y datanode3 de S1— y sobre los que Spark reparte las tareas.

Animación: spark-submit --master yarn se lanza desde namenode y los ejecutores calculan en datanode1, datanode2 y datanode3

Como ya se explicó en "Instalar PySpark en este notebook", una SparkSession no puede cambiar de master una vez construida: por eso esta sección no reutiliza la sesión local de las secciones anteriores, sino que escribe un script independiente y lo lanza con spark-submit --master yarn como un proceso nuevo.

Qué Python usan el driver y los ejecutores¶

El envío a YARN usa el mismo entorno Python que el kernel de este notebook, /opt/tcdm/venv, instalado en la imagen de namenode y comprobado con %pip install -r requirements.txt al principio de la sesión. De él salen tanto spark-submit como el intérprete del driver, el único proceso que necesita el paquete pyspark instalado para construir la SparkSession y lanzar la aplicación.

Los ejecutores son otra historia: se ejecutan en los DataNodes, cuyas imágenes sólo traen el python3 de base (sin PySpark ni ningún paquete instalado por pip). spark-submit --master yarn distribuye a través de YARN los ficheros de soporte de PySpark (pyspark.zip y py4j.zip) y los .jar de Spark, así que a cada DataNode le basta con ese intérprete base para ejecutar el código ya distribuido. Por eso el envío pide spark.pyspark.python=python3 a secas: el python3 del sistema en cada DataNode, que es la misma versión de Python que la de /opt/tcdm/venv porque ambas imágenes parten de la misma base.

El script que se enviará a YARN¶

Escribimos un fichero .py independiente. Hace dos cosas: calcula el beneficio neto total de web_sales por año (uniendo con date_dim, la misma clave de negocio de siempre) y, además, genera un pequeño informe de qué host ha procesado cada partición. mapPartitionsWithIndex ejecuta una función en cada partición y permite registrar allí socket.gethostname(). Así se comprueba que YARN reparte el trabajo entre varios DataNodes.

In [ ]:
from pathlib import Path

JOB_SCRIPT: str = r'''
# Trabajo independiente para YARN, lanzado con spark-submit --master yarn.
#
# No es la SparkSession del notebook: Spark no permite cambiar el master de
# una sesión ya construida, así que esta aplicación crea la suya propia.
import socket

from pyspark.sql import SparkSession
from pyspark.sql.functions import count, countDistinct
from pyspark.sql.functions import sum as spark_sum
from pyspark.sql.types import IntegerType, LongType, StringType, StructField, StructType

BUSINESS_OUTPUT = "/user/luser/s4/yarn-job/business-result"
EXECUTOR_REPORT_OUTPUT = "/user/luser/s4/yarn-job/executor-report"

spark = SparkSession.builder.appName("TCDM_S4_YARN_business").master("yarn").getOrCreate()

date_dim = spark.read.parquet("/datalake/raw/tpcds/date_dim")
web_sales = spark.read.parquet("/datalake/raw/tpcds/web_sales")

sales_by_year = (
    web_sales.join(date_dim, web_sales.ws_sold_date_sk == date_dim.d_date_sk)
    .groupBy("d_year")
    .agg(
        spark_sum("ws_net_profit").alias("beneficio_neto_total"),
        countDistinct("ws_order_number").alias("pedidos_distintos"),
        count("*").alias("lineas_de_venta"),
    )
    .orderBy("d_year")
)
sales_by_year.coalesce(1).write.mode("overwrite").csv(BUSINESS_OUTPUT, header=True)

# Registra el host que ha ejecutado cada partición.
def describe_partition(partition_id, rows):
    rows = list(rows)
    yield (partition_id, socket.gethostname(), len(rows))


reports = (
    web_sales.select("ws_order_number").rdd.mapPartitionsWithIndex(describe_partition).collect()
)

report_schema = StructType(
    [
        StructField("partition", IntegerType(), False),
        StructField("executor_host", StringType(), False),
        StructField("records", LongType(), False),
    ]
)
(
    spark.createDataFrame(reports, report_schema)
    .orderBy("partition")
    .coalesce(1)
    .write.mode("overwrite")
    .csv(EXECUTOR_REPORT_OUTPUT, header=True)
)

distinct_hosts = {report[1] for report in reports}
print(f"Hosts remotos que procesaron alguna partición: {sorted(distinct_hosts)}")
spark.stop()
'''

Path("/tmp/s4_yarn_job.py").write_text(JOB_SCRIPT)
print("Escrito /tmp/s4_yarn_job.py")

Cuántos ejecutores pedir¶

El número de ejecutores que se solicita debe derivarse de los NodeManagers que YARN anuncia como RUNNING, no de un número fijo elegido a mano. La variable de entorno TCDM_SPARK_EXECUTORS, si está definida, permite forzar otro valor, por ejemplo para probar con menos ejecutores.

In [ ]:
import subprocess

result: subprocess.CompletedProcess[str] = subprocess.run(
    [
        "bash",
        "-lc",
        "yarn node -list -all 2>/dev/null | awk '$2 == \"RUNNING\" {count++} END {print count + 0}'",
    ],
    capture_output=True,
    text=True,
    check=True,
)
running_nodemanagers: int = int(result.stdout.strip() or 0)
print(f"NodeManagers en estado RUNNING: {running_nodemanagers}")

assert running_nodemanagers >= 2, "Se esperan al menos dos NodeManagers RUNNING"
In [ ]:
import os

requested_executors: int = int(os.environ.get("TCDM_SPARK_EXECUTORS", running_nodemanagers))
print(f"Solicitando {requested_executors} ejecutores Spark a YARN.")

El envío con spark-submit¶

Los valores de memoria y núcleos están ajustados a los recursos de este clúster:

  • --executor-cores 1 --executor-memory 512m junto con spark.executor.memoryOverhead=256m suman 768 MiB por ejecutor, dentro de los 2560 MiB que cada NodeManager anuncia a YARN (recuerda: de los 3072 MiB del contenedor DataNode, 512 MiB quedan reservados para el propio DataNode, el NodeManager y otros procesos auxiliares).
  • spark.yarn.am.cores y spark.yarn.am.memory/memoryOverhead son la petición de recursos del ApplicationMaster, un contenedor YARN aparte de los ejecutores.
  • spark.scheduler.minRegisteredResourcesRatio=1.0 junto con maxRegisteredResourcesWaitingTime=60s hacen que Spark espere a tener todos los ejecutores solicitados (o hasta 60 segundos) antes de empezar a repartir tareas, evitando medir la distribución de trabajo con sólo parte de los ejecutores ya registrados.
  • spark.pyspark.python=python3 es el intérprete de los ejecutores (el python3 de base de cada DataNode); spark.pyspark.driver.python apunta al Python de /opt/tcdm/venv, el único sitio donde hace falta el paquete pyspark instalado.

Usamos el spark-submit de /opt/tcdm/venv/bin, el mismo entorno que el kernel, para que coincida exactamente con la versión de PySpark que usa el resto del notebook.

Mientras el trabajo está en marcha, su driver —que también vive en namenode— vuelve a servir la interfaz de Spark, y la pestaña Executors muestra ahora un ejecutor por contenedor YARN, cada uno en su DataNode. Sobre YARN esa interfaz se consulta a través del ResourceManager: en http://localhost:8088, el enlace ApplicationMaster de la aplicación; si el navegador acaba en una dirección con el nombre resourcemanager, sustitúyelo por localhost. El trabajo dura poco y la interfaz desaparece cuando termina: la evidencia permanente de qué host procesó cada partición es el informe que escribe el propio script.

La orden imprime muchas líneas de registro de Spark y puede tardar un minuto o más. La que importa está casi al final: Hosts remotos que procesaron alguna partición: [...].

In [ ]:
spark_submit_cmd: str = f"""/opt/tcdm/venv/bin/spark-submit \
    --master yarn \
    --deploy-mode client \
    --num-executors {requested_executors} \
    --executor-cores 1 \
    --executor-memory 512m \
    --conf spark.executor.memoryOverhead=256m \
    --conf spark.yarn.am.cores=1 \
    --conf spark.yarn.am.memory=256m \
    --conf spark.yarn.am.memoryOverhead=256m \
    --conf spark.scheduler.minRegisteredResourcesRatio=1.0 \
    --conf spark.scheduler.maxRegisteredResourcesWaitingTime=60s \
    --conf spark.pyspark.python=python3 \
    --conf spark.pyspark.driver.python=/opt/tcdm/venv/bin/python \
    /tmp/s4_yarn_job.py"""

print(spark_submit_cmd)
In [ ]:
!export HADOOP_CONF_DIR=/opt/bd/hadoop/etc/hadoop && {spark_submit_cmd}

Comprobar las evidencias¶

Leemos de vuelta los dos resultados que ha escrito el trabajo YARN y comprobamos la aplicación en la CLI de YARN, la misma que el ResourceManager (http://localhost:8088) muestra de forma gráfica.

In [ ]:
!hdfs dfs -cat /user/luser/s4/yarn-job/business-result/part-*.csv

Ese resultado debe ser el mismo que calculó la sesión local en la sección de Spark SQL. Lo comprobamos fila a fila, con las filas que guardamos entonces en sales_by_year_local_rows:

In [ ]:
import csv
import io
from decimal import Decimal

yarn_csv: str = subprocess.run(
    ["hdfs", "dfs", "-cat", "/user/luser/s4/yarn-job/business-result/part-*.csv"],
    capture_output=True,
    text=True,
    check=True,
).stdout

yarn_by_year: dict[int, tuple[int, int, Decimal]] = {
    int(row["d_year"]): (
        int(row["pedidos_distintos"]),
        int(row["lineas_de_venta"]),
        Decimal(row["beneficio_neto_total"]),
    )
    for row in csv.DictReader(io.StringIO(yarn_csv))
}
local_by_year: dict[int, tuple[int, int, Decimal]] = {
    row["d_year"]: (row["pedidos_distintos"], row["lineas_de_venta"], row["beneficio_neto_total"])
    for row in sales_by_year_local_rows
}

assert yarn_by_year == local_by_year, "El resultado sobre YARN no coincide con el local"
print(f"Local y YARN coinciden en los {len(local_by_year)} años.")
In [ ]:
!hdfs dfs -cat /user/luser/s4/yarn-job/executor-report/part-*.csv
In [ ]:
!yarn application -list -appStates ALL

Por último, extraemos del informe de particiones la lista de hosts distintos que han hecho trabajo real y comprobamos que al menos dos DataNodes remotos participan cuando los recursos lo permiten. Si sólo aparece uno, repite el envío: con tan pocos datos el planificador puede concentrar el trabajo en un único nodo.

In [ ]:
from subprocess import CompletedProcess

report_result: CompletedProcess[str] = subprocess.run(
    ["bash", "-lc", "hdfs dfs -cat /user/luser/s4/yarn-job/executor-report/part-*.csv"],
    capture_output=True,
    text=True,
    check=True,
)
report_lines: list[str] = [line for line in report_result.stdout.splitlines()[1:] if line.strip()]
executor_hosts: set[str] = {line.split(",")[1] for line in report_lines}

print(f"Hosts que procesaron al menos una partición: {sorted(executor_hosts)}")
print(f"Ejecutores remotos que trabajaron: {len(executor_hosts)}")
assert len(executor_hosts) >= 2, "Se esperan al menos dos hosts remotos distintos"

Si todo ha ido bien, el resultado de negocio es el mismo que ya vimos en modo local en la sección de Spark SQL —beneficio neto, pedidos distintos y líneas de venta por año, comparados más arriba con un assert—, sólo que esta vez calculado por ejecutores que corren en, al menos, dos DataNodes distintos, coordinados por un driver que vive en namenode.

Recapitulación

El trabajo ocurrió en los DataNodes

  • El informe registra qué host procesó cada partición: al menos dos

DataNodes remotos hicieron trabajo real

  • El número de ejecutores se derivó de los NodeManagers RUNNING, no de

un número elegido a mano

  • El resultado de negocio es el mismo que en local: el motor cambia, la

métrica no

Hilo de negocio con TPC-DS¶

A lo largo de la sesión se ha mantenido la misma disciplina que en S2, sólo que ahora ejecutada por Spark, en local y sobre YARN:

  • contar líneas de web_sales con count(*) es distinto de contar pedidos con count(DISTINCT ws_order_number), y cada consulta de esta sesión ha dicho explícitamente cuál de las dos calculaba;
  • las sumas de importes o unidades siempre han indicado qué columna representa la medida (ws_net_profit, ws_ext_sales_price, ss_ext_sales_price, ws_net_paid...), en lugar de sumar "lo que hubiera";
  • web_sales se ha unido con date_dim mediante ws_sold_date_sk = d_date_sk en todas las agregaciones por año, con time_dim mediante ws_sold_time_sk = t_time_sk en la distribución horaria, y con customer/customer_address sólo cuando la pregunta lo necesitaba;
  • ningún identificador surrogate (ws_bill_customer_sk, c_customer_sk, ca_address_sk...) se ha presentado como si fuera el nombre de una persona o una región administrativa real: son claves sintéticas del generador TPC-DS.

La sección de YARN calculó el beneficio neto por año exactamente con la misma métrica y las mismas claves que la sección de Spark SQL calculó en local: el motor de ejecución cambia, pero la respuesta de negocio debería ser la misma consultando el mismo dataset inmutable.

Preguntas para interpretar la experiencia¶

  • ¿Qué diferencia concreta hay entre una transformación y una acción? Da un ejemplo de esta sesión de cada una.
  • ¿Por qué select("*") y una selección de dos columnas producen el mismo resultado filtrado pero planes físicos distintos? ¿Qué cambia realmente entre ambos?
  • El número de particiones de ejecución de un DataFrame, ¿depende del número de ficheros Parquet, del número de bloques HDFS o del número de ejecutores disponibles? Justifica la respuesta con algo que hayas medido en el notebook.
  • ¿Por qué un join normal necesita un shuffle y un broadcast join puede evitarlo? ¿Qué condición debe cumplir la tabla que se distribuye entera?
  • ¿Qué controla spark.sql.shuffle.partitions y en qué se diferencia de las particiones de lectura que vimos con date_dim y web_sales?
  • ¿Qué decide Adaptive Query Execution por sí solo, y qué sigue estando bajo el control explícito de quien escribe la consulta (por ejemplo, cuántos ejecutores pedir a YARN)?
  • ¿Por qué la sección de YARN no pudo reutilizar la SparkSession local del resto del notebook? ¿Qué tuvo que crearse aparte y por qué?
  • En la sección de UDFs, ¿qué información concreta pierde Spark al no poder "mirar dentro" de una función Python, y qué consecuencia tiene eso sobre el rendimiento?

Evidencias para la siguiente sesión¶

Antes de la siguiente sesión tendrás una reunión individual breve con el profesor para revisar el trabajo de esta sesión. Esa reunión combina una demostración en vivo sobre tu propio ordenador y una memoria escrita breve.

Qué mostrar en el ordenador durante la reunión¶

Ten preparadas estas evidencias:

  1. El esquema (printSchema) y el número de filas de date_dim o web_sales, junto con su número de particiones de ejecución.
  2. El plan de .explain("formatted") de al menos una consulta con join (por ejemplo, el broadcast join con customer_address).
  3. El número de particiones de un DataFrame antes y después de un repartition o un coalesce.
  4. La aplicación YARN completada (yarn application -list -appStates ALL) y la lista de al menos dos hosts distintos que procesaron alguna partición (de ahí se deduce cuántos ejecutores remotos trabajaron).
  5. El resultado de negocio (beneficio neto por año de web_sales) obtenido en local con Spark SQL y el mismo resultado obtenido sobre YARN, comprobando que coinciden.
  6. Que sabes explicar, con tus propias palabras, la diferencia entre una partición de ejecución de Spark, una partición Hive y un bloque HDFS.

Memoria escrita (una o dos páginas)¶

Trae también un documento breve —una o dos páginas, no hace falta más— que no se limite a pegar capturas de las salidas anteriores: debe explicar con tus propias palabras qué papel tienen el driver, los ejecutores y el gestor de recursos, y qué cambia entre local[2] y YARN; la diferencia entre una transformación y una acción, y qué se lee en un plan de explain; qué es una partición de ejecución y en qué se distingue de un fichero Parquet, de un bloque HDFS y de una partición Hive; y por qué el resultado de negocio es el mismo en local y sobre YARN.

Qué es importante de cara al examen final¶

El examen no pide recordar la sintaxis exacta de una función. Debes poder explicar:

  • qué hacen el driver, los ejecutores y el gestor de recursos, y qué son un job, una stage y una task;
  • la evaluación perezosa: qué es una transformación, qué es una acción y por qué Spark construye un plan antes de ejecutar;
  • la diferencia entre dependencias estrechas y anchas, qué es un shuffle y qué operaciones lo provocan;
  • qué se comprueba en un plan: selección de columnas, filtros empujados a la lectura y tipo de join (incluido el broadcast join);
  • qué controlan repartition, coalesce y spark.sql.shuffle.partitions, qué decide por sí solo AQE y por qué importan los ficheros pequeños;
  • cuándo compensa persistir un DataFrame y por qué una UDF de Python es más lenta que una función integrada;
  • la diferencia entre contar líneas de venta y contar pedidos distintos.
Recapitulación

Sesión 4

  • Spark lee y escribe los mismos Parquet; el driver nunca materializa la

tabla entera

  • Transformaciones perezosas, acciones que ejecutan, planes que se leen

con explain

  • Local y YARN son el mismo programa: cambian el gestor de recursos y

dónde viven los ejecutores

Siguiente bloque: buenas prácticas de data lake y particionado físico

con el mismo motor (S5)

Siguiente paso¶

S4 enseña Spark como motor. S5 utilizará ese motor para crear una tabla refinada como colección de Parquet particionados. Todavía se conservará una zona fuente inmutable y no se presentará la salida como una tabla gestionada.