History Server: настройка event log на S3/HDFS и retention политика

Полный разбор Spark Event Log инфраструктуры: архитектура LiveListenerBus, конфигурация записи логов на S3/MinIO и HDFS, настройка демона Spark History Server, политики retention и ротация логов длинных стриминг-джобов, S3 Lifecycle Rules против встроенного клинера.

platform

Предыдущие четыре урока этого модуля посвящены чтению Spark UI в реальном времени - во время активного выполнения джобы. Но в production большинство расследований происходит постфактум: инженер получает алерт в 03:17, ночная ETL-джоба уже упала, кластер уже уничтожен, а Spark UI пропал вместе с ним. Именно для таких случаев и существует Spark History Server - автономный HTTP-сервис, который сохраняет полную историю выполнения завершённых приложений и позволяет просматривать тот же интерфейс Spark UI спустя дни, недели и месяцы после окончания работы джобы.

Этот урок разбирает не просто «как нажать кнопки в конфиге», а полную архитектуру системы: откуда берутся event logs, почему их нельзя хранить на локальном диске, как правильно настроить запись в S3 или HDFS, как сконфигурировать сам History Server, и как не дать системе захлебнуться через полгода эксплуатации.


Что такое Spark History Server и зачем он нужен

Проблема: Spark UI умирает вместе с приложением

Когда Spark-приложение запущено, его Web UI доступен по адресу вида http://driver-host:4040. Там можно видеть текущее состояние джобов, stage'ов, executor'ов, SQL-планы - всё, что разбиралось в предыдущих четырёх уроках. Но этот UI существует исключительно в памяти процесса драйвера: как только SparkContext завершается (штатно или аварийно), Web UI на порту 4040 немедленно пропадает. Нет никакого способа вернуться к нему позже.

Это фундаментальная проблема для production-эксплуатации. Инженер данных не сидит за экраном во время каждого ночного запуска ETL - он обнаруживает проблему утром по алерту или жалобе. К этому моменту джоба давно завершилась, кластер мог быть уничтожен, и единственный вопрос, который нужно ответить - «что именно пошло не так» - остаётся без ответа, если не было предусмотрено средств для сохранения истории.

Что такое Spark History Server

Spark History Server - это отдельный, постоянно работающий HTTP-сервис, включённый в дистрибутив Apache Spark. Он не является частью ни драйвера, ни executor'ов: это самостоятельная программа, которая:

  • Читает event log файлы из центрального хранилища (HDFS, S3 или другой файловой системы)
  • Парсит их и восстанавливает полную картину выполнения каждого приложения
  • Предоставляет тот же самый веб-интерфейс, что и живой Spark UI, - с теми же вкладками Jobs, Stages, Storage, Executors, SQL, Environment

По умолчанию History Server слушает порт 18080 и хранит UI одновременно для многих завершённых приложений. Пользователь видит список всех приложений и может открыть любое - совершенно независимо от того, работает ли оно сейчас или завершилось неделю назад.

Для каких задач используется History Server

Постфактум-расследование инцидентов. Это главный сценарий. Ночная ETL-джоба упала в 03:17. Утром инженер открывает History Server, находит приложение по имени или времени запуска, смотрит вкладку Stages - и за две минуты определяет, в каком именно stage упала задача, было ли это OOM на конкретном executor'е или проблема с данными на конкретном партишне.

