Saltar al contenido
Lakehouse LabLakehouse LabPreparación Databricks Data Engineer
Módulo 05 · Associate + Professional

Spark

Contenido abierto

Puedes leer todo sin registrarte. Solo crearemos un perfil anónimo cuando decidas guardar tu progreso.

Associate + Professional

Catalyst, particiones, joins y shuffles

Razona sobre el plan lógico y físico antes de ajustar configuraciones o añadir cómputo.

Lectura pública
Al terminar podrás
  • Leer explain formatted
  • Detectar Exchange, skew y spill
  • Elegir broadcast, repartition o coalesce con evidencia
Ver fuentes y revisión

Metadatos editoriales

Última revisión
21 jul 2026
Nivel
Associate + Professional
Ruta relacionada
core
Dominios blueprint
Troubleshooting and Optimization · Spark execution
Estado
Revisión editorial interna
Fuentes principales
Optimización en Databricks · Debugging con Spark UI
Reportar un error
01
Modelo mental

Lazy evaluation y DAG de Spark

El plan físico y sus métricas indican dónde se mueve y procesa el dato.

Objetivo
El plan físico y sus métricas indican dónde se mueve y procesa el dato.
Duración estimada
17 min aprox.
Dificultad
Associate + Professional
Prerrequisitos
m04
Reportar un error en esta lección

explain formatted separa plan lógico y físico; Exchange suele indicar shuffle, y el tipo de join muestra la estrategia elegida.

Spark UI confirma con bytes, duración y distribución por tarea. Optimiza una hipótesis cada vez y conserva una línea base.

Modelo mental

Un plan de Spark es una explicación ejecutable de cómo una intención relacional se convierte en trabajo distribuido. El plan lógico conserva operaciones como filtro, proyección y join; Catalyst lo analiza y optimiza; el plan físico elige scans, algoritmos de join, exchanges y agregaciones concretas. explain formatted muestra estructura y estimaciones, pero no sustituye métricas reales. Spark UI revela stages, tareas, bytes, shuffle, spill y distribución temporal después de ejecutar. La lectura competente conecta ambos: un Exchange anticipa redistribución, mientras las métricas confirman su volumen y equilibrio. Optimizar significa formular una hipótesis causal, cambiar una sola variable y comparar contra una línea base semánticamente equivalente.

Catalyst

Optimizador de consultas de Spark que analiza expresiones y transforma planes lógicos en alternativas físicas.

Permite razonar sobre por qué código distinto puede producir ejecución equivalente.
Exchange

Operador físico que redistribuye datos entre executors y suele crear una frontera de stage.

Señala un shuffle potencialmente costoso que debe confirmarse con métricas.
Plan adaptativo

Plan que AQE puede revisar durante ejecución usando estadísticas observadas en runtime.

Explica cambios de estrategia y particiones que no aparecen en el plan inicial.
PySparkPlan físico
result.explain("formatted")

Busca Exchange, Scan y el operador de join; después contrasta en Spark UI.

Puntos clave

  • Exchange señala redistribución
  • Plan y métricas se complementan
  • Mide antes y después

Evita

  • Ajustar configuraciones sin línea base
  • Interpretar solo el plan lógico

Recuerdo activo

¿Qué operador suele delimitar un shuffle?

Borrador privado · solo en este navegador
02
Implementación

Catalyst y optimizaciones de consulta

Particiones determinan paralelismo, tamaño de tareas y número de archivos.

Objetivo
Particiones determinan paralelismo, tamaño de tareas y número de archivos.
Duración estimada
17 min aprox.
Dificultad
Associate + Professional
Prerrequisitos
m04
Reportar un error en esta lección

spark.sql.shuffle.partitions controla particiones tras shuffles SQL; spark.default.parallelism influye en operaciones RDD y ciertos orígenes.

repartition provoca shuffle y puede aumentar o redistribuir; coalesce suele reducir sin shuffle completo. El objetivo es tareas suficientemente numerosas y archivos de tamaño razonable.

Modelo mental

Una partición de Spark es la unidad de datos que una tarea procesa secuencialmente. El número y distribución de particiones delimitan paralelismo, overhead, memoria por tarea y archivos de salida. Muy pocas dejan cores ociosos y crean tareas grandes; demasiadas producen planificación, conexiones y archivos pequeños. spark.sql.shuffle.partitions fija un punto de partida para shuffles SQL, aunque AQE puede fusionar particiones posteriores. repartition introduce una redistribución para aumentar o equilibrar; coalesce suele reducir aprovechando la distribución existente. No existe un número universal. Se dimensiona a partir de volumen comprimido, recursos, operadores y distribución, y se valida con métricas por tarea y tamaño de archivos.

