Spark Metrics System: Sources, Sinks и namespace таксономия

Глубокое погружение в систему метрик Apache Spark: архитектура Dropwizard Metrics, таксономия namespace, все Sources (Driver, Executor, JVM, Shuffle), все Sinks (Prometheus, Graphite, JMX, Slf4j), конфигурация metrics.properties, интеграция с Kubernetes PodMonitor и практика: стек Spark + Prometheus + Grafana на Docker Compose.

platform

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

Именно для этого в Apache Spark существует Spark Metrics System — встроенная подсистема сбора и экспорта числовых метрик, построенная поверх библиотеки Dropwizard Metrics. Она работает параллельно с Spark UI, не конкурирует с ним, а дополняет его: пока Spark UI показывает детальную картину одного приложения в виде интерактивных таблиц, Metrics System экспортирует тысячи числовых показателей в системы мониторинга (Prometheus, Graphite, Grafana), где они хранятся годами и позволяют строить тренды, capacity planning-отчёты и алерты.


Зачем Spark своя система метрик

Ограничения Spark UI как инструмента мониторинга

Spark UI превосходно справляется со своей задачей — интерактивной диагностикой конкретного приложения. Но у него есть фундаментальные ограничения, которые делают его непригодным для production-мониторинга кластера:

Данные живут в памяти драйвера. Web UI на порту 4040 существует только пока работает SparkContext. Как только приложение завершается, все метрики исчезают. Нельзя построить graf JVM heap usage за последние 30 дней, анализируя только Spark UI.

Нет временных рядов. Spark UI показывает агрегированные значения за время выполнения stage или task — например, «shuffle write: 4.2 GB». Но он не хранит историю: как менялся heap memory каждые 15 секунд? Когда именно произошёл GC pause на 800 мс? Эти вопросы Spark UI не может ответить.

Один UI — одно приложение. На одном Yarn-кластере параллельно могут работать десятки Spark-приложений. У каждого — свой порт 4040 (или 4041, 4042...). Нет единой точки, где можно увидеть все приложения одновременно и сравнить их потребление ресурсов.

Нет алертинга. Spark UI — пассивный инструмент: он показывает данные, но не может послать уведомление, если heap usage executor'а превысил 90%.

Что даёт Spark Metrics System

Spark Metrics System решает все эти проблемы за счёт другой архитектурной концепции: она непрерывно экспортирует числовые метрики во внешние системы хранения (Prometheus, Graphite, InfluxDB) через стандартизированные протоколы. Это даёт:

  • Долгосрочное хранение метрик в формате временных рядов (дни, месяцы, годы)
  • Единую точку наблюдения за всеми приложениями и компонентами кластера
  • Алертинг на основе числовых пороговых значений
  • Дашборды с историческими трендами и сравнением нескольких джоб
  • Автоматический capacity planning на основе накопленных данных

Схема выше показывает ключевую архитектурную особенность: Metrics System живёт внутри каждого JVM-процесса — и в Driver, и в каждом Executor. Это не централизованный агент, а встроенный инструментарий, который экспортирует метрики своего процесса наружу.


Архитектура: Dropwizard Metrics под капотом

Что такое Dropwizard Metrics

Apache Spark не изобретал систему метрик с нуля. Начиная с версии 1.1, Spark использует библиотеку Dropwizard Metrics (раньше называлась Codahale Metrics) — одну из самых зрелых и широко используемых библиотек инструментирования для JVM.

Dropwizard Metrics предоставляет:

  • MetricRegistry — центральный реестр всех метрик процесса. Каждая метрика имеет уникальное строковое имя в виде пути (например, spark.driver.jvm.heap.used).
  • Типы метрик: Gauge (мгновенное значение), Counter (накапливаемый счётчик), Meter (события в секунду), Timer (время + процентили), Histogram (распределение значений).
  • Sink API — абстрактный интерфейс для экспорта значений из реестра во внешние системы.

Spark дополняет эту базу двумя собственными концепциями:

  • Source — объект, который знает, как собрать метрики из конкретного компонента Spark, и регистрирует их в MetricRegistry.
  • Sink — объект, который знает, как отправить собранные метрики во внешнюю систему (Prometheus, Graphite и т.д.).

Жизненный цикл метрики

Последовательность показывает важный момент: метрики регистрируются один раз при старте и дальше живут в реестре. Sink не «пушит» данные постоянно — он либо отвечает на pull-запросы (Prometheus), либо сам периодически пушит их (Graphite). Регистрация метрики в Dropwizard — это регистрация функции-геттера (Gauge), которая вызывается в момент экспорта. Нет постоянного копирования данных.

Driver vs Executor: два независимых мира

Одна из ключевых вещей, которую нужно понять сразу: Driver и каждый Executor — это отдельные JVM-процессы. У каждого из них своя копия MetricRegistry. Spark Metrics System работает независимо в каждом из этих процессов.

Это означает:

  • Если у тебя 10 executor'ов, у тебя 11 независимых MetricRegistry (10 + driver)
  • Каждый экспортирует свои метрики отдельно
  • Prometheus должен знать адреса всех этих процессов и скрейпить каждый
  • Имена метрик включают идентификатор executor'а, чтобы различать их в Prometheus

