Spark

Junior

¿Cómo leo un event log de Spark sin abrir la Spark UI?

La Spark UI no es más que un renderizador de este archivo. Sabiendo qué tipos de evento importan de verdad, puedes obtener las mismas cifras más rápido - o pasarlas directamente a tu automatización.

Leer en:
SparkListenerTaskEnd SparkListenerStageCompleted SparkListenerJobEnd
Qué tipos de evento llevan las métricas que realmente necesitas, aproximadamente en ese orden - no es un desglose medido.

Un event log de Spark es un archivo JSON delimitado por líneas - un objeto de evento por línea, escrito en tiempo real mientras el job se ejecuta. La UI que usas lee exactamente ese mismo archivo.

Dónde vive el archivo, y qué contiene

Configura spark.eventLog.enabled=true y spark.eventLog.dir apuntando a algo duradero (disco local, S3, HDFS) y Spark escribirá una línea de JSON por evento mientras la aplicación se ejecuta - sin herramientas adicionales. Cada línea tiene un campo "Event" que nombra su tipo: SparkListenerApplicationStart, SparkListenerJobStart, SparkListenerStageSubmitted, SparkListenerTaskEnd, SparkListenerStageCompleted, y una docena más que casi ningún job necesita.

# las 10 tareas más largas, directamente del event log cat app-events.log | jq -r 'select(.Event == "SparkListenerTaskEnd") | [."Task Info"."Task ID", ."Task Metrics"."Executor Run Time"] | @tsv' | sort -k2 -n -r | head -10

Los tres tipos de evento que llevan casi todo lo que necesitas

SparkListenerTaskEnd es el más importante: su objeto "Task Metrics" tiene el tiempo de ejecución en el executor, memoryBytesSpilled, diskBytesSpilled, y los bytes de shuffle read/write - por tarea, que es la granularidad que realmente muestra el sesgo (skew). SparkListenerStageCompleted te da el resumen a nivel de stage sin tener que reagregar cada tarea tú mismo, y SparkListenerJobEnd te dice si todo terminó correctamente. Casi cualquier pregunta diagnóstica - ¿este job está haciendo spill?, ¿una tarea es mucho más lenta que el resto?, ¿este stage llegó a terminar? - se puede responder solo con esos tres.

Dónde esto deja de funcionar, y por qué la UI sigue existiendo

Dos límites honestos: el event log se escribe después de los hechos, así que no sirve para un job que todavía se está ejecutando (la vista en vivo de la UI lee el estado en memoria, no el archivo). Y se vuelve grande - un job con 50.000 tareas produce aproximadamente una línea JSON por tarea solo para eventos TaskEnd, así que cargar el archivo entero con json.load() es el enfoque equivocado pasados unos cientos de MB; en su lugar, procésalo línea por línea.

Esto es literalmente lo que hace el motor de reglas de opti-pipe con él: recorre los eventos de fin de tarea en streaming en lugar de cargar el archivo entero, extrae las métricas numéricas de arriba, y nunca toca nada que revele tu SQL real o tu código de DataFrame - el event log no lleva eso de todos modos, solo las operaciones que produjo el propio planificador de DAG de Spark.

Mira cómo se ve esto en tu propio pipeline.

Sube un event log real de Spark, un run_results.json de dbt, o una exportación de métricas de Flink, y recibe recomendaciones concretas que requieren tu aprobación - no otra regla general.