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.
- Leer explain formatted
- Detectar Exchange, skew y spill
- Elegir broadcast, repartition o coalesce con evidencia
01Modelo mentalLazy evaluation y DAG de Spark
El plan físico y sus métricas indican dónde se mueve y procesa el dato.
+
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
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.
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.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 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.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 navegador02ImplementaciónCatalyst y optimizaciones de consulta
Particiones determinan paralelismo, tamaño de tareas y número de archivos.
+
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
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.
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.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.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.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 navegador03OperaciónParticiones y paralelismo
Broadcast evita mover el lado grande cuando el otro cabe con seguridad en ejecutores.
+
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
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.
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.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.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.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 navegador04DiagnósticoShuffles, skew y spills
Skew, shuffle y spill son síntomas diferentes y requieren evidencia distinta.
+
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
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.
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.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.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.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 navegador05Decisión de diseñoEstrategias de join y broadcast
El diagnóstico eficaz separa fallos de arranque, librerías, driver y ejecutores.
+
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
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.
Etapa concreta del ciclo del workload en la que deja de progresar correctamente.
Dirige hacia la superficie de evidencia adecuada y evita cambios irrelevantes.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.Caso más pequeño que mantiene la causa y condiciones esenciales del fallo.
Permite probar hipótesis rápidamente sin perder causalidad.# 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