Разделение Driver и Executor определяет, какие метрики где доступны. На стороне Driver есть информация о планировщике (DAGScheduler, количество активных job/stage), о cluster manager (статус executor'ов), о SQL планировщике. На стороне каждого Executor — JVM heap, GC, потребление CPU, shuffle read/write, BlockManager.


Sources: откуда берутся метрики

Source — это класс, который знает, как получить значения метрик из конкретного компонента Spark, и регистрирует их в MetricRegistry. Каждый Source реализует интерфейс org.apache.spark.metrics.source.Source.

Полная карта Sources по компонентам

JvmSource: метрики Java Virtual Machine

JvmSource — наиболее универсальный Source, который работает одинаково в Driver и Executor. Он использует стандартный MXBean API JVM и регистрирует метрики следующих категорий:

Heap memory — самые важные метрики для диагностики OOM и настройки кластера:

Метрика Смысл Что означает в практике
heap.init Начальный размер heap при старте JVM Определяется -Xms
heap.used Текущее использование heap Растёт по мере создания объектов
heap.committed Память, зарезервированная у ОС Всегда >= used
heap.max Максимальный размер heap Определяется -Xmx

Non-Heap memory — память JVM вне heap (метаспейс, стеки потоков, JIT-компилятор):

Метрика Смысл
nonHeap.used Используемая не-heap память
nonHeap.committed Зарезервированная не-heap память

Рост nonHeap.used при загрузке большого количества UDF-классов — нормальное явление. Ненормальный рост может указывать на утечку ClassLoader в Spark Streaming.

Garbage Collection — критически важны для анализа производительности:

Метрика Смысл
gc.G1-Young-Generation.count Количество minor GC за время жизни JVM
gc.G1-Young-Generation.time Суммарное время minor GC в мс
gc.G1-Old-Generation.count Количество full GC
gc.G1-Old-Generation.time Суммарное время full GC в мс

GC time — один из главных индикаторов здоровья Spark-кластера. Если gc.G1-Old-Generation.time растёт быстро, executor тратит значительное время на сборку мусора вместо полезной работы. При значении GC time > 10% от общего времени работы executor'а следует рассматривать увеличение heap или настройку GC-параметров.

Threads — информация о потоках JVM:

Метрика Смысл
threads.count Общее число живых потоков
threads.daemon.count Число daemon-потоков
threads.deadlock.count Число потоков в deadlock (должно быть 0!)
threads.blocked.count Число потоков в состоянии BLOCKED

Рост threads.count свыше нескольких сотен и появление threads.deadlock.count > 0 — признаки критической проблемы, требующей немедленного внимания.

ExecutorSource: метрики выполнения задач

ExecutorSource живёт только в процессах executor'ов (не в driver) и собирает метрики о выполнении задач и операциях ввода-вывода.

Thread pool — о параллелизме на executor:

Метрика Смысл
threadpool.activeTasks Задачи, выполняющиеся прямо сейчас
threadpool.completeTasks Суммарно завершённых задач
threadpool.currentPool_size Размер thread pool
threadpool.maxPool_size Максимальный размер pool (= spark.executor.cores)

Если activeTasks постоянно равно maxPool_size, executor загружен на 100%. Если значительно меньше — часть ядер простаивает (или задачи ждут данных).

Filesystem I/O — чтение/запись в разные файловые системы:

Метрика Смысл
filesystem.hdfs.read_bytes Байт прочитано из HDFS
filesystem.hdfs.write_bytes Байт записано в HDFS
filesystem.s3a.read_bytes Байт прочитано из S3
filesystem.s3a.write_bytes Байт записано в S3
filesystem.file.read_bytes Байт прочитано с локального диска
filesystem.hdfs.read_ops Операций чтения HDFS
filesystem.hdfs.largeRead_ops «Больших» операций чтения (>65536 байт)
filesystem.hdfs.write_ops Операций записи HDFS

Эти метрики позволяют понять, где именно находится узкое место ввода-вывода. Если s3a.read_bytes растёт быстро, но hdfs.read_bytes почти нулевой — приложение читает данные из S3. Сравнение read_bytes с пропускной способностью сети помогает понять, достигается ли теоретический предел.

Shuffle метрики — особенно важны для джоб с тяжёлыми join/groupBy:

Метрика Смысл
shuffle.localBlocksFetched Блоки, прочитанные локально (быстро)
shuffle.remoteBlocksFetched Блоки, полученные по сети (медленно)
shuffle.remoteBytesRead Байт прочитано по сети при shuffle
shuffle.remoteBytesReadToDisk Байт, вылившихся на диск при shuffle
shuffle.localBytesRead Байт прочитано локально
shuffle.writeTime Время записи shuffle данных в нс
shuffle.shuffleBytesWritten Байт записано как shuffle output
shuffle.shuffleRecordsWritten Записей в shuffle output

Если remoteBytesReadToDisk постоянно растёт — executor не хватает памяти для shuffle буферов, и данные спиливаются на диск. Это сигнал увеличить spark.executor.memory или spark.shuffle.memoryFraction.

BlockManagerSource: метрики хранилища RDD/DataFrame кеша

BlockManagerSource присутствует и в Driver, и в каждом Executor. Он отслеживает состояние BlockManager — компонента, отвечающего за кеширование RDD/DataFrame блоков в памяти и на диске.

Метрика Смысл
memory.maxMem_MB Максимум памяти под BlockManager
memory.remainingMem_MB Свободная память BlockManager
memory.memUsed_MB Используемая память BlockManager
memory.diskSpaceUsed_MB Занято на диске (spill)
memory.offHeapMaxMem_MB Off-heap максимум (если включён)
memory.offHeapMemUsed_MB Off-heap использование

Важно понимать, что BlockManager использует ту часть executor memory, которая определяется spark.memory.fraction (по умолчанию 0.6 от spark.executor.memory). Если remainingMem_MB стабильно близко к нулю — кеш вытесняет данные, и RDD перевычисляются заново. Это объясняет, почему джоба, активно использующая .cache(), работает медленнее ожидаемого.

DAGSchedulerSource: метрики планировщика (только Driver)

DAGSchedulerSource регистрирует метрики о состоянии задач планировщика DAG. Доступен только на Driver:

Метрика Смысл
stage.waitingStages Stage'ы в очереди (ждут завершения зависимостей)
stage.runningStages Stage'ы, выполняющиеся прямо сейчас
stage.failedStages Stage'ы, завершившиеся с ошибкой
job.allJobs Суммарное количество всех job
job.activeJobs Job, выполняющиеся прямо сейчас

Если stage.waitingStages постоянно велико при малом stage.runningStages — приложение страдает от нехватки executor'ов или от высокой задержки в очереди YARN/Kubernetes.


Namespace таксономия: как именуются метрики

Структура полного имени метрики

Каждая метрика в Spark Metrics System имеет уникальное имя, которое состоит из четырёх частей, разделённых точкой:

[namespace].[instance].[source].[metricName]
  • namespace — идентификатор приложения или кастомное имя (задаётся через spark.metrics.namespace)
  • instance — тип компонента: driver, executor, master, worker, shuffleService, applicationMaster
  • source — имя Source, который зарегистрировал метрику: jvm, executor, BlockManager, DAGScheduler и т.д.
  • metricName — конкретная метрика внутри Source (может быть многосегментной: heap.used, gc.G1-Young-Generation.time)

Пример реальных имён метрик, которые видит Prometheus:

application_1718000000000_0001.driver.jvm.heap.used
application_1718000000000_0001.driver.BlockManager.memory.remainingMem_MB
application_1718000000000_0001.driver.DAGScheduler.stage.runningStages
application_1718000000000_0001.0.jvm.heap.used
application_1718000000000_0001.1.jvm.gc.G1-Old-Generation.time

Проблема динамического App ID

Самая болезненная проблема Spark Metrics System в production — динамически генерируемый Application ID. Когда YARN или Kubernetes создаёт новое приложение, оно получает уникальный App ID вида application_1718000000000_0001. Следующий запуск той же джобы получит application_1718000000000_0002.

Это означает, что все имена метрик меняются при каждом перезапуске джобы. Для системы мониторинга, которая хранит данные по именам метрик, это катастрофа:

  • Grafana-дашборды нельзя сделать с фиксированными именами метрик — каждый раз будет новый префикс
  • Prometheus не может накапливать историю по одной «логической» джобе — каждый запуск создаёт новую группу метрик
  • Alerting-правила невозможно написать на конкретный App ID

Решение: spark.metrics.namespace

Параметр spark.metrics.namespace позволяет заменить динамический App ID на фиксированное строковое значение. Это стандартное решение проблемы для production:

conf = SparkConf()
conf.set("spark.metrics.namespace", "my_etl_pipeline")

Или в spark-defaults.conf:

spark.metrics.namespace  my_etl_pipeline

После этого все метрики получат предсказуемый префикс:

my_etl_pipeline.driver.jvm.heap.used
my_etl_pipeline.driver.DAGScheduler.stage.runningStages
my_etl_pipeline.0.executor.shuffle.remoteBytesRead

Теперь Grafana-дашборд, написанный для my_etl_pipeline.*, будет работать для всех запусков этой джобы — и показывать единый временной ряд.

Конфликты при параллельных запусках

Если одновременно запущены несколько экземпляров джобы с одинаковым spark.metrics.namespace, их метрики в Prometheus будут неразличимы. В Prometheus они будут суммироваться/конфликтовать. Для параллельных запусков нужно добавить уникальный суффикс (имя tenant'а, дата, instance ID) или использовать Prometheus labels через конфигурацию PodMonitor/ServiceMonitor.

Специальный синтаксис с переменной имени приложения

Начиная со Spark 3.2, в значении spark.metrics.namespace можно использовать переменную ${spark.app.name}:

spark.metrics.namespace  ${spark.app.name}

Это дёшево решает проблему: имя приложения обычно стабильно между запусками (daily-etl, ml-training-pipeline), в отличие от App ID. Пробелы в именах заменяются на _, специальные символы не допускаются.


Sinks: куда отправляются метрики

Sink — это компонент, который читает значения из MetricRegistry и доставляет их во внешнюю систему. Все Sinks конфигурируются через файл metrics.properties.

Файл metrics.properties: структура и расположение

Файл $SPARK_HOME/conf/metrics.properties — главный конфигурационный файл Spark Metrics System. По умолчанию его нет, и система не экспортирует ничего (метрики собираются во внутренний реестр, но никуда не отправляются).

Пример полного файла с несколькими Sinks:

# Файл: $SPARK_HOME/conf/metrics.properties
# Применяется ко всем компонентам: driver, executor, master, worker

# Какие Sources включить для каждого компонента
*.source.jvm.class=org.apache.spark.metrics.source.JvmSource

# PrometheusServlet: pull-based HTTP endpoint
*.sink.prometheusServlet.class=org.apache.spark.metrics.sink.PrometheusServlet
*.sink.prometheusServlet.path=/metrics/prometheus

# Graphite: push-based TCP sink (закомментировано)
#*.sink.graphite.class=org.apache.spark.metrics.sink.GraphiteSink
#*.sink.graphite.host=graphite.internal
#*.sink.graphite.port=2003
#*.sink.graphite.period=10
#*.sink.graphite.unit=seconds
#*.sink.graphite.prefix=spark

# Slf4j: вывод в логи (для отладки)
#*.sink.slf4j.class=org.apache.spark.metrics.sink.Slf4jSink
#*.sink.slf4j.period=60
#*.sink.slf4j.unit=seconds

# JMX: для jconsole/jvisualvm
#*.sink.jmx.class=org.apache.spark.metrics.sink.JmxSink

Синтаксис ключей:

  • * — glob, применяется к всем компонентам (driver, executor, master, worker)
  • driver — только для процесса Driver
  • executor — только для процессов Executor
  • master — только для Spark Standalone Master
  • worker — только для Spark Standalone Worker

Например, если нужно включить Prometheus только для executor'ов:

executor.sink.prometheusServlet.class=org.apache.spark.metrics.sink.PrometheusServlet
executor.sink.prometheusServlet.path=/metrics/prometheus

PrometheusServlet: pull-модель для современного стека

PrometheusServlet — самый важный Sink в современных Spark-деплойментах на Kubernetes и YARN. Он реализует pull-модель: Prometheus сам приходит к каждому Spark-процессу и запрашивает метрики по HTTP.

Формат вывода соответствует стандарту OpenMetrics (расширенный Prometheus exposition format). Пример ответа на /metrics/prometheus:

# TYPE metrics_executor_cpuTime_total untyped
metrics_executor_cpuTime_total{app_id="application_1718000000000_0001",executor_id="1"} 1.234567891E9

# TYPE metrics_jvm_heap_used untyped
metrics_jvm_heap_used{app_id="application_1718000000000_0001",executor_id="1"} 2147483648.0

# TYPE metrics_jvm_gc_G1_Young_Generation_time_total untyped
metrics_jvm_gc_G1_Young_Generation_time_total{app_id="application_1718000000000_0001",executor_id="1"} 4521.0

PrometheusServlet добавляет метки (labels) app_id и executor_id автоматически, начиная со Spark 3.0. Это стандартный способ различать метрики разных executor'ов в Prometheus.

Адрес endpoint'а — не порт 4040. Порт 4040 — это Spark UI Driver. Executor'ы запускают HTTP-сервер на случайном порту, определяемом параметрами:

spark.executor.port                0  (случайный)
spark.ui.port                      4040 (только driver)
spark.executor.ui.port             0  (случайный)

Это создаёт проблему service discovery для Prometheus — он должен как-то узнать порты всех Executor-процессов. Именно для этого нужна интеграция с Kubernetes PodMonitor (рассмотрим ниже).

GraphiteSink: push-модель для Legacy-стеков

GraphiteSink реализует push-модель: Spark сам периодически отправляет метрики на Graphite-сервер по TCP. Это старый подход, который был стандартным до появления Prometheus. Используется в стеках, где Graphite + Grafana уже развёрнуты.

*.sink.graphite.class=org.apache.spark.metrics.sink.GraphiteSink
*.sink.graphite.host=graphite.prod.internal
*.sink.graphite.port=2003
*.sink.graphite.period=10
*.sink.graphite.unit=seconds
*.sink.graphite.prefix=spark.prod

Параметры:

Параметр Значение Смысл
host hostname Адрес Graphite Carbon relay
port 2003 TCP порт Carbon plaintext protocol
period 10 Интервал пуша в секундах
unit seconds Единица измерения интервала
prefix spark.prod Префикс, добавляемый ко всем именам метрик

В отличие от Prometheus, где Spark является сервером (отвечает на запросы), в Graphite Spark является клиентом (сам инициирует соединение). Это проще с точки зрения network policies — нужно только одностороннее соединение от Spark-процесса к Graphite.

ConsoleSink и Slf4jSink: для отладки

ConsoleSink и Slf4jSink не предназначены для production — они нужны для отладки конфигурации Metrics System.

ConsoleSink выводит метрики в stdout:

*.sink.console.class=org.apache.spark.metrics.sink.ConsoleSink
*.sink.console.period=30
*.sink.console.unit=seconds

Slf4jSink пишет в логи через SLF4J:

*.sink.slf4j.class=org.apache.spark.metrics.sink.Slf4jSink
*.sink.slf4j.period=60
*.sink.slf4j.unit=seconds

С Slf4jSink в логах executor'а каждые 60 секунд появляется блок вида:

INFO MetricsSystem: type=GAUGE, name=spark.1.jvm.heap.used, value=2147483648.0
INFO MetricsSystem: type=GAUGE, name=spark.1.jvm.gc.G1-Young-Generation.count, value=142.0

Это удобно для быстрой проверки: «работают ли метрики вообще», не поднимая Prometheus.

JmxSink: интеграция с JConsole и Java Mission Control

JmxSink экспортирует все метрики через JMX (Java Management Extensions). Это позволяет подключаться к процессу Spark с помощью стандартных JVM-инструментов: JConsole, JVisualVM, Java Mission Control.

*.sink.jmx.class=org.apache.spark.metrics.sink.JmxSink

После включения все метрики Spark Metrics System появляются как MBean'ы в пространстве имён metrics. Это полезно при профилировании на локальной машине — можно открыть JConsole, подключиться к Spark-процессу и видеть все метрики в реальном времени без какой-либо дополнительной инфраструктуры.


Практика: Prometheus Sink на Kubernetes

Проблема service discovery для Executor'ов

В Kubernetes каждый Executor работает в отдельном Pod'е. Pod получает динамический IP-адрес, а порт HTTP-сервера с метриками (PrometheusServlet) — случайный номер, определяемый при запуске. Prometheus не может знать заранее, какие IP:PORT нужно скрейпить.

Решение — использовать механизм Kubernetes Service Discovery в Prometheus в связке с Prometheus Operator (CRD PodMonitor / ServiceMonitor).

Конфигурация metrics.properties для Kubernetes

Создадим ConfigMap с metrics.properties:

apiVersion: v1
kind: ConfigMap
metadata:
  name: spark-metrics-config
  namespace: spark
data:
  metrics.properties: |
    # Включить JVM метрики для всех компонентов
    *.source.jvm.class=org.apache.spark.metrics.source.JvmSource

    # PrometheusServlet - единый endpoint для pull
    *.sink.prometheusServlet.class=org.apache.spark.metrics.sink.PrometheusServlet
    *.sink.prometheusServlet.path=/metrics/prometheus

Монтируем ConfigMap в Spark pods через SparkApplication (если используем Spark Operator):

apiVersion: sparkoperator.k8s.io/v1beta2
kind: SparkApplication
metadata:
  name: my-etl-job
  namespace: spark
spec:
  sparkConf:
    "spark.metrics.namespace": "my_etl_job"
    "spark.ui.prometheus.enabled": "true"

  driver:
    configMaps:
      - name: spark-metrics-config
        path: /opt/spark/conf

  executor:
    configMaps:
      - name: spark-metrics-config
        path: /opt/spark/conf

Аннотации для Prometheus scrape

Самый простой способ настроить scrape без Prometheus Operator — использовать Pod annotations. Добавим аннотации в SparkApplication:

spec:
  driver:
    annotations:
      prometheus.io/scrape: "true"
      prometheus.io/port:   "4040"
      prometheus.io/path:   "/metrics/prometheus"

  executor:
    annotations:
      prometheus.io/scrape: "true"
      prometheus.io/path:   "/metrics/prometheus"

PodMonitor для Prometheus Operator

Если в кластере развёрнут Prometheus Operator (kube-prometheus-stack), используем PodMonitor CRD — он более гибок, чем аннотации:

apiVersion: monitoring.coreos.com/v1
kind: PodMonitor
metadata:
  name: spark-pods
  namespace: monitoring
  labels:
    release: kube-prometheus-stack
spec:
  namespaceSelector:
    matchNames:
      - spark

  selector:
    matchExpressions:
      - key: spark-role
        operator: In
        values: [driver, executor]

  podMetricsEndpoints:
    - port: spark-ui
      path: /metrics/prometheus
      interval: 15s
      relabelings:
        - sourceLabels: [__meta_kubernetes_pod_label_spark_role]
          targetLabel: spark_role
        - sourceLabels: [__meta_kubernetes_pod_name]
          targetLabel: pod

    - port: spark-exec-ui
      path: /metrics/prometheus
      interval: 15s
      relabelings:
        - sourceLabels: [__meta_kubernetes_pod_label_spark_role]
          targetLabel: spark_role
        - sourceLabels: [__meta_kubernetes_pod_label_spark_exec_id]
          targetLabel: executor_id

Чтобы порты стали «именованными», нужно добавить их в SparkApplication:

spec:
  driver:
    ports:
      - name: spark-ui
        containerPort: 4040
        protocol: TCP

  executor:
    ports:
      - name: spark-exec-ui
        containerPort: 4041
        protocol: TCP

Фиксированный порт executor'а

Spark 3.4+ поддерживает spark.executor.ui.port — фиксированный порт для HTTP-сервера executor'а. Установив его в константу (например, 4041), можно избежать проблемы с динамическими портами:

spark.executor.ui.port=4041

Тогда PodMonitor может ссылаться на конкретный порт без магии service discovery.

Конфигурация на YARN

На YARN (без Kubernetes) ситуация с service discovery другая. YARN не предоставляет механизма, аналогичного Kubernetes PodMonitor. Решение — использовать push-based подход с GraphiteSink, или настроить Prometheus с file-based service discovery.

Вариант 1: GraphiteSink (рекомендуется для YARN)

# metrics.properties для YARN-деплоймента
*.source.jvm.class=org.apache.spark.metrics.source.JvmSource

*.sink.graphite.class=org.apache.spark.metrics.sink.GraphiteSink
*.sink.graphite.host=graphite.yarn-cluster.internal
*.sink.graphite.port=2003
*.sink.graphite.period=15
*.sink.graphite.unit=seconds
*.sink.graphite.prefix=spark

Вариант 2: Prometheus с file_sd_config

Написать скрипт, который периодически запрашивает YARN Resource Manager REST API (/ws/v1/cluster/apps), находит запущенные Spark-приложения и записывает их адреса в JSON-файл для Prometheus file_sd:

# prometheus.yml
scrape_configs:
  - job_name: spark_on_yarn
    file_sd_configs:
      - files:
          - /etc/prometheus/spark_targets/*.json
        refresh_interval: 30s
[
  {
    "targets": ["node1.yarn:4040", "node2.yarn:4041"],
    "labels": {
      "job": "spark_etl",
      "app_id": "application_1718000000000_0001"
    }
  }
]

Мониторинг здоровья Lakehouse: ключевые метрики и алерты

Какие метрики важны в production

Из тысяч метрик, которые экспортирует Spark, в реальной production-эксплуатации Lakehouse (Iceberg + HDFS/S3 + Spark) нас интересует несколько десятков. Остановимся на самых важных и на том, что они означают на практике.

Alerting-правила в формате Prometheus

Запишем ключевые алерты в формате Prometheus Alerting Rules:

# spark_alerts.yaml - Prometheus Alerting Rules
groups:
  - name: spark_executor_health
    interval: 1m
    rules:

      # JVM Heap: критически высокое заполнение
      - alert: SparkExecutorHighHeapUsage
        expr: |
          (
            metrics_jvm_heap_used{spark_role="executor"}
            /
            metrics_jvm_heap_max{spark_role="executor"}
          ) > 0.85
        for: 5m
        labels:
          severity: warning
          team: data-platform
        annotations:
          summary: "Spark executor heap usage > 85%"
          description: |
            Executor {{ $labels.executor_id }} приложения {{ $labels.app_id }}
            использует {{ $value | humanizePercentage }} heap памяти.
            Риск OOM. Рассмотрите увеличение spark.executor.memory.

      # JVM Heap: почти OOM
      - alert: SparkExecutorCriticalHeapUsage
        expr: |
          (
            metrics_jvm_heap_used{spark_role="executor"}
            /
            metrics_jvm_heap_max{spark_role="executor"}
          ) > 0.95
        for: 2m
        labels:
          severity: critical
          team: data-platform
        annotations:
          summary: "Spark executor heap usage > 95% - imminent OOM"
          description: |
            Executor {{ $labels.executor_id }} приложения {{ $labels.app_id }}
            использует {{ $value | humanizePercentage }} heap.
            OOM неизбежен в течение минут.

      # GC Time: слишком много времени на сборку мусора
      - alert: SparkExecutorHighGCTime
        expr: |
          rate(
            metrics_jvm_gc_G1_Old_Generation_time_total{spark_role="executor"}[5m]
          ) > 100
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: "Spark executor spending too much time in Full GC"
          description: |
            Executor {{ $labels.executor_id }} выполняет Full GC
            со скоростью {{ $value }}ms/s.
            Если значение больше 100ms/s (> 10% wall clock),
            executor тратит значительное время на GC вместо задач.

      # Shuffle Spill: данные льются на диск
      - alert: SparkShuffleSpillToDisk
        expr: |
          rate(
            metrics_executor_shuffle_remoteBytesReadToDisk_total[5m]
          ) > 0
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Spark shuffle spilling to disk"
          description: |
            Executor {{ $labels.executor_id }} сливает shuffle данные на диск.
            Скорость: {{ $value | humanize }}B/s.
            Увеличьте spark.executor.memory или spark.shuffle.memoryFraction.

      # Failed Stages
      - alert: SparkStagesFailing
        expr: |
          increase(
            metrics_DAGScheduler_stage_failedStages_total{spark_role="driver"}[5m]
          ) > 0
        labels:
          severity: critical
        annotations:
          summary: "Spark job has failing stages"
          description: |
            Приложение {{ $labels.app_id }} имеет падающие stage'и.
            Откройте Spark UI / History Server для диагностики.

PromQL-запросы для Grafana дашборда

Несколько готовых PromQL-запросов для типичных панелей Grafana:

Процент использования Heap по всем executor'ам:

(
  metrics_jvm_heap_used{spark_role="executor", app_id=~"$app_id"}
  /
  metrics_jvm_heap_max{spark_role="executor", app_id=~"$app_id"}
) * 100

GC time rate для executor — сколько мс/с тратится на GC:

rate(metrics_jvm_gc_G1_Old_Generation_time_total{
  spark_role="executor",
  app_id=~"$app_id"
}[2m])

Суммарный shuffle network traffic (read) по всем executor'ам:

sum by (app_id) (
  rate(
    metrics_executor_shuffle_remoteBytesRead_total{
      spark_role="executor",
      app_id=~"$app_id"
    }[2m]
  )
)

Активных задач — для проверки параллелизма:

sum by (app_id) (metrics_executor_threadpool_activeTasks{app_id=~"$app_id"})

BlockManager: сколько памяти осталось для кеширования:

metrics_BlockManager_memory_remainingMem_MB{
  spark_role="executor",
  app_id=~"$app_id"
}

Лабораторная работа: Spark + Prometheus + Grafana на Docker Compose

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

Структура проекта

spark-monitoring-lab/
├── docker-compose.yml
├── conf/
│   └── metrics.properties
├── prometheus/
│   ├── prometheus.yml
│   └── rules/
│       └── spark_alerts.yml
├── grafana/
│   ├── provisioning/
│   │   ├── datasources/
│   │   │   └── prometheus.yml
│   │   └── dashboards/
│   │       └── dashboard.yml
│   └── dashboards/
│       └── spark_overview.json
└── jobs/
    └── sample_job.py

docker-compose.yml

version: '3.8'

services:
  # Apache Spark Master
  spark-master:
    image: bitnami/spark:3.5.1
    environment:
      - SPARK_MODE=master
      - SPARK_MASTER_WEBUI_PORT=8080
    ports:
      - "8080:8080"
      - "7077:7077"
    volumes:
      - ./conf/metrics.properties:/opt/bitnami/spark/conf/metrics.properties:ro
      - ./jobs:/jobs:ro

  # Apache Spark Worker
  spark-worker:
    image: bitnami/spark:3.5.1
    environment:
      - SPARK_MODE=worker
      - SPARK_MASTER_URL=spark://spark-master:7077
      - SPARK_WORKER_WEBUI_PORT=8081
      - SPARK_WORKER_MEMORY=2G
      - SPARK_WORKER_CORES=2
    ports:
      - "8081:8081"
    volumes:
      - ./conf/metrics.properties:/opt/bitnami/spark/conf/metrics.properties:ro
    depends_on:
      - spark-master

  # Prometheus
  prometheus:
    image: prom/prometheus:v2.51.0
    ports:
      - "9090:9090"
    volumes:
      - ./prometheus/prometheus.yml:/etc/prometheus/prometheus.yml:ro
      - ./prometheus/rules:/etc/prometheus/rules:ro
      - prometheus_data:/prometheus
    command:
      - "--config.file=/etc/prometheus/prometheus.yml"
      - "--storage.tsdb.retention.time=30d"
      - "--web.enable-lifecycle"

  # Grafana
  grafana:
    image: grafana/grafana:10.4.1
    ports:
      - "3000:3000"
    environment:
      - GF_SECURITY_ADMIN_PASSWORD=admin
      - GF_USERS_ALLOW_SIGN_UP=false
    volumes:
      - ./grafana/provisioning:/etc/grafana/provisioning:ro
      - ./grafana/dashboards:/var/lib/grafana/dashboards:ro
      - grafana_data:/var/lib/grafana
    depends_on:
      - prometheus

volumes:
  prometheus_data:
  grafana_data:

conf/metrics.properties

# JVM метрики для всех компонентов
*.source.jvm.class=org.apache.spark.metrics.source.JvmSource

# PrometheusServlet - pull endpoint
*.sink.prometheusServlet.class=org.apache.spark.metrics.sink.PrometheusServlet
*.sink.prometheusServlet.path=/metrics/prometheus

# Slf4j для отладки (каждые 30 секунд)
*.sink.slf4j.class=org.apache.spark.metrics.sink.Slf4jSink
*.sink.slf4j.period=30
*.sink.slf4j.unit=seconds

prometheus/prometheus.yml

global:
  scrape_interval: 15s
  evaluation_interval: 15s

rule_files:
  - "rules/*.yml"

scrape_configs:
  # Spark Driver (порт 4040)
  - job_name: spark_driver
    static_configs:
      - targets:
          - spark-master:4040
        labels:
          spark_role: driver
          cluster: local-lab
    metrics_path: /metrics/prometheus
    scrape_interval: 15s

  # Spark Master (порт 8080)
  - job_name: spark_master
    static_configs:
      - targets:
          - spark-master:8080
        labels:
          spark_role: master
    metrics_path: /metrics/prometheus

  # Spark Worker (порт 8081)
  - job_name: spark_worker
    static_configs:
      - targets:
          - spark-worker:8081
        labels:
          spark_role: worker
    metrics_path: /metrics/prometheus

Executor scrape в Standalone

В режиме Spark Standalone driver и worker запускаются в заранее известных процессах на известных хостах. Executor'ы запускаются worker-процессом на worker-хосте и используют случайные порты. Для лабораторной работы будем смотреть метрики через Driver (он агрегирует часть данных) и через Worker. В production на Kubernetes используйте PodMonitor как описано выше.

jobs/sample_job.py

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
import time

spark = (
    SparkSession.builder
    .appName("metrics-lab-job")
    .config("spark.metrics.namespace", "metrics_lab")
    .getOrCreate()
)

sc = spark.sparkContext
sc.setLogLevel("WARN")

print("=== Spark Metrics Lab Job Started ===")
print(f"Spark UI: {sc.uiWebUrl}")
print(f"Metrics: {sc.uiWebUrl}/metrics/prometheus")

# Генерируем нагрузку на JVM heap
print("\n[1] Создаём большой DataFrame и кешируем его...")
df = spark.range(0, 10_000_000).toDF("id")
df = df.withColumn("category", (F.col("id") % 100).cast("string"))
df = df.withColumn("value", F.rand() * 1000)
df.cache()
df.count()
print(f"    Закешировано строк: {df.count():,}")

# Shuffle нагрузка: groupBy + join
print("\n[2] Тяжёлый groupBy (shuffle нагрузка)...")
agg = df.groupBy("category").agg(
    F.count("id").alias("cnt"),
    F.avg("value").alias("avg_val"),
    F.sum("value").alias("sum_val"),
)
agg.cache()
print(f"    Категорий: {agg.count()}")

# Self-join для максимального shuffle
print("\n[3] Self-join для создания shuffle трафика...")
joined = df.alias("a").join(
    df.alias("b"),
    F.col("a.category") == F.col("b.category"),
    "inner"
).select(F.col("a.id"), F.col("b.value"))
print(f"    Строк после join: {joined.limit(100).count()}")

# Пауза для наблюдения метрик
print("\n[4] Спим 120 секунд - наблюдайте метрики в Grafana...")
print(f"    Prometheus: http://localhost:9090")
print(f"    Grafana:    http://localhost:3000")
time.sleep(120)

df.unpersist()
agg.unpersist()
print("\n=== Job Completed ===")
spark.stop()

Запуск лабораторной работы

# 1. Запустить стек мониторинга
docker compose up -d

# 2. Подождать, пока поднимутся все сервисы (~30 секунд)
docker compose ps

# 3. Запустить Spark job
docker compose exec spark-master \
  /opt/bitnami/spark/bin/spark-submit \
    --master spark://spark-master:7077 \
    --conf spark.metrics.namespace=metrics_lab \
    --conf spark.eventLog.enabled=false \
    /jobs/sample_job.py

# 4. Открыть в браузере
# Spark UI:    http://localhost:4040
# Prometheus:  http://localhost:9090
# Grafana:     http://localhost:3000  (admin/admin)

Проверка метрик в Prometheus

Открываем Prometheus UI по адресу http://localhost:9090 и переходим в Graph. Вводим запрос:

metrics_jvm_heap_used

Должны появиться метрики с labels instance, job, spark_role. Если ничего нет — проверяем:

# Убедиться, что endpoint доступен
curl http://localhost:4040/metrics/prometheus | head -30

# Проверить статус scrape в Prometheus
# Меню: Status -> Targets
# Все targets должны быть в состоянии UP

Если endpoint возвращает ошибку или страницу Spark UI вместо метрик — проверьте, что metrics.properties смонтирован корректно:

docker compose exec spark-master \
  cat /opt/bitnami/spark/conf/metrics.properties

Полезные PromQL запросы для исследования

После запуска job введите следующие запросы в Prometheus:

# Heap использование Driver в %
(metrics_jvm_heap_used{job="spark_driver"} /
 metrics_jvm_heap_max{job="spark_driver"}) * 100

# GC время - сколько ms/s тратится на сборку мусора
rate(metrics_jvm_gc_G1_Young_Generation_time_total[2m])

# Активные потоки
metrics_jvm_threads_count

# BlockManager - сколько памяти занято под кеш
metrics_BlockManager_memory_memUsed_MB

Расширенная конфигурация: тонкая настройка

Фильтрация метрик: включить только нужные

По умолчанию Spark Metrics System экспортирует все зарегистрированные метрики. Если нужно ограничить набор (например, снизить cardinality в Prometheus), можно использовать фильтрацию на уровне metric_relabel_configs в Prometheus:

scrape_configs:
  - job_name: spark_executors
    metric_relabel_configs:
      # Оставить только heap, GC и shuffle метрики
      - source_labels: [__name__]
        regex: "metrics_(jvm_heap|jvm_gc|executor_shuffle|BlockManager).*"
        action: keep

      # Удалить слишком детальные метрики memory pool
      - source_labels: [__name__]
        regex: "metrics_jvm_memory_pool.*"
        action: drop

      # Убрать executor_id из labels для агрегации (опционально)
      - action: labeldrop
        regex: executor_id

Настройка производительности: scrape interval

Очень частый scrape (менее 10 секунд) создаёт нагрузку на HTTP-сервер каждого Spark-процесса, которая при тысячах executor'ов может стать заметной. Рекомендуемые значения:

Сценарий Scrape interval Retention
Debugging / инцидент 5s 1-2 дня
Production мониторинг 15-30s 30 дней
Capacity planning 60s 1 год
Тревожные алерты 15s 15 дней

Для Alertmanager алерты вычисляются на интервале evaluation_interval (обычно 1 мин), поэтому scrape_interval должен быть <= evaluation_interval.

Кастомные метрики через Accumulator

Spark позволяет создавать собственные метрики через AccumulatorV2. Это полезно, когда нужно экспортировать бизнес-метрики (количество обработанных записей, ошибок, событий) вместе с системными метриками:

from pyspark.sql import SparkSession
from pyspark.sql import functions as F

spark = SparkSession.builder.appName("custom-metrics").getOrCreate()
sc = spark.sparkContext

# Простой счётчик через встроенный Accumulator
records_processed = sc.accumulator(0)
errors_count = sc.accumulator(0)

def process_row(row):
    try:
        result = int(row["value"]) * 2
        records_processed.add(1)
        return result
    except Exception:
        errors_count.add(1)
        return None

df = spark.range(1000).toDF("value")
rdd = df.rdd.map(lambda r: process_row(r))
rdd.count()

# Аккумуляторы доступны на driver после action
print(f"Обработано записей: {records_processed.value}")
print(f"Ошибок: {errors_count.value}")

Аккумуляторы и Metrics System

Аккумуляторы не экспортируются через PrometheusServlet автоматически. Значения аккумуляторов доступны только на Driver после завершения action. Для экспорта бизнес-метрик в Prometheus в реальном времени используют отдельный Prometheus Pushgateway или кастомный Source на Scala/Java.


Интеграция с Apache Iceberg: специфические метрики

При работе с Iceberg-таблицами важно отслеживать не только метрики Spark, но и метрики самого Iceberg-движка, которые доступны через JMX.

Активация метрик Iceberg

spark = (
    SparkSession.builder
    .config("spark.sql.extensions",
            "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions")
    .config("spark.sql.catalog.local",
            "org.apache.iceberg.spark.SparkCatalog")
    .config("spark.sql.catalog.local.type", "hadoop")
    .config("spark.sql.catalog.local.warehouse", "/tmp/iceberg-warehouse")
    .getOrCreate()
)

Ключевые метрики Iceberg через JMX

После включения JmxSink в metrics.properties в JConsole появятся:

Метрика Смысл
iceberg.table.scan.total-planning-duration Время планирования scan (file listing)
iceberg.table.scan.result-data-files Число файлов данных после прунинга
iceberg.table.scan.skipped-data-files Число файлов, пропущенных прунингом
iceberg.table.commit.duration Время коммита транзакции
iceberg.table.commit.added-data-files Файлов добавлено в коммите
iceberg.table.commit.removed-data-files Файлов удалено (expire/compact)

Отношение skipped к total files — ключевой показатель эффективности partition pruning и Z-order:

# Эффективность прунинга (чем выше - тем лучше)
iceberg_table_scan_skipped_data_files_total /
(iceberg_table_scan_result_data_files_total + iceberg_table_scan_skipped_data_files_total)

Если это значение меньше 0.3 (то есть менее 30% файлов пропускается), стоит пересмотреть стратегию партиционирования или добавить Z-order кластеризацию.


Диагностика проблем с Metrics System

Метрики не появляются в Prometheus

# 1. Проверить, что endpoint отвечает
curl -s http://spark-driver:4040/metrics/prometheus | grep "^metrics_" | head -5

# 2. Если пустой ответ - проверить логи Spark при старте
# Должны быть строки:
# INFO MetricsSystem: Registered source jvm
# INFO MetricsSystem: Registered sink prometheusServlet

grep "MetricsSystem" $SPARK_HOME/logs/spark-*.out | grep -E "(Registered|ERROR)"

# 3. Задать путь к конфигу явно через SparkConf
# (по умолчанию Spark ищет в $SPARK_HOME/conf/metrics.properties)
spark.conf.set("spark.metrics.conf", "/custom/path/metrics.properties")

Метрики есть, но Prometheus не scrapes

В Prometheus UI откройте Status -> Targets. Если target в состоянии DOWN:

  • connection refused — Firewall или network policy блокирует доступ к порту. Driver/Worker ещё не стартовали.
  • context deadline exceeded — Слишком долгий ответ. Увеличьте scrape_timeout:
scrape_configs:
  - job_name: spark_driver
    scrape_timeout: 30s
    ...

Взрывной рост cardinality в Prometheus

Частая проблема при большом количестве executor'ов — взрывной рост cardinality (количества уникальных label combinations) в Prometheus TSDB. Это замедляет Prometheus и увеличивает потребление памяти.

Если cardinality проблема, применяйте агрессивную фильтрацию через metric_relabel_configs и убирайте executor_id из labels там, где важна агрегация, а не индивидуальные показатели.


Итоги: когда какой инструмент использовать

Мы рассмотрели полную архитектуру Spark Metrics System. Подведём итог — когда использовать каждый инструмент мониторинга:

Spark Metrics System — это не замена Spark UI, а дополнение к нему на уровне кластера. Spark UI отвечает на вопросы «что происходит» внутри одного приложения. Metrics System + Prometheus + Grafana отвечают на вопросы «как ведёт себя система во времени» и «когда нужно поднять алерт».

Ключевые выводы этого урока:

  • Dropwizard Metrics — основа Spark Metrics System; понимание типов метрик (Gauge, Counter, Timer) важно для интерпретации экспортируемых данных
  • Каждый JVM-процесс (Driver, каждый Executor) имеет независимый MetricRegistry и отдельный endpoint
  • spark.metrics.namespace — обязательный параметр для production; без него метрики меняют имена при каждом перезапуске
  • PrometheusServlet — стандарт для Kubernetes; GraphiteSink — для YARN и legacy-стеков
  • Ключевые метрики для алертов: jvm.heap.used/max, gc.G1-Old-Generation.time, shuffle.remoteBytesReadToDisk, stage.failedStages
  • PodMonitor в Kubernetes решает проблему service discovery для executor'ов с динамическими портами
  • Iceberg scan pruning ratio — отдельный важный показатель здоровья Lakehouse, дополняющий метрики Spark