Анализ деградации производительности. История Server хранит метрики по всем джобам. Если время выполнения ETL выросло с 20 минут до 45 минут за последние две недели, инженер может сравнить два запуска: открыть приложение от понедельника прошлой недели и приложение от сегодня, найти разницу в stage-метриках и установить, что стало причиной замедления (например, утроилось shuffle write на одном из stage'ов из-за изменившегося распределения данных).

Capacity planning и оптимизация. Анализ реального потребления ресурсов (executor cores, память, shuffle объём) по историческим данным помогает правильно подбирать параметры кластера для новых или изменившихся джоб. REST API History Server позволяет делать это программно.

Аудит и compliance. В некоторых организациях требуется фиксировать, кто и что запускал, сколько ресурсов потребил, когда завершилось. History Server хранит эту информацию в structured-виде, доступном через API.

Командная разработка и code review. Один разработчик делает оптимизацию, запускает джобу, отправляет коллеге ссылку на приложение в History Server: http://history-server:18080/history/application_001/jobs/. Коллега открывает тот же UI и видит точно те же метрики - не скриншот, а живой интерактивный интерфейс.

Как History Server соотносится с Spark UI и другими инструментами мониторинга

Важно чётко понимать место History Server в общей экосистеме наблюдаемости:

Инструмент Что показывает Когда доступен Тип данных
Spark UI (порт 4040) Текущее состояние живого приложения Только пока запущен драйвер Реальное время
Spark History Server (порт 18080) Полная история завершённых приложений Постоянно, независимо от приложений Исторические данные по приложениям
Prometheus + Grafana Агрегированные метрики кластера в виде временных рядов Постоянно Метрики с временной меткой
SparkMeasure Детальные метрики stage/task на уровне кода В рантайме и постфактум через файл Структурированные метрики

History Server - это не замена Prometheus/Grafana: там метрики агрегированы по всему кластеру и показывают «сколько памяти использовалось в 14:30». History Server показывает «в этом конкретном приложении, запущенном в 14:00, stage 5 занял 28 минут, а task 142 пережил OOM-kill и был перезапущен». Это разные уровни детализации, которые дополняют друг друга.

Как History Server восстанавливает UI из файла

Механизм работы прозрачен, если понять его шаги:

  1. Spark-приложение во время выполнения записывает каждое событие (start job, end stage, end task...) в event log файл на общем хранилище. Один файл на одно приложение. Файл растёт построчно.

  2. History Server периодически сканирует директорию логов и обнаруживает новые завершённые файлы (по наличию строки SparkListenerApplicationEnd в конце файла).

  3. При первом обращении пользователя к конкретному приложению History Server читает весь файл, воспроизводит его событие за событием (replay), строя в памяти полную модель состояния: какие job'ы были, какие stage'ы, сколько задач, какие метрики. Это называется log replay.

  4. Результат replay сохраняется в локальный кеш на диске сервера (LevelDB/RocksDB). При следующих обращениях данные читаются уже из кеша, без повторного парсинга файла.

  5. Пользователь видит в браузере тот же интерфейс, что и в живом Spark UI - потому что History Server использует абсолютно тот же frontend-код.

Именно поэтому качество данных в History Server напрямую зависит от того, насколько полно и корректно был записан event log. Если при записи логов возникли ошибки (потеря сети, нехватка места), часть событий может быть потеряна, и соответствующие метрики в UI будут неполными.


Архитектура Event Log: жизненный цикл Spark-события

Как рождается event внутри Spark

Каждое значимое событие в жизни Spark-приложения - запуск job'а, начало и конец stage, создание и завершение task, изменение состояния executor'а - порождает событие (event) внутри JVM драйвера. Это не метрика и не лог-строка в традиционном смысле, а структурированный объект одного из нескольких десятков типов: SparkListenerJobStart, SparkListenerStageCompleted, SparkListenerTaskEnd, SparkListenerExecutorAdded и так далее.

Внутри драйвера события обрабатывает компонент LiveListenerBus - шина асинхронной доставки событий (event bus). Когда TaskScheduler завершает задачу, когда DAGScheduler открывает новый Stage, когда HeartbeatReceiver получает сигнал от executor'а - все они постят событие в очередь LiveListenerBus. Оттуда события доставляются параллельно всем зарегистрированным слушателям (listeners) в отдельных потоках.

Ключевой компонент для нас - EventLoggingListener. Это и есть тот слушатель, который при включённом spark.eventLog.enabled = true перехватывает каждое событие, сериализует его в JSON и дописывает строку в event log файл. Формат предельно прост: один JSON-объект на строку, файл растёт по ходу выполнения приложения.

Физический формат event log файла

Event log - это обычный текстовый файл (по умолчанию без расширения), где каждая строка - самостоятельный JSON-объект с полем Event, указывающим тип события. Первой строкой всегда идёт SparkListenerLogStart с версией Spark:

{"Event":"SparkListenerLogStart","Spark Version":"3.5.1"}
{"Event":"SparkListenerApplicationStart","App Name":"ETL_orders_pipeline","App ID":"application_1234567890_0001","Timestamp":1718000000000,"User":"spark"}
{"Event":"SparkListenerEnvironmentUpdate","JVM Information":{...},"Spark Properties":{...},"System Properties":{...},"Classpath Entries":{...}}
{"Event":"SparkListenerExecutorAdded","Timestamp":1718000001000,"Executor ID":"1","Executor Info":{...}}
{"Event":"SparkListenerJobStart","Job ID":0,"Submission Time":1718000002000,"Stage IDs":[0,1,2],"Properties":{...}}
{"Event":"SparkListenerStageSubmitted","Stage Info":{"Stage ID":0,"Stage Attempt ID":0,"Stage Name":"read at ETLJob.py:42","Number of Tasks":200,...}}
...
{"Event":"SparkListenerTaskEnd","Stage ID":0,"Stage Attempt ID":0,"Task Type":"ResultTask","Task End Reason":{"Reason":"Success"},"Task Info":{...},"Task Executor Metrics":{...},"Task Metrics":{...}}
...
{"Event":"SparkListenerApplicationEnd","Timestamp":1718003600000}

Последняя строка SparkListenerApplicationEnd - сигнал для History Server о том, что приложение завершилось корректно. Если этой строки нет (драйвер упал, OOM, kill -9), History Server всё равно отобразит приложение, но пометит его как incomplete (незавершённое).

Асинхронный буфер: почему запись логов не тормозит основной код

Важнейшее архитектурное свойство EventLoggingListener: он работает в отдельном потоке с буферизованным выводом, полностью изолированным от основного потока выполнения DAGScheduler и TaskScheduler. Это означает, что даже если сеть до S3 временно деградировала или MinIO отвечает с задержкой в 200ms, эта задержка не влияет на выполнение задач executor'ами. Буфер поглощает всплески, и события доставляются в хранилище по мере возможности.

Размер этого буфера настраивается параметром spark.eventLog.buffer.kb (по умолчанию 100 КБ). На высоко-нагруженных кластерах с тысячами task'ов в секунду рекомендуется увеличивать его до 1024-4096 КБ, чтобы уменьшить частоту операций PUT в объектное хранилище и избежать throttling.


Почему event logs нельзя хранить локально

Эфемерная природа современных кластеров

В эпоху Kubernetes, AWS EMR, Yandex DataProc и других managed-платформ концепция «кластер - это постоянная инфраструктура» устарела. Современный Spark-кластер - это эфемерная структура: он создаётся под конкретную джобу или пачку джоб, выполняет работу и уничтожается. Pod'ы завершаются, виртуальные машины возвращаются в пул, локальные диски отчищаются.

Если event log хранится на локальном диске ноды (например, в /tmp/spark-events/), то вместе с удалением ноды или pod'а исчезает и вся история выполнения джобы. При этом важно понимать: Spark History Server - это отдельный процесс, который читает логи из файловой системы. Если History Server живёт на одной машине с драйвером (что иногда делают в учебных целях), то после уничтожения этой машины теряется всё.

Решение всегда одно: писать event logs в централизованное хранилище, доступное независимо от жизненного цикла кластера. На практике это:

  • HDFS - для on-premise кластеров с постоянными Hadoop-инсталляциями
  • S3 / S3-совместимые хранилища (MinIO, Ceph RGW, Yandex Object Storage) - для cloud и Kubernetes-окружений
  • GCS / Azure Blob Storage - для соответствующих облаков

Конфигурация записи event logs в S3/MinIO

Базовые параметры

Минимальная конфигурация для включения event logging в S3:

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("ETL_orders_pipeline")

    # Включение event logging
    .config("spark.eventLog.enabled", "true")
    .config("spark.eventLog.dir", "s3a://my-data-bucket/spark-events")

    # S3A connector - доступ к MinIO/S3
    .config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000")
    .config("spark.hadoop.fs.s3a.access.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.secret.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.path.style.access", "true")
    .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")

    .getOrCreate()
)

Параметр spark.eventLog.dir должен указывать на уже существующую директорию в хранилище. Spark не создаёт её автоматически - если директория отсутствует, при запуске приложения появится ошибка FileNotFoundException и event logging будет отключён (само приложение при этом продолжит работу).

Сжатие логов: зачем и как

Event log файл крупного ETL-приложения с тысячами задач запросто вырастает до 500 МБ - 2 ГБ несжатого JSON. Каждое событие содержит полные метаданные executor'а, детальные метрики каждой задачи, informação о shuffle. Это важная диагностическая информация, но платить за её хранение и передачу без сжатия расточительно.

Параметр spark.eventLog.compress = true включает сжатие «на лету»: EventLoggingListener пропускает JSON-строки через кодек перед записью в хранилище. По умолчанию используется lz4 - он обеспечивает хороший баланс между скоростью сжатия (почти не нагружает CPU драйвера) и степенью сжатия (как правило 5-10x для JSON, то есть 2 ГБ превращаются в 200-400 МБ).

# Дополнительно к базовой конфигурации:

.config("spark.eventLog.compress", "true")
# Кодек: lz4 (по умолчанию, быстрый) или zstd (чуть медленнее, лучше сжатие)
.config("spark.eventLog.compression.codec", "lz4")
# Размер буфера записи: увеличить на высоконагруженных кластерах
.config("spark.eventLog.buffer.kb", "1024")

