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.
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.
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.