Spark

Junior

Как читать Spark event log без открытия Spark UI?

Spark UI — это просто визуализация того же файла. Зная, какие типы событий реально важны, вы получите те же цифры быстрее — или сразу передадите их в автоматизацию.

Читать на:
SparkListenerTaskEnd SparkListenerStageCompleted SparkListenerJobEnd
Какие типы событий несут нужные метрики, примерно в этом порядке — не измеренная разбивка.

Spark event log — это файл с построчным JSON: по одному объекту события на строку, записываемый в реальном времени по ходу выполнения задачи. UI, к которому вы привыкли, читает тот же самый файл.

Где лежит файл и что в нём

Установите spark.eventLog.enabled=true и spark.eventLog.dir на что-то надёжное (локальный диск, S3, HDFS) — и Spark будет писать по одной строке JSON на каждое событие по ходу выполнения приложения, без дополнительных инструментов. У каждой строки есть поле "Event", называющее её тип: SparkListenerApplicationStart, SparkListenerJobStart, SparkListenerStageSubmitted, SparkListenerTaskEnd, SparkListenerStageCompleted и ещё десяток других, которые почти никогда не нужны.

# 10 самых долгих задач прямо из 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

Три типа событий, которые несут почти всё нужное

SparkListenerTaskEnd — самый важный: его объект "Task Metrics" содержит время выполнения на executor'е, memoryBytesSpilled, diskBytesSpilled и объёмы shuffle read/write — на уровне задачи, а это как раз та детализация, которая показывает перекос. SparkListenerStageCompleted даёт сводку по стейджу без необходимости агрегировать все задачи самостоятельно, а SparkListenerJobEnd говорит, завершилось ли всё успешно. Почти любой диагностический вопрос — сбрасывается ли задача на диск, сильно ли одна задача медленнее остальных, завершился ли вообще этот стейдж — можно решить, используя только эти три.

Где это не работает, и почему UI всё ещё нужен

Два честных ограничения: event log записывается постфактум, поэтому он бесполезен для ещё выполняющейся задачи (живое представление UI читает состояние в памяти, а не файл). И он бывает большим — задача с 50 000 задач-тасков даёт примерно одну JSON-строку на таск только для событий TaskEnd, так что загружать весь файл через json.load() — неверный подход после нескольких сотен МБ; вместо этого читайте построчно.

Именно так с этим работает движок правил opti-pipe: он проходит по событиям окончания задач потоково, а не загружает файл целиком, извлекает числовые метрики выше и никогда не касается ничего, что раскрыло бы ваш реальный SQL или код DataFrame — event log этого и не содержит, только операции, сгенерированные собственным DAG-планировщиком Spark.

Посмотрите, как это выглядит на вашем собственном пайплайне.

Загрузите реальный event log Spark, run_results.json от dbt или экспорт метрик Flink — и получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.