Сжатый файл имеет расширение .lz4 или .zstd. History Server автоматически распознаёт формат и декомпрессирует при чтении. Это важно иметь в виду: History Server при открытии лога должен будет декомпрессировать весь файл, что добавляет время при первом открытии очень крупных логов.

Кодек zstd (Zstandard) появился в Spark 3.x и даёт лучшее соотношение сжатие/скорость по сравнению с lz4, но требует наличия соответствующей нативной библиотеки в JVM. На большинстве современных Spark-дистрибутивов она включена по умолчанию.

Буфер записи и проблема S3 throttling

Объектное хранилище (S3, MinIO) принципиально отличается от файловой системы: у него нет понятия «открытый файловый дескриптор с буферизованной записью». Каждый PUT объекта - отдельная HTTP-операция. Если буфер EventLoggingListener слишком мал, он будет делать сотни маленьких PUT-запросов в секунду, что:

  1. Создаёт нагрузку на S3 API (у AWS есть лимит ~3500 PUT/сек на prefix, после чего начинается throttling с HTTP 503)
  2. Увеличивает стоимость (каждый PUT-запрос тарифицируется)
  3. Создаёт ненужный сетевой overhead на драйвере

Параметр spark.eventLog.buffer.kb задаёт размер in-memory буфера в килобайтах. При заполнении буфера EventLoggingListener делает один PUT-запрос. На высоконагруженных кластерах разумный диапазон - 512 КБ - 4096 КБ:

.config("spark.eventLog.buffer.kb", "4096")  # 4 МБ буфер - для больших кластеров

Аутентификация без хардкода секретов

Хардкод ключей в конфигурации SparkSession допустим только в локальной разработке. В production используются следующие подходы:

Для AWS EMR / EC2 с IAM-ролями - никаких ключей не нужно, S3A автоматически подхватывает IAM Instance Profile:

.config("spark.hadoop.fs.s3a.aws.credentials.provider",
        "com.amazonaws.auth.InstanceProfileCredentialsProvider")

Для Kubernetes с Service Account и IRSA (IAM Roles for Service Accounts):

.config("spark.hadoop.fs.s3a.aws.credentials.provider",
        "com.amazonaws.auth.WebIdentityTokenCredentialsProvider")

Для MinIO / self-hosted S3 - через переменные окружения. S3A умеет читать AWS_ACCESS_KEY_ID и AWS_SECRET_ACCESS_KEY из окружения без указания в конфиге:

export AWS_ACCESS_KEY_ID="minioadmin"
export AWS_SECRET_ACCESS_KEY="minioadmin"
.config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000")
.config("spark.hadoop.fs.s3a.path.style.access", "true")
# Ключи не указываем - берутся из переменных окружения

Через Kubernetes Secrets - монтируем Secret как файл или переменные окружения в pod, а конфиг читается через spark.hadoop.fs.s3a.access.key из ${env:MINIO_ACCESS_KEY}.

Полная production-конфигурация для MinIO на Kubernetes

Вот реалистичный пример SparkSession-конфигурации, которую можно встретить в production Kubernetes-окружении с MinIO:

spark = (
    SparkSession.builder
    .appName("orders-etl-daily")

    # === Event Log ===
    .config("spark.eventLog.enabled", "true")
    .config("spark.eventLog.dir", "s3a://spark-platform/event-logs")
    .config("spark.eventLog.compress", "true")
    .config("spark.eventLog.compression.codec", "zstd")
    .config("spark.eventLog.buffer.kb", "2048")

    # === S3A / MinIO ===
    .config("spark.hadoop.fs.s3a.endpoint", "http://minio-svc.storage.svc.cluster.local:9000")
    .config("spark.hadoop.fs.s3a.access.key", "${env:MINIO_ACCESS_KEY}")
    .config("spark.hadoop.fs.s3a.secret.key", "${env:MINIO_SECRET_KEY}")
    .config("spark.hadoop.fs.s3a.path.style.access", "true")
    .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
    # Multipart upload для крупных файлов (включает нормальные append-операции на S3)
    .config("spark.hadoop.fs.s3a.multipart.size", "67108864")    # 64 МБ часть
    .config("spark.hadoop.fs.s3a.fast.upload", "true")
    # Пул соединений для параллельных запросов с executor'ов
    .config("spark.hadoop.fs.s3a.connection.maximum", "100")
    .config("spark.hadoop.fs.s3a.threads.max", "20")

    .getOrCreate()
)

Конфигурация записи event logs в HDFS

Специфика on-premise стека

Для кластеров, где Spark работает поверх HDFS (Hadoop on-premise, CDH, HDP), event logging через hdfs:// протокол является стандартным подходом. В отличие от S3, HDFS поддерживает настоящий append к существующему файлу, что упрощает буферизацию - EventLoggingListener может держать файл открытым и дописывать в него по мере поступления событий.

spark = (
    SparkSession.builder
    .appName("orders-etl-daily")
    .config("spark.eventLog.enabled", "true")
    .config("spark.eventLog.dir", "hdfs://namenode:8020/var/log/spark/apps")
    .config("spark.eventLog.compress", "true")
    .getOrCreate()
)

Если HDFS настроен с HA (High Availability) через nameservice, то используется логическое имя:

.config("spark.eventLog.dir", "hdfs://hadoop-cluster/var/log/spark/apps")
# + соответствующие hadoop.conf параметры для HA resolution

Права доступа: почему важен sticky bit

На HDFS-кластере event logs от разных пользователей и приложений записываются в одну общую директорию. Если директория /var/log/spark/apps доступна всем на запись без sticky bit, любой пользователь сможет удалить чужие логи - намеренно или случайно.

Правильная настройка прав доступа:

# Создаём директорию
hdfs dfs -mkdir -p /var/log/spark/apps

# Устанавливаем sticky bit и широкие права (аналог chmod 1777 на Linux)
hdfs dfs -chmod 1777 /var/log/spark/apps

# Проверяем
hdfs dfs -ls -d /var/log/spark/
# drwxrwxrwt   - spark supergroup    0 2024-06-01 00:00 /var/log/spark/apps

Sticky bit (1 в начале маски прав) означает, что даже при правах rwxrwxrwx, удалить файл из директории может только его владелец (или root/superuser HDFS). Это стандартная практика для shared-директорий, применяемая с тех же причин, что и /tmp в Linux.

Параметры отказоустойчивости HDFS-клиента

При кратковременных сетевых проблемах между Spark-драйвером и HDFS NameNode EventLoggingListener может получить IOException. По умолчанию Spark при таком сбое молча прекращает запись событий и продолжает работу (что лучше, чем падение всего приложения). Этим поведением управляет параметр:

# Если logging провалился - продолжить приложение без логирования
.config("spark.extraListeners", "")  # EventLoggingListener добавляется автоматически

Более тонкая настройка - параметры HDFS FailoverProxy и retry:

