Spark
JuniorКак читать Spark event log без открытия Spark UI?
Spark UI — это просто визуализация того же файла. Зная, какие типы событий реально важны, вы получите те же цифры быстрее — или сразу передадите их в автоматизацию.
Spark event log — это файл с построчным JSON: по одному объекту события на строку, записываемый в реальном времени по ходу выполнения задачи. UI, к которому вы привыкли, читает тот же самый файл.
Где лежит файл и что в нём
Установите spark.eventLog.enabled=true и spark.eventLog.dir на что-то надёжное
(локальный диск, S3, HDFS) — и Spark будет писать по одной строке JSON на каждое событие по ходу выполнения
приложения, без дополнительных инструментов. У каждой строки есть поле "Event", называющее её тип:
SparkListenerApplicationStart, SparkListenerJobStart,
SparkListenerStageSubmitted, SparkListenerTaskEnd,
SparkListenerStageCompleted и ещё десяток других, которые почти никогда не нужны.
Три типа событий, которые несут почти всё нужное
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 — и
получите конкретные рекомендации, которые нужно одобрить, а не ещё одно эмпирическое правило.