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.
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,applicationMastersource— имя 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— только для процесса Driverexecutor— только для процессов Executormaster— только для Spark Standalone Masterworker— только для 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