# Количество ретраев при недоступности NameNode
.config("spark.hadoop.ipc.client.connect.max.retries", "10")
.config("spark.hadoop.ipc.client.connect.retry.interval", "1000")  # мс между ретраями

Настройка демона Spark History Server

История Server как отдельный долгоживущий процесс

Spark History Server - это отдельная JVM-программа, запускаемая скриптом $SPARK_HOME/sbin/start-history-server.sh. Она не зависит от наличия запущенных Spark-приложений и работает постоянно - читает event log файлы из хранилища, парсит их и предоставляет веб-интерфейс (по умолчанию на порту 18080).

В Kubernetes он деплоится как обычный Deployment с одним pod'ом и Service поверх него. В on-premise окружениях - как системный сервис (systemd unit) на выделенной машине.

Конфигурация History Server хранится в файле spark-defaults.conf на сервере (или передаётся через переменные окружения):

# /etc/spark/conf/spark-defaults.conf на машине History Server

# Откуда читаем логи
spark.history.fs.logDirectory=s3a://spark-platform/event-logs

# Как часто сканируем директорию на предмет новых логов
spark.history.fs.update.interval=30s

# Максимум приложений в кеше UI (остальные перечитываются с диска по запросу)
spark.history.retainedApplications=200

# Путь к локальному кешу метаданных (RocksDB/LevelDB)
spark.history.store.path=/var/spark-history-store

# Максимальный размер кеша метаданных
spark.history.store.maxDiskUsage=10g

# S3A конфигурация (если логи на S3)
spark.hadoop.fs.s3a.endpoint=http://minio:9000
spark.hadoop.fs.s3a.access.key=minioadmin
spark.hadoop.fs.s3a.secret.key=minioadmin
spark.hadoop.fs.s3a.path.style.access=true
spark.hadoop.fs.s3a.impl=org.apache.hadoop.fs.s3a.S3AFileSystem

Цикл обновления: баланс между свежестью и нагрузкой на S3

Параметр spark.history.fs.update.interval определяет, как часто History Server сканирует директорию логов на предмет новых завершённых приложений. Это критичный параметр с нетривиальным tradeoff:

Слишком маленький интервал (например, 5s) - History Server будет делать S3 LIST-запрос каждые 5 секунд. При тысячах файлов в бакете каждый LIST возвращает постраничный результат, и у AWS S3 есть лимиты на LIST (~5500 LIST-запросов/сек на prefix). Постоянное сканирование тысяч файлов:

  • Создаёт лишние S3 API-расходы (LIST-запросы тарифицируются)
  • Может конкурировать с другими процессами, читающими тот же бакет
  • При большой истории логов сканирование самого списка занимает секунды

Слишком большой интервал (например, 5m) - инженер, ожидающий результатов только что завершившейся джобы, будет 5 минут смотреть на «Приложение не найдено».

Практическая рекомендация: для большинства production-кластеров разумны значения 30s - 60s. Если History Server работает с HDFS (а не S3), можно смело ставить 10s - 15s, потому что HDFS LIST на одну директорию значительно дешевле, чем S3 LIST с тысячами объектов.

Локальный кеш метаданных: зачем нужен RocksDB/LevelDB

Без кеша History Server был бы обязан каждый раз, когда пользователь открывает страницу приложения, заново скачивать и парсить весь event log файл. Для большого приложения с 1 ГБ сжатых логов это означало бы 10-30 секунд ожидания на каждый клик.

История Server использует локальный key-value store (в Spark 3.x - это LevelDB, начиная с Spark 3.3+ опционально RocksDB) для хранения распарсенных метаданных приложений. После первичного парсинга данные сохраняются на локальный SSD сервера и при последующих запросах читаются уже оттуда без обращения к S3.

# Путь к кешу - должен быть на быстром локальном диске (SSD предпочтительно)
spark.history.store.path=/var/spark-history-store

# Максимальный объём кеша на диске
spark.history.store.maxDiskUsage=20g

# Для больших инсталляций можно выбрать RocksDB (требует Spark 3.3+)
spark.history.store.hybridStore.enabled=true
spark.history.store.hybridStore.diskBackend=ROCKSDB