Partición

Segmento lógico de un dataset que una tarea Spark procesa en un executor.

Conecta distribución de datos con paralelismo, memoria y duración.
Dependencia ancha

Relación en la que una partición de salida necesita datos de múltiples particiones de entrada.

Suele requerir shuffle y crear una nueva etapa de ejecución.
repartition

Operación que redistribuye datos mediante shuffle para crear una nueva partición física.

Puede equilibrar o aumentar paralelismo, pero añade coste que debe amortizarse.
PySparkReparto consciente
spark.conf.set("spark.sql.shuffle.partitions", 200)
balanced = events.repartition(200, "event_date")

El valor 200 es punto de prueba, no receta; mide tamaños y duración.

Puntos clave

  • Muy pocas particiones limitan paralelismo
  • Demasiadas crean overhead
  • repartition y coalesce no son equivalentes

Evita

  • Copiar un número fijo a cualquier volumen
  • coalesce a 1 antes de cada escritura

Recuerdo activo

¿Qué ajuste controla particiones de shuffle SQL?

Borrador privado · solo en este navegador
03
Operación

Particiones y paralelismo

Broadcast evita mover el lado grande cuando el otro cabe con seguridad en ejecutores.

Objetivo
Broadcast evita mover el lado grande cuando el otro cabe con seguridad en ejecutores.
Duración estimada
17 min aprox.
Dificultad
Associate + Professional
Prerrequisitos
m04
Reportar un error en esta lección

BroadcastHashJoin distribuye una relación pequeña a cada executor y elimina el shuffle de la grande. El umbral automático o broadcast hint influyen en Catalyst.

Una estimación obsoleta puede emitir un broadcast demasiado grande y causar presión de memoria. SortMergeJoin es razonable para dos lados grandes.

Modelo mental

Un broadcast join cambia quién se mueve. En un shuffle join ambos lados se redistribuyen por clave; en BroadcastHashJoin el lado pequeño se recopila y distribuye a los executors para construir una tabla hash local, de modo que las particiones del lado grande pueden leerse donde están. La ventaja depende del tamaño serializado real y de la memoria disponible en cada executor, no de que el dataset se llame dimensión. Catalyst puede elegir broadcast con estadísticas y umbral, y un hint puede influir, pero forzar una relación creciente puede causar OOM. Cuando ambos lados son grandes, sort-merge suele ser una estrategia razonable aunque implique shuffle.

BroadcastHashJoin

Join que replica una relación pequeña y construye una tabla hash local en cada executor.

Evita redistribuir el lado grande cuando la memoria soporta la réplica.
Estadística de tamaño

Estimación del volumen de una relación usada por el optimizador para comparar estrategias.

Una estimación obsoleta puede conducir a un plan físico inadecuado.
Hint

Indicación declarativa que influye en la estrategia seleccionada por Catalyst.

Debe usarse con evidencia porque puede imponerse sobre una decisión adaptativa más segura.
PySparkBroadcast explícito
from pyspark.sql.functions import broadcast
enriched = facts.join(broadcast(dim.select("id", "segment")), "id")

Confirma el tamaño real de dim y el operador del plan.

Puntos clave

  • Broadcast mueve el lado pequeño
  • Las estadísticas importan
  • Dos lados grandes suelen requerir shuffle

Evita

  • Broadcast de una dimensión no acotada
  • Forzar hint sin revisar memoria

Recuerdo activo

¿Qué dataset se replica en BroadcastHashJoin?

Borrador privado · solo en este navegador
04
Diagnóstico

Shuffles, skew y spills

Skew, shuffle y spill son síntomas diferentes y requieren evidencia distinta.

Objetivo
Skew, shuffle y spill son síntomas diferentes y requieren evidencia distinta.
Duración estimada
17 min aprox.
Dificultad
Associate + Professional
Prerrequisitos
m04
Reportar un error en esta lección

Skew aparece cuando pocas tareas duran o leen mucho más que el resto. AQE puede dividir particiones sesgadas, pero una clave nula dominante o diseño incorrecto puede exigir filtrado, salting o preagregación.

Spill indica que una operación usa disco por falta de memoria de ejecución; OOM puede proceder de collect, broadcast excesivo o particiones enormes. Añadir memoria sin corregir el patrón solo pospone el fallo.