Важная деталь: кеш хранит именно метаданные (список приложений, stage'ов, метрики), но не сырой JSON event log. Если history-server перезапустился, ему не нужно перечитывать все исторические логи - он начнёт показывать приложения из кеша немедленно, и будет дополнять кеш только новыми приложениями, появившимися с момента последнего выключения.

Конфигурация для Docker / Kubernetes деплоя

# kubernetes/history-server-deployment.yaml
apiVersion: apps/v1
kind: Deployment
metadata:
  name: spark-history-server
  namespace: spark-platform
spec:
  replicas: 1
  selector:
    matchLabels:
      app: spark-history-server
  template:
    metadata:
      labels:
        app: spark-history-server
    spec:
      containers:
        - name: history-server
          image: apache/spark:3.5.1
          command: ["/opt/spark/bin/spark-class"]
          args: ["org.apache.spark.deploy.history.HistoryServer"]
          ports:
            - containerPort: 18080
          env:
            - name: SPARK_HISTORY_OPTS
              value: >-
                -Dspark.history.fs.logDirectory=s3a://spark-platform/event-logs
                -Dspark.history.fs.update.interval=30s
                -Dspark.history.store.path=/spark-history-store
                -Dspark.history.store.maxDiskUsage=10g
                -Dspark.hadoop.fs.s3a.endpoint=http://minio-svc:9000
                -Dspark.hadoop.fs.s3a.access.key=$(MINIO_ACCESS_KEY)
                -Dspark.hadoop.fs.s3a.secret.key=$(MINIO_SECRET_KEY)
                -Dspark.hadoop.fs.s3a.path.style.access=true
            - name: MINIO_ACCESS_KEY
              valueFrom:
                secretKeyRef:
                  name: minio-credentials
                  key: access-key
            - name: MINIO_SECRET_KEY
              valueFrom:
                secretKeyRef:
                  name: minio-credentials
                  key: secret-key
          volumeMounts:
            - name: history-store
              mountPath: /spark-history-store
      volumes:
        - name: history-store
          persistentVolumeClaim:
            claimName: spark-history-store-pvc
---
apiVersion: v1
kind: Service
metadata:
  name: spark-history-server
  namespace: spark-platform
spec:
  selector:
    app: spark-history-server
  ports:
    - port: 18080
      targetPort: 18080
  type: ClusterIP

Обратите внимание: PersistentVolumeClaim для /spark-history-store критически важен. Если использовать emptyDir, при перезапуске pod'а кеш будет потерян и сервер начнёт повторно парсить все логи - это может занять часы на больших инсталляциях.


Политики удержания (Retention) и ротация логов

Проблема бесконечного роста

Без механизма очистки directory с event logs будет бесконечно расти. Представьте: тысяча джоб в день, каждая пишет лог размером 100-500 МБ (несжатый) или 20-100 МБ (сжатый). Через год это:

  • Несжатые: 1000 × 300 МБ × 365 = ~110 ТБ
  • Сжатые (lz4): ~22 ТБ

Это не просто стоимость хранения. С ростом числа файлов в S3-бакете LIST-операции History Server становятся всё медленнее (S3 листинг пагинируется по 1000 объектов на страницу). При миллионе файлов полный LIST занимает сотни секунд. JVM History Server может получить OOM при попытке загрузить в память метаданные миллиона приложений.

Встроенный клинер (Cleaner): настройка автоматической очистки

Spark History Server имеет встроенный механизм очистки старых записей. Важно понимать: клинер удаляет записи из кеша History Server, и опционально может также удалять сами файлы логов из хранилища.

# Включаем автоматическую очистку
spark.history.fs.cleaner.enabled=true

# Как часто запускать очистку
spark.history.fs.cleaner.interval=1d

# Максимальный возраст логов (приложения старше этого срока удаляются)
spark.history.fs.cleaner.maxAge=7d

# Минимальное количество завершённых приложений, которые ВСЕГДА сохраняются
# (даже если они старше maxAge)
spark.history.fs.cleaner.maxNum=10000

Параметр spark.history.fs.cleaner.maxAge=7d означает: любое приложение, завершившееся более 7 дней назад, будет удалено из кеша History Server и из файловой системы логов. 7 дней - хороший практический выбор: обычно инженеры расследуют инциденты в течение нескольких дней после происшествия, и недельная история покрывает большинство сценариев.

Пример расчёта: 7 дней × 1000 джоб/день × 50 МБ (сжатый лог) = 350 ГБ постоянного хранилища. Это управляемый объём для большинства production-инсталляций.

Ротация логов: решение проблемы вечно работающих стриминговых джоб

Обычный ETL имеет чёткое начало и конец - приложение запустилось, обработало батч, завершилось. Event log остаётся финальным снимком.

Совсем иначе с Structured Streaming джобами, которые работают непрерывно неделями и месяцами. У них нет SparkListenerApplicationEnd. Event log растёт всё время работы: каждый micro-batch добавляет события stage start/end, task metrics. За неделю непрерывной работы log легко вырастает до 50-100 ГБ.

Параметры ротации логов в Spark:

# В конфигурации SparkSession стримингового приложения:

# Включаем ротацию
.config("spark.eventLog.rotation.enabled", "true")

# Максимальный размер одного файла лога до ротации
# После достижения этого размера текущий файл закрывается,
# открывается новый с суффиксом _2, _3, ...
.config("spark.eventLog.rotation.maxFileSize", "512m")

С включённой ротацией History Server правильно группирует файлы одного приложения: вы видите одну запись «Streaming Orders Job» с треком всех исторических данных, при этом каждый файл лога ограничен 512 МБ. Клинер может удалять старые фрагменты, сохраняя только актуальные.

Важный нюанс: History Server может показывать только завершённые ротированные файлы. Текущий, ещё незакрытый файл последнего периода виден как «incomplete» - это нормальное поведение для работающего стримингового приложения.


Облачный Retention: S3 Lifecycle Policies против встроенного клинера

Почему встроенный клинер - не лучшее решение для S3

Встроенный клинер Spark History Server работает так: он читает метаданные о старых приложениях из своего кеша, затем вызывает S3 API для удаления файлов. Это создаёт несколько проблем:

  1. Нагрузка на API: удаление 1000 файлов = 1000 DELETE-запросов к S3. Если у вас накопилось 100 000 устаревших файлов, клинер будет работать долго и создавать нагрузку на API-лимиты.

  2. Зависимость от доступности: клинер работает, только пока запущен History Server. Если сервер упал или был отключён на неделю, очистка не происходит.

  3. Ограниченная гибкость: клинер не умеет перекладывать файлы в «холодное» хранилище - он только удаляет.

  4. Single point of knowledge: History Server должен «знать» обо всех файлах, которые нужно удалить, через свой кеш. Это создаёт граничные случаи при сбоях.

S3/MinIO Lifecycle Rules: правильное решение на инфраструктурном уровне

Объектные хранилища имеют встроенный механизм управления жизненным циклом объектов (Lifecycle Rules / Lifecycle Policies). Это серверная логика на стороне самого хранилища - она работает независимо от Spark, History Server и любого другого клиента.

Пример настройки через AWS CLI:

# Создаём файл политики
cat > lifecycle-policy.json << 'EOF'
{
    "Rules": [
        {
            "ID": "spark-event-logs-retention",
            "Status": "Enabled",
            "Filter": {
                "Prefix": "spark-events/"
            },
            "Expiration": {
                "Days": 14
            }
        },
        {
            "ID": "spark-event-logs-archive",
            "Status": "Enabled",
            "Filter": {
                "Prefix": "spark-events/"
            },
            "Transitions": [
                {
                    "Days": 7,
                    "StorageClass": "STANDARD_IA"
                }
            ]
        }
    ]
}
EOF

# Применяем к бакету
aws s3api put-bucket-lifecycle-configuration \
    --bucket spark-platform \
    --lifecycle-configuration file://lifecycle-policy.json

Эта конфигурация делает две вещи:

  1. Через 7 дней файлы автоматически переносятся в S3 Standard-IA (Infrequent Access) - хранение дешевле, чтение дороже, что идеально для логов, к которым редко обращаются
  2. Через 14 дней файлы автоматически удаляются

Для MinIO аналогичный механизм настраивается через mc (MinIO Client):

# Применяем lifecycle policy к MinIO-бакету
mc ilm import myminio/spark-platform << 'EOF'
{
    "Rules": [
        {
            "Expiration": {
                "Days": 14
            },
            "ID": "spark-logs-14d",
            "Filter": {
                "Prefix": "spark-events/"
            },
            "Status": "Enabled"
        }
    ]
}
EOF

# Проверяем
mc ilm ls myminio/spark-platform

Синхронизация: что происходит, когда S3 удаляет файл «за спиной» History Server

Это реальный граничный случай: Lifecycle Rule удалила файл, а History Server всё ещё хранит его метаданные в своём кеше и показывает в UI. Пользователь кликает на приложение - History Server пытается прочитать удалённый лог и получает NoSuchKeyException (или FileNotFoundException через S3A).

По умолчанию History Server не падает при этом, а возвращает ошибку в UI: «Event log file is no longer available». Это нормальное поведение. Чтобы History Server убирал такие «призраки» из кеша при следующем сканировании, важно сохранить включённым встроенный клинер с более длинным maxAge, чем период Lifecycle Rule:

# Если Lifecycle удаляет через 14 дней, ставим cleaner на 13 дней
# Тогда History Server успеет сам убрать запись из кеша до того,
# как S3 удалит файл - и не будет "призраков"
spark.history.fs.cleaner.enabled=true
spark.history.fs.cleaner.maxAge=13d

Таким образом, рекомендуемая архитектура - двухслойная защита:

  1. History Server Cleaner (maxAge = N-1 дней) - убирает метаданные из кеша
  2. S3 Lifecycle Rule (maxAge = N дней) - физически удаляет файлы

Практический демо-блок

Настройка окружения

Для воспроизведения демо нужны: Docker Compose с MinIO и доступ к Spark (можно через pyspark в той же сети).

from pyspark.sql import SparkSession
import time

spark = (
    SparkSession.builder
    .appName("history-server-demo")
    .config("spark.jars.packages",
            "org.apache.hadoop:hadoop-aws:3.3.4,"
            "com.amazonaws:aws-java-sdk-bundle:1.12.262")
    .config("spark.eventLog.enabled", "true")
    .config("spark.eventLog.dir", "s3a://spark-platform/event-logs")
    .config("spark.eventLog.compress", "true")
    .config("spark.eventLog.compression.codec", "zstd")
    .config("spark.eventLog.buffer.kb", "512")
    .config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000")
    .config("spark.hadoop.fs.s3a.access.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.secret.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.path.style.access", "true")
    .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem")
    .getOrCreate()
)

Кейс 1: Наблюдение за ростом event log в реальном времени

from pyspark.sql import functions as F
import subprocess

APP_ID = spark.sparkContext.applicationId
print(f"App ID: {APP_ID}")
print(f"Event log будет записан в: s3a://spark-platform/event-logs/{APP_ID}")

# Генерируем данные и выполняем несколько actions
# - каждый action создаёт новые события в логе
for i in range(5):
    df = spark.range(1_000_000).withColumn("val", F.rand())
    count = df.groupBy((F.col("id") % 10).alias("bucket")).agg(F.sum("val")).count()
    print(f"Итерация {i+1}: {count} агрегированных строк")
    time.sleep(2)

# После каждого Action в event log добавляются:
# SparkListenerJobStart, SparkListenerStageSubmitted (x N),
# SparkListenerTaskStart/End (x число задач),
# SparkListenerStageCompleted, SparkListenerJobEnd

Кейс 2: Проверка размера event log файла

import boto3
from botocore.client import Config

# Подключаемся к MinIO через boto3
s3_client = boto3.client(
    "s3",
    endpoint_url="http://minio:9000",
    aws_access_key_id="minioadmin",
    aws_secret_access_key="minioadmin",
    config=Config(signature_version="s3v4"),
    region_name="us-east-1",
)

# Ищем файл нашего приложения
APP_ID = spark.sparkContext.applicationId
response = s3_client.list_objects_v2(
    Bucket="spark-platform",
    Prefix=f"event-logs/{APP_ID}",
)

for obj in response.get("Contents", []):
    size_mb = obj["Size"] / (1024 * 1024)
    print(f"Файл: {obj['Key']}, Размер: {size_mb:.2f} МБ, Дата: {obj['LastModified']}")

# Пример вывода:
# Файл: event-logs/application_1234_0001.zstd, Размер: 0.43 МБ, Дата: 2024-06-15 10:23:41+00:00

Кейс 3: Имитация ротации логов для стримингового приложения

# Конфигурация с ротацией
spark_streaming = (
    SparkSession.builder
    .appName("streaming-with-rotation")
    .config("spark.eventLog.enabled", "true")
    .config("spark.eventLog.dir", "s3a://spark-platform/event-logs")
    .config("spark.eventLog.compress", "true")
    .config("spark.eventLog.rotation.enabled", "true")
    # Для демо ставим очень маленький размер - 1 МБ
    # В production разумно 256m-1g
    .config("spark.eventLog.rotation.maxFileSize", "1m")
    .config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000")
    .config("spark.hadoop.fs.s3a.access.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.secret.key", "minioadmin")
    .config("spark.hadoop.fs.s3a.path.style.access", "true")
    .getOrCreate()
)

APP_ID = spark_streaming.sparkContext.applicationId

# Генерируем много событий - лог будет разрезан на части
for i in range(100):
    spark_streaming.range(100_000).count()
    if i % 10 == 0:
        # Смотрим, сколько файлов уже создано
        resp = s3_client.list_objects_v2(
            Bucket="spark-platform",
            Prefix=f"event-logs/{APP_ID}",
        )
        files = resp.get("Contents", [])
        print(f"Итерация {i}: {len(files)} файлов ротации")

Кейс 4: Проверка, что History Server видит приложение

import requests
import json

HISTORY_SERVER_URL = "http://spark-history-server:18080"

# REST API History Server для получения списка приложений
response = requests.get(f"{HISTORY_SERVER_URL}/api/v1/applications")
apps = response.json()

print(f"Всего приложений в History Server: {len(apps)}")
for app in apps[:5]:
    print(f"  ID: {app['id']}, Имя: {app['name']}, "
          f"Статус: {'completed' if app['completed'] else 'running'}, "
          f"Длительность: {app.get('duration', 0) // 1000} сек")

# Получаем детали конкретного приложения
app_id = apps[0]["id"]
stages_response = requests.get(f"{HISTORY_SERVER_URL}/api/v1/applications/{app_id}/stages")
stages = stages_response.json()
print(f"\nStages для {app_id}: {len(stages)}")
for stage in stages[:3]:
    print(f"  Stage {stage['stageId']}: {stage['name']}, "
          f"Tasks: {stage['numTasks']}, "
          f"Duration: {stage.get('executorRunTime', 0) // 1000} сек")

History Server предоставляет полноценный REST API (совместимый со Spark UI REST API), который позволяет автоматически анализировать производительность без браузера - строить отчёты, алерты, интеграции с внешними системами мониторинга.

Кейс 5: Очистка старых логов вручную

from datetime import datetime, timedelta, timezone

# Удаляем все event logs старше 7 дней
RETENTION_DAYS = 7
cutoff_date = datetime.now(timezone.utc) - timedelta(days=RETENTION_DAYS)

print(f"Удаляем логи старше: {cutoff_date.isoformat()}")

# Получаем список всех объектов
paginator = s3_client.get_paginator("list_objects_v2")
pages = paginator.paginate(Bucket="spark-platform", Prefix="event-logs/")

to_delete = []
for page in pages:
    for obj in page.get("Contents", []):
        if obj["LastModified"] < cutoff_date:
            to_delete.append({"Key": obj["Key"]})
            print(f"  К удалению: {obj['Key']} ({obj['LastModified'].isoformat()})")

if to_delete:
    # Удаляем пачками по 1000 (лимит S3 API)
    for i in range(0, len(to_delete), 1000):
        batch = to_delete[i:i+1000]
        s3_client.delete_objects(
            Bucket="spark-platform",
            Delete={"Objects": batch, "Quiet": True},
        )
    print(f"\nУдалено {len(to_delete)} файлов")
else:
    print("Нет файлов для удаления")

Производственный кейс

Команда Data Platform крупной e-commerce компании развернула Spark на Kubernetes с MinIO как объектным хранилищем. История Server работал нормально первые два месяца, но на третий месяц начал запускаться по 25-30 минут - инженеры жаловались, что не могут быстро проверить результаты недавних джоб. В конце концов сервер начал получать OutOfMemoryError при старте.

Расследование показало корневую причину: за три месяца в бакете накопилось 850 000 файлов event logs. Ни spark.history.fs.cleaner.enabled, ни Lifecycle Rules настроено не было - «руки не дошли». При каждом старте History Server делал полный LIST бакета (850 000 объектов = 850 страниц по 1000 объектов = 850 LIST-запросов к MinIO), затем пытался загрузить метаданные всех приложений в heap JVM, что и вызывало OOM.

Исправление потребовало трёх шагов:

Шаг 1 - экстренная очистка: написали скрипт (аналог Кейса 5 выше), который удалил все логи старше 30 дней. Объём уменьшился с 850 000 файлов до 120 000. History Server снова запустился нормально.

Шаг 2 - настройка Lifecycle Rules в MinIO:

mc ilm import myminio/spark-platform << 'EOF'
{
  "Rules": [{
    "Expiration": {"Days": 14},
    "ID": "spark-logs-14d",
    "Filter": {"Prefix": "event-logs/"},
    "Status": "Enabled"
  }]
}
EOF

Шаг 3 - настройка History Server Cleaner и кеша:

spark.history.fs.cleaner.enabled=true
spark.history.fs.cleaner.maxAge=13d
spark.history.fs.cleaner.interval=1d
spark.history.store.path=/data/spark-history-store
spark.history.store.maxDiskUsage=20g
spark.history.fs.update.interval=60s

После этих изменений сервер запускался за 40 секунд вместо 25 минут, OOM прекратились, а количество файлов в бакете стабилизировалось на уровне 20 000-25 000 (две недели активных джоб).

Дополнительно команда настроила мониторинг: cron-скрипт раз в сутки запрашивает количество объектов в prefix event-logs/ и отправляет метрику в Prometheus. При превышении 50 000 файлов алерт уходит на дежурного инженера платформы - ещё до того, как History Server начнёт ощутимо деградировать.


Типичные заблуждения

«History Server - это резервный Spark UI, который работает параллельно с основным». Неверно: это разные инструменты. Spark UI работает прямо в процессе драйвера и отражает текущее состояние запущенного приложения в реальном времени. History Server - отдельный сервис, читающий файлы завершённых приложений. Когда приложение завершилось и драйвер остановился, Spark UI умирает вместе с ним, а History Server остаётся.

«Включение event logging замедляет мои джобы». Как правило нет: EventLoggingListener работает в отдельном потоке с буфером. При разумных настройках (достаточный буфер, сжатие, быстрая сеть до хранилища) overhead не превышает 1-2% даже на интенсивных джобах. Проблема возникает только при очень маленьком буфере + медленном хранилище + огромном числе event'ов в секунду (тысячи задач/сек).

«Если History Server не показывает приложение - значит логи не записались». Не обязательно: могло не пройти время сканирования spark.history.fs.update.interval. Подождите 30-60 секунд и обновите страницу. Также проверьте наличие файла в хранилище напрямую (через boto3 или hdfs dfs -ls) - это быстрее, чем перезапускать History Server.

«spark.history.fs.cleaner.maxAge - это то же самое, что spark.history.retainedApplications». Нет: maxAge задаёт максимальный возраст по времени (7 дней, 14 дней), после которого приложение удаляется. retainedApplications задаёт максимальное количество приложений, которые держатся в памяти кеша History Server одновременно. Это дополняющие, а не взаимозаменяющие настройки.

«Сжатие event logs включать не нужно - диск дешевый». Аргумент не в стоимости диска. Сжатые логи передаются по сети в 5-10 раз меньшим объёмом (меньше трафика от драйвера до S3 = меньше latency записи = меньше нагрузки на сеть), и S3/MinIO LIST-операции работают быстрее, когда файлов меньше по размеру. Для крупного кластера это ощутимая экономия.

«Встроенный клинер History Server и S3 Lifecycle - дублирующие решения, нужно выбрать одно». Лучше использовать оба в связке: Cleaner убирает метаданные из кеша сервера (чтобы не было «призраков» в UI), а Lifecycle физически удаляет объекты из S3. Как показывает производственный кейс выше, только одного из них недостаточно.


Производственный чек-лист

  • Event logging включён во всех production Spark-приложениях (spark.eventLog.enabled = true) - не только в batch ETL, но и в стриминговых джобах. Без этого постфактум-расследование инцидента практически невозможно.

  • Директория для логов создана в хранилище до запуска приложения - Spark не создаёт её автоматически и при отсутствии просто молча отключает логирование, не бросая исключение в основной поток.

  • Сжатие логов включено (spark.eventLog.compress = true) с кодеком zstd или lz4. Для нового кластера выбирайте zstd - он лучше по степени сжатия при сопоставимой скорости.

  • Буфер записи достаточен (spark.eventLog.buffer.kb >= 512) - на высоконагруженных кластерах с тысячами task'ов значение 1024-4096 КБ снижает количество S3 PUT-запросов и риск throttling.

  • Ротация логов включена для стриминговых джоб (spark.eventLog.rotation.enabled = true, spark.eventLog.rotation.maxFileSize = 256m или 512m) - без ротации стриминговый лог вырастет до десятков гигабайт и сделает History Server неработоспособным для этого приложения.

  • Retention настроен двухслойно: History Server Cleaner (maxAge = N-1 дней) + S3/MinIO Lifecycle Rule (Expiration: N дней). Значение N выбирается исходя из реальных потребностей расследования инцидентов (типично 14-30 дней).

  • Кеш метаданных History Server хранится на Persistent Volume (PVC в Kubernetes, отдельный диск в on-premise). При потере кеша сервер вынужден перепарсировать все логи заново.

  • spark.history.fs.update.interval настроен разумно: 30-60 секунд для S3, 10-30 секунд для HDFS. Значения менее 15 секунд для S3 оправданы только если количество файлов в бакете не превышает нескольких тысяч.

  • Мониторинг количества файлов в event log директории настроен с алертом. Рост числа объектов - ранний индикатор проблем с retention, обнаруживаемый до деградации History Server.

  • REST API History Server используется в CI/CD или post-job-отчётах для автоматической проверки производительности - не только браузерный UI. GET /api/v1/applications/{id}/stages возвращает полные метрики stage'ов в JSON.


Мостик к следующему уроку

История Server решает задачу постфактум-анализа завершённых приложений. Но в production нужно также понимать, что происходит с кластером прямо сейчас - сколько памяти используют executor'ы, каков garbage collection overhead в данный момент, не накапливается ли shuffle spill. Для этого нужна метрическая система с временными рядами и алертингом.

Следующий урок разбирает Spark Metrics System - механизм, по которому Spark экспортирует численные показатели (не события, а метрики) во внешние системы: JMX, Prometheus, Graphite. Это нижний уровень наблюдаемости, комплементарный History Server: там где History Server показывает «что было», Metrics System показывает «что есть сейчас» и «как менялось во времени».


Домашнее задание

  1. Разверните MinIO локально через Docker Compose и настройте запись event logs в бакет spark-events. Создайте директорию spark-events в MinIO заранее. Запустите тестовую PySpark-джобу и убедитесь, что файл лога появился в бакете - проверьте через MinIO Console или boto3.

  2. Включите сжатие event logs через zstd и запустите ту же джобу дважды - один раз без сжатия, один раз с. Сравните размеры файлов. Ответьте: какой процент сжатия получился? Попробуйте также lz4 и сравните с zstd.

  3. Напишите скрипт на Python (используя boto3), который сканирует бакет spark-events и выводит сводную таблицу: количество event log файлов по дням за последние 7 дней, суммарный размер за каждый день. Это имитация простого мониторинга роста логов.

  4. Настройте и запустите Spark History Server в Docker-контейнере, указав ему на тот же MinIO-бакет. Откройте UI на порту 18080 и убедитесь, что завершённые джобы из заданий 1-2 отображаются. Изучите вкладки Jobs, Stages, Executors для одного из приложений.

  5. Включите ротацию логов (spark.eventLog.rotation.enabled = true, spark.eventLog.rotation.maxFileSize = 1m для учебного эффекта) и запустите джобу с 500+ Action'ами. Проверьте в MinIO, что один App ID породил несколько файлов с суффиксами. Убедитесь, что History Server всё равно правильно группирует их под одним приложением.

  6. Настройте Lifecycle Policy в MinIO на удаление объектов старше 1 дня (используйте mc ilm import). Дождитесь срабатывания policy (проверить через mc ilm ls). Объясните, почему History Server может временно показывать «призраки» приложений, чьи файлы уже удалены, и как настроить Cleaner, чтобы это минимизировать.

  7. Используя REST API History Server (GET /api/v1/applications), напишите скрипт, который выводит список всех приложений за последние 24 часа с их названием, длительностью и статусом (завершено / с ошибкой). Это базовый строительный блок для автоматического post-job-отчёта.


Итоги

Event Log - это хроника всей жизни Spark-приложения в виде JSON, которая пишется асинхронно через EventLoggingListener в LiveListenerBus. Асинхронность гарантирует, что запись логов не влияет на производительность основных вычислений.

Хранение logов на локальном диске недопустимо в production из-за эфемерной природы кластеров. Для cloud/Kubernetes окружений используется S3/MinIO, для on-premise - HDFS. Директория должна существовать до запуска приложения.

Сжатие (zstd/lz4) и достаточный буфер - обязательные настройки. Без сжатия крупное приложение пишет несколько гигабайт JSON, создавая лишнюю нагрузку на сеть и хранилище. Маленький буфер вызывает S3 throttling.

Spark History Server - отдельный долгоживущий сервис, периодически сканирующий директорию логов и предоставляющий тот же UI, что и живой Spark. Локальный кеш метаданных (LevelDB/RocksDB) критически важен для производительности - без него каждый клик пользователя означает повторный парсинг файла логов.

Без Retention-политики History Server деградирует при накоплении сотен тысяч файлов. Рекомендуется двухслойный подход: History Server Cleaner убирает метаданные из кеша за N-1 дней, S3/MinIO Lifecycle Rule физически удаляет файлы за N дней.

Для стриминговых джоб обязательна ротация логов (spark.eventLog.rotation.enabled): без неё один файл вырастает до десятков гигабайт за недели работы, делая его парсинг History Server'ом нереалистичным.


Краткий глоссарий терминов урока

LiveListenerBus - асинхронная шина событий внутри Spark-драйвера, доставляющая SparkListener-события (Job/Stage/Task/Executor) всем подписанным слушателям в отдельных потоках.

EventLoggingListener - подписчик LiveListenerBus, сериализующий события в JSON и записывающий их в event log файл при включённом spark.eventLog.enabled = true.

Event Log - файл (один на приложение) в формате «один JSON-объект на строку», растущий по ходу выполнения Spark-приложения и завершающийся записью SparkListenerApplicationEnd.

Spark History Server - отдельный HTTP-сервис, читающий event log файлы из хранилища и предоставляющий Spark UI для завершённых приложений на порту 18080.

spark.eventLog.dir - параметр, указывающий директорию хранения event logs. Поддерживает протоколы file://, hdfs://, s3a://.

spark.eventLog.compress - включает сжатие event log «на лету». Кодек задаётся через spark.eventLog.compression.codec (lz4 по умолчанию, zstd рекомендуется для новых кластеров).

spark.eventLog.buffer.kb - размер in-memory буфера EventLoggingListener перед записью в хранилище. Большее значение снижает число S3 PUT-запросов.

spark.eventLog.rotation.enabled - включает ротацию event log для долгоработающих приложений. Максимальный размер сегмента задаётся через spark.eventLog.rotation.maxFileSize.

S3A FileSystem - реализация Hadoop FileSystem API для доступа к S3-совместимым хранилищам (Amazon S3, MinIO, Ceph RGW). Используется через spark.hadoop.fs.s3a.* конфигурацию.

Sticky Bit (HDFS) - бит прав доступа, разрешающий удаление файла из директории только его владельцу. Устанавливается через hdfs dfs -chmod 1777 на shared-директориях event logs.

spark.history.fs.update.interval - периодичность сканирования History Server директории логов на предмет новых приложений.

spark.history.store.path - путь к локальному кешу метаданных History Server (LevelDB/RocksDB), хранящему распарсенные данные приложений для быстрого доступа без повторного чтения файлов логов.

History Server Cleaner - встроенный механизм очистки старых записей из кеша History Server. Настраивается через spark.history.fs.cleaner.*.

S3 Lifecycle Rule / MinIO Lifecycle Policy - серверное правило объектного хранилища для автоматического удаления (Expiration) или перехода в другой storage class (Transition) объектов по истечении заданного срока. Работает независимо от Spark и History Server.

retainedApplications - параметр History Server, ограничивающий число приложений в in-memory части кеша (менее используемые вытесняются на диск в store.path).

REST API History Server - HTTP API (/api/v1/applications/...) для программного доступа к данным о завершённых приложениях, stage'ах, задачах и метриках - без браузерного UI.

Multipart Upload (S3A) - механизм загрузки крупных файлов частями параллельно. Для event logs настраивается через spark.hadoop.fs.s3a.multipart.size и spark.hadoop.fs.s3a.fast.upload.