Modelo mental

Shuffle describe movimiento de datos entre particiones; skew describe desigualdad; spill describe uso de disco cuando el estado en memoria no cabe. Pueden coexistir, pero no son sinónimos. Un shuffle grande puede estar equilibrado y terminar correctamente. Una clave dominante produce una o pocas tareas largas aunque el resto termine pronto. Spill generalizado puede indicar particiones demasiado grandes, agregaciones extensas o memoria insuficiente. AQE adapta particiones y puede mitigar algunos joins sesgados, pero no corrige una clave nula que debería tratarse semánticamente ni una llamada collect que desborda el driver. El diagnóstico compara distribución por tarea y conecta cada síntoma con el operador del plan.

Skew

Distribución extrema en la que unas pocas particiones contienen mucho más trabajo que sus pares.

Explica colas largas que no mejoran añadiendo workers.
Spill

Escritura temporal a disco de estado de ejecución que no cabe en la memoria disponible.

Indica presión de memoria, aunque su severidad debe cuantificarse.
AQE

Adaptive Query Execution, mecanismo que revisa decisiones físicas con estadísticas observadas durante el run.

Puede fusionar particiones o mitigar skew sin cambiar la lógica declarada.
SQLDetectar claves dominantes
SELECT join_key, count(*) AS rows
FROM main.silver.events
GROUP BY join_key
ORDER BY rows DESC
LIMIT 20;

Compara con la distribución de duración por tarea en Spark UI.

Puntos clave

  • Skew es distribución desigual
  • Spill no equivale siempre a fallo
  • AQE adapta el plan en runtime

Evita

  • Desactivar AQE sin causa
  • Resolver OOM aumentando driver si falla un executor

Recuerdo activo

¿Qué señal distingue skew?

Borrador privado · solo en este navegador
05
Decisión de diseño

Estrategias de join y broadcast

El diagnóstico eficaz separa fallos de arranque, librerías, driver y ejecutores.

Objetivo
El diagnóstico eficaz separa fallos de arranque, librerías, driver y ejecutores.
Duración estimada
17 min aprox.
Dificultad
Associate + Professional
Prerrequisitos
m04
Reportar un error en esta lección

Un cluster que no inicia se investiga en eventos de compute y configuración; un import error en logs de tarea apunta a dependencias; driver OOM suele relacionarse con collect o metadatos.

Parte del mensaje exacto, correlaciona run, task y cluster, y reproduce con el menor input. Cambiar simultáneamente runtime, nodos y código destruye la evidencia.

Modelo mental

Diagnosticar empieza por localizar la fase del fallo: aprovisionamiento, inicialización, carga de dependencias, planificación, ejecución distribuida o commit de salida. Un cluster que nunca alcanza RUNNING no tiene todavía stages útiles; se investiga en eventos y configuración. Un ModuleNotFoundError pertenece al entorno de la tarea. Driver OOM y executor OOM apuntan a memorias y patrones distintos. Una consulta lenta exige plan, Spark UI o Query Profile. El mensaje exacto, run_id, task_key, compute y timestamp forman la evidencia mínima. Cambiar runtime, tamaño y código a la vez destruye causalidad. El método reproduce con el menor input que conserva el síntoma y modifica una variable por experimento.

Fase de fallo

Etapa concreta del ciclo del workload en la que deja de progresar correctamente.

Dirige hacia la superficie de evidencia adecuada y evita cambios irrelevantes.
Driver OOM

Agotamiento de memoria en el proceso coordinador, a menudo por resultados locales o metadatos excesivos.

Requiere un diagnóstico diferente del fallo de una tarea en executor.
Reproducción mínima

Caso más pequeño que mantiene la causa y condiciones esenciales del fallo.

Permite probar hipótesis rápidamente sin perder causalidad.
PythonEvitar materializar en driver
# Evita: rows = large_df.collect()
summary = large_df.groupBy("status").count()
display(summary)

Agrega de forma distribuida antes de devolver un resultado pequeño.

Puntos clave

  • Localiza primero la fase del fallo
  • Driver y executor tienen causas distintas
  • Cambia una variable por experimento

Evita

  • Reinstalar librerías sin leer el error
  • Aumentar el driver para un problema de skew

Recuerdo activo

¿Dónde empiezas ante un cluster que nunca arrancó?

Borrador privado · solo en este navegador
5 lecciones pendientes

Vista de lectura · sin ejecución

Learn Databricks

Learn Databricks · commit 08c378c

README.md