Observability: структурированное логирование и трассировка Spark job
Три столпа observability в data engineering: structured logging (JSON logs, python-json-logger), distributed tracing (correlation IDs, SparkContext local properties), Spark Listeners для автоматических метрик, интеграция с Prometheus и Grafana
Введение: проблема «чёрного ящика» в распределённых вычислениях¶
Представьте: в 2:00 ночи Airflow сигнализирует, что задача silver_orders упала. Вы открываете логи и видите несколько тысяч строк вида:
24/01/15 02:14:37 INFO DAGScheduler: Stage 12 (collect at OrdersJob.scala:87) finished in 14.8 s
24/01/15 02:14:38 WARN TaskSchedulerImpl: Initial job has not accepted any resources; check your cluster UI
24/01/15 02:14:52 ERROR SparkContext: Error initializing SparkContext.
java.lang.OutOfMemoryError: GC overhead limit exceeded
at org.apache.spark.sql.catalyst.expressions...
(47 more frames)
Где именно упало? На каком Executor? На каком этапе трансформации? Сколько строк успело обработать? Какой датасет был в обработке? Как это связано с задачей extract_raw_orders, которая запустилась раньше? Стандартные log4j-логи Spark не дают ответа ни на один из этих вопросов.
Observability - инженерная дисциплина, позволяющая понять внутреннее состояние сложной распределённой системы по её внешним сигналам. Не просто «смотреть логи», а иметь структурированный поток телеметрии, по которому можно ответить на вопрос: «что именно произошло, почему, и как долго это нарастало?»
В data engineering обычные DevOps-инструменты не работают напрямую. У Spark-пайплайна нет «одного процесса» - есть Driver, десятки Executor'ов, shuffle, объектное хранилище, внешние JDBC-источники. Ошибка в задаче Executor'а может быть следствием неправильного плана запроса на Driver'е. OOM на воркере - следствием data skew. Медленный shuffle - следствием неверного выбора числа партиций две стадии назад.
Часть 1. Три столпа Observability¶
1.1 Мониторинг vs Observability: разница принципиальная¶
Мониторинг работает с заранее известными вопросами. Вы знаете, что хотите отслеживать (CPU, memory, job duration), создаёте dashboard и алерт. Если что-то выходит за порог - алерт срабатывает.
Observability работает с неизвестными вопросами. Вы не знаете заранее, что сломается. Вместо этого создаёте систему, которая позволит ответить на любой вопрос о состоянии системы - включая тот, который вы ещё не придумали. Для этого нужно три вида телеметрии:
1.2 Специфика Spark как объекта наблюдения¶
Spark-пайплайн - исключительно сложный объект для мониторинга по нескольким причинам.
JVM + Python разрыв. PySpark работает в двух языковых средах: Driver запускает Python-процесс, который через Py4J взаимодействует с JVM. Executor'ы запускают JVM-процессы, которые через Arrow или pickle сериализуют данные для Python UDF. Stack trace из JVM и из Python - это разные потоки, которые нужно вручную связывать.
Распределённость. Один Spark Job = Driver + N×Executor'ов. При падении Executor'а Driver получает уведомление, но оригинальный stack trace находится в логах конкретного Executor'а на конкретном K8s pod'е, который к этому моменту уже убит.
Отложенная материализация (Lazy evaluation). Ошибка в трансформации проявляется при Action, а не при определении. Логи при этом указывают на Action (.count(), .write()), а не на трансформацию с реальной проблемой.
Многоуровневость. Один запрос в Spark SQL разбивается на Job → Stage → Task. Task выполняется параллельно на всех Executor'ах. Метрика одной Task'и (время выполнения, shuffle bytes) не агрегируется автоматически в метрику всего Job'а.
1.3 Observability в контексте Medallion Architecture¶
Каждый слой Medallion требует своего набора сигналов:
| Слой | Что мониторить | Ключевые метрики |
|---|---|---|
| Bronze | Здоровье ingestion | rows ingested/sec, schema drift events, parse error rate, source lag |
| Silver | Качество трансформаций | dedup rate, null rate, quarantine rate, transformation duration |
| Gold | Стабильность бизнес-метрик | AOV delta, revenue delta, unique users delta, freshness lag |
| Инфраструктура | Производительность кластера | shuffle bytes, GC pressure, spill to disk, executor utilization |
Часть 2. Структурированное логирование¶
2.1 Почему текстовые логи не работают в production¶
Классический текстовый лог выглядит так:
2024-01-15 02:14:22 - Processing batch for date 2024-01-14, loaded 5234891 rows from Bronze
2024-01-15 02:15:44 - Deduplication complete: 5198234 unique rows remain
2024-01-15 02:17:01 - Writing to Silver...
Попробуйте ответить на вопросы: какой пайплайн? какой run_id? сколько длилась дедупликация? сколько дубликатов? Для текстовых логов нужен ручной парсинг. Grep по 50GB лог-файлов. Нет возможности построить dashboard.
JSON-лог той же операции выглядит иначе:
{
"timestamp": "2024-01-15T02:14:22.341Z",
"level": "INFO",
"pipeline": "bronze_to_silver_orders",
"run_id": "airflow_run_2024-01-15T00:00:00",
"stage": "bronze_read",
"layer": "bronze",
"table": "orders",
"processing_date": "2024-01-14",
"rows_read": 5234891,
"duration_ms": 82341,
"trace_id": "7f3a2b1c9d4e5f6a",
"host": "spark-driver-xyz"
}
Этот лог:
- Парсится автоматически в ELK/Grafana Loki/Datadog
- Фильтруется по любому полю без regex
- Агрегируется:
avg(duration_ms) GROUP BY pipeline, processing_date - Строит алерты:
rows_read < 1000000 → critical alert
2.2 Настройка log4j2 для JSON-вывода в Spark¶
Spark использует log4j2 для JVM-логов. По умолчанию формат - текст. Переключаем на JSON:
# log4j2.properties - добавить в spark.driver.extraJavaOptions
# или передать через --conf spark.driver.extraJavaOptions=-Dlog4j2.configurationFile=...
status = error
name = PropertiesConfig
# Appender: структурированный JSON в stdout
appender.console.type = Console
appender.console.name = JsonConsole
appender.console.layout.type = JsonTemplateLayout
appender.console.layout.eventTemplateUri = classpath:EcsLayout.json
# Корневой логгер
rootLogger.level = WARN
rootLogger.appenderRef.console.ref = JsonConsole
# Spark-специфичные логгеры
logger.spark.name = org.apache.spark
logger.spark.level = WARN
# Наш application-код логируем на INFO
logger.app.name = com.company.dataplatform
logger.app.level = INFO
Подключение через spark-submit:
spark-submit \
--conf "spark.driver.extraJavaOptions=-Dlog4j2.configurationFile=/opt/spark/conf/log4j2.properties" \
--conf "spark.executor.extraJavaOptions=-Dlog4j2.configurationFile=/opt/spark/conf/log4j2.properties" \
jobs/silver_orders.py
2.3 Структурированный логгер на Python¶
Для Python-части (Driver-код, Python UDF) используем python-json-logger:
pip install python-json-logger
import logging
import json
import os
import time
from datetime import datetime
from pythonjsonlogger import jsonlogger
class SparkPipelineLogger:
"""
Структурированный логгер для PySpark-пайплайна.
Каждая запись - валидный JSON с предустановленным контекстом.
Контекст (pipeline_name, run_id, trace_id) устанавливается один раз
при инициализации и автоматически добавляется ко всем записям.
"""
def __init__(
self,
pipeline_name: str,
run_id: str,
trace_id: str = None,
layer: str = None,
):
self.pipeline_name = pipeline_name
self.run_id = run_id
self.trace_id = trace_id or self._generate_trace_id()
self.layer = layer
# Настраиваем logger с JSON форматтером
self._logger = logging.getLogger(f"pipeline.{pipeline_name}")
self._logger.setLevel(logging.DEBUG)
if not self._logger.handlers:
handler = logging.StreamHandler()
formatter = jsonlogger.JsonFormatter(
fmt="%(asctime)s %(levelname)s %(name)s %(message)s",
datefmt="%Y-%m-%dT%H:%M:%S.%fZ",
)
handler.setFormatter(formatter)
self._logger.addHandler(handler)
def _generate_trace_id(self) -> str:
import uuid
return uuid.uuid4().hex[:16]
def _base_context(self) -> dict:
"""Контекст, который добавляется к каждой записи."""
return {
"pipeline": self.pipeline_name,
"run_id": self.run_id,
"trace_id": self.trace_id,
"layer": self.layer,
"host": os.environ.get("HOSTNAME", "unknown"),
"spark_app_id": os.environ.get("SPARK_APP_ID", "unknown"),
}
def info(self, message: str, **extra):
self._logger.info(message, extra={**self._base_context(), **extra})
def warning(self, message: str, **extra):
self._logger.warning(message, extra={**self._base_context(), **extra})
def error(self, message: str, **extra):
self._logger.error(message, extra={**self._base_context(), **extra})
def stage_start(self, stage: str, **extra) -> float:
"""Логируем начало стадии и возвращаем timestamp для расчёта длительности."""
self.info(
f"Stage started: {stage}",
event="stage_start",
stage=stage,
**extra,
)
return time.time()
def stage_end(self, stage: str, start_time: float, **extra):
"""Логируем конец стадии с длительностью в миллисекундах."""
duration_ms = int((time.time() - start_time) * 1000)
self.info(
f"Stage completed: {stage}",
event="stage_end",
stage=stage,
duration_ms=duration_ms,
**extra,
)
return duration_ms
def rows_processed(
self,
stage: str,
rows_in: int,
rows_out: int,
table: str = None,
**extra,
):
"""Специализированный лог для статистики обработки строк."""
self.info(
f"Rows processed in {stage}",
event="rows_processed",
stage=stage,
rows_in=rows_in,
rows_out=rows_out,
rows_dropped=rows_in - rows_out,
drop_ratio=round((rows_in - rows_out) / rows_in, 4) if rows_in > 0 else 0,
table=table,
**extra,
)
def quality_gate(
self,
gate_name: str,
passed: bool,
metric_value,
threshold,
**extra,
):
"""Лог результата Quality Gate - для мониторинга и алертинга."""
level = "info" if passed else "warning"
getattr(self, level)(
f"Quality gate {'PASSED' if passed else 'FAILED'}: {gate_name}",
event="quality_gate",
gate_name=gate_name,
passed=passed,
metric_value=metric_value,
threshold=threshold,
**extra,
)
2.4 Использование логгера в пайплайне¶
def run_bronze_to_silver(
spark: SparkSession,
batch_date: str,
run_id: str,
trace_id: str = None,
) -> None:
log = SparkPipelineLogger(
pipeline_name="bronze_to_silver_orders",
run_id=run_id,
trace_id=trace_id,
layer="silver",
)
# ── Шаг 1: чтение Bronze ──
t0 = log.stage_start("bronze_read", table="orders", date=batch_date)
bronze_df = spark.read.format("delta").load(
f"s3a://datalake/bronze/orders/date={batch_date}/"
)
bronze_count = bronze_df.count()
log.stage_end("bronze_read", t0, rows_read=bronze_count, table="orders")
log.rows_processed("bronze_read", rows_in=bronze_count, rows_out=bronze_count)
# ── Шаг 2: дедупликация ──
t1 = log.stage_start("deduplication")
deduped_df = bronze_df.dropDuplicates(["order_id"])
deduped_count = deduped_df.count()
dup_count = bronze_count - deduped_count
log.stage_end("deduplication", t1, rows_deduplicated=dup_count)
log.rows_processed(
"deduplication",
rows_in=bronze_count,
rows_out=deduped_count,
table="orders",
)
# ── Шаг 3: quality gate ──
null_count = deduped_df.filter(F.col("order_id").isNull()).count()
log.quality_gate(
gate_name="null_order_id_check",
passed=(null_count == 0),
metric_value=null_count,
threshold=0,
table="orders",
)
if null_count > 0:
raise ValueError(f"Null order_id detected: {null_count} rows")
# ── Шаг 4: запись Silver ──
t2 = log.stage_start("silver_write", table="orders")
(
deduped_df
.write.mode("overwrite")
.format("delta")
.partitionBy("processing_date")
.save("s3a://datalake/silver/orders/")
)
log.stage_end("silver_write", t2, rows_written=deduped_count)
log.info(
"Pipeline complete",
event="pipeline_complete",
bronze_rows=bronze_count,
silver_rows=deduped_count,
duplicates_removed=dup_count,
)
2.5 Анатомия идеального лог-события¶
{
"timestamp": "2024-01-15T02:15:44.891Z",
"level": "INFO",
"event": "rows_processed",
"message": "Rows processed in deduplication",
"pipeline": "bronze_to_silver_orders",
"run_id": "scheduled__2024-01-15T00:00:00+00:00",
"trace_id": "7f3a2b1c9d4e5f6a",
"layer": "silver",
"stage": "deduplication",
"table": "orders",
"rows_in": 5234891,
"rows_out": 5198234,
"rows_dropped": 36657,
"drop_ratio": 0.007,
"host": "spark-driver-orders-abc123",
"spark_app_id": "application_1705276800000_0042"
}
Каждое поле несёт смысловую нагрузку:
event- машиночитаемое имя типа события для группировки в Kibana/Grafanatrace_id- сквозной идентификатор, связывающий все логи одного запуска через Airflow, Spark Driver и Executor'ыrows_dropped+drop_ratio- числовые метрики, которые можно агрегировать в dashboardspark_app_id- связывает лог с Spark UI и History Server
Часть 3. Distributed Tracing и Correlation IDs¶
3.1 Проблема потери контекста¶
Рассмотрим цепочку событий в реальном pipeline:
Airflow DAG: daily_orders → запускает task: extract_orders
extract_orders: spark-submit jobs/extract.py → Application ID: app_001
→ Stage 0: read Kafka (OK)
→ Stage 1: write Bronze (OK)
Airflow DAG: daily_orders → запускает task: transform_silver
transform_silver: spark-submit jobs/transform.py → Application ID: app_002
→ Stage 0: read Bronze (OK)
→ Stage 1: dedup (FAILED after 14 min)
Без correlation ID вопросы «какой Airflow run запустил app_002?» и «как связать падение Stage 1 с метриками Stage 0?» требуют ручной работы - сопоставлять timestamps, dag_run_id из Airflow с application_id из Spark History Server.
Correlation ID / Trace ID - единый идентификатор, который проходит через все слои системы и связывает все события одного логического запуска.
3.2 Схема проброса Trace ID¶
3.3 Генерация Trace ID в Airflow и передача в Spark¶
# dags/orders_pipeline.py
import uuid
from airflow.models import Variable
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
def generate_trace_id(**context) -> str:
"""
Генерируем trace_id один раз для всего DAG-запуска.
Используем комбинацию dag_run_id + короткий UUID для читаемости.
"""
dag_run_id = context["run_id"]
short_uuid = uuid.uuid4().hex[:8]
# Формат: airflow_run_id + короткий суффикс
return f"{dag_run_id[:20]}_{short_uuid}"
with DAG(dag_id="orders_pipeline", ...) as dag:
# Генерация trace_id через BashOperator или XCom
trace_id = "{{ run_id | replace(':', '') | replace('+', '') | truncate(20, False, '') }}_{{ macros.uuid.uuid4().hex[:8] }}"
extract_task = SparkSubmitOperator(
task_id="extract_orders",
application="jobs/extract_orders.py",
application_args=["--date", "{{ ds }}", "--run-id", "{{ run_id }}"],
conf={
# Trace ID передаётся как Spark conf - доступен в коде через SparkConf
"spark.app.trace_id": "{{ run_id }}",
"spark.app.pipeline_name": "orders_pipeline",
"spark.app.dag_run_id": "{{ run_id }}",
"spark.app.task_id": "extract_orders",
},
)
transform_task = SparkSubmitOperator(
task_id="transform_silver",
application="jobs/transform_silver.py",
application_args=["--date", "{{ ds }}", "--run-id", "{{ run_id }}"],
conf={
"spark.app.trace_id": "{{ run_id }}",
"spark.app.pipeline_name": "orders_pipeline",
"spark.app.task_id": "transform_silver",
},
)
extract_task >> transform_task
3.4 Чтение Trace ID в Spark и проброс на Executor'ы¶
Ключевой механизм: LocalProperties в SparkContext. Свойства, установленные через sc.setLocalProperty(), автоматически копируются из Driver-потока во все Task'и на Executor'ах. Это встроенный в Spark механизм распространения контекста.
from pyspark.sql import SparkSession
from pyspark import SparkContext
def setup_trace_context(spark: SparkSession) -> dict:
"""
Читаем trace_id и другой контекст из SparkConf
(куда Airflow передал их через --conf).
Затем устанавливаем в LocalProperties для автоматической
передачи во все Executor Task'и.
"""
sc: SparkContext = spark.sparkContext
conf = sc.getConf()
trace_id = conf.get("spark.app.trace_id", "no_trace")
pipeline_name = conf.get("spark.app.pipeline_name", "unknown")
dag_run_id = conf.get("spark.app.dag_run_id", "unknown")
task_id = conf.get("spark.app.task_id", "unknown")
# Устанавливаем в LocalProperties Spark —
# это магия: свойства автоматически копируются в каждую Task на Executor'е
sc.setLocalProperty("trace_id", trace_id)
sc.setLocalProperty("pipeline_name", pipeline_name)
sc.setLocalProperty("dag_run_id", dag_run_id)
sc.setLocalProperty("task_id", task_id)
# Spark UI тоже умеет показывать эти свойства в Job Description
sc.setJobDescription(
f"pipeline={pipeline_name} | run={dag_run_id[:20]} | trace={trace_id[:8]}"
)
return {
"trace_id": trace_id,
"pipeline_name": pipeline_name,
"dag_run_id": dag_run_id,
"task_id": task_id,
"spark_app_id": sc.applicationId,
}
def main():
spark = SparkSession.builder.appName("TransformSilverOrders").getOrCreate()
# Устанавливаем контекст трассировки
ctx = setup_trace_context(spark)
# Создаём логгер с trace_id
log = SparkPipelineLogger(
pipeline_name=ctx["pipeline_name"],
run_id=ctx["dag_run_id"],
trace_id=ctx["trace_id"],
layer="silver",
)
log.info(
"Job started",
event="job_start",
spark_app_id=ctx["spark_app_id"],
task_id=ctx["task_id"],
)
# ... дальнейшая логика пайплайна
3.5 Использование Trace ID внутри Python UDF (Executor'ы)¶
На Executor'ах LocalProperties доступны через TaskContext:
from pyspark import TaskContext
from pyspark.sql.functions import udf
from pyspark.sql.types import StringType
def create_traced_udf(func):
"""
Обёртка для UDF, которая добавляет trace_id из TaskContext в логи.
TaskContext.get().getLocalProperty("trace_id") возвращает значение,
которое было установлено на Driver через sc.setLocalProperty().
"""
import logging
def traced_func(*args):
# Получаем trace_id из контекста текущей Task
ctx = TaskContext.get()
trace_id = ctx.getLocalProperty("trace_id") if ctx else "no_trace"
partition_id = ctx.partitionId() if ctx else -1
try:
return func(*args)
except Exception as e:
# Логируем с trace_id - теперь можно найти по ID в ELK
logger = logging.getLogger("executor.udf")
logger.error(
f"UDF failed | trace_id={trace_id} | "
f"partition={partition_id} | error={str(e)}"
)
return None
return traced_func
@create_traced_udf
def normalize_phone(phone: str) -> str:
"""Нормализация номера телефона с трассировкой."""
import re
if not phone:
return None
digits = re.sub(r"[^\d]", "", phone)
return f"+7{digits[-10:]}" if len(digits) >= 10 else None
normalize_phone_udf = udf(normalize_phone, StringType())
Часть 4. Spark Listeners: автоматический сбор метрик¶
4.1 Архитектура Spark Listener API¶
Spark предоставляет внутренний event bus. При каждом значимом событии (старт Job, конец Stage, конец Task, завершение SQL-запроса) Spark публикует событие, которое обрабатывают все зарегистрированные Listener'ы. Это не требует изменения кода самого пайплайна - Listener подключается снаружи.
4.2 QueryExecutionListener: метрики записи DataFrame¶
QueryExecutionListener - самый полезный listener для data engineers. Он перехватывает каждую операцию Spark SQL/DataFrame и предоставляет доступ к метрикам выполнения: количество строк, байт, файлов.
from pyspark.sql import SparkSession
from pyspark.sql.execution.listener import QueryExecutionListener
import json
import logging
logger = logging.getLogger("spark.query_listener")
class DataPipelineQueryListener(QueryExecutionListener):
"""
Кастомный QueryExecutionListener для сбора метрик каждой Spark SQL операции.
Автоматически срабатывает после каждого .write(), .count() и других Actions.
Не требует изменения кода пайплайна - подключается один раз при инициализации.
"""
def __init__(self, pipeline_name: str, run_id: str, trace_id: str = None):
self.pipeline_name = pipeline_name
self.run_id = run_id
self.trace_id = trace_id or "no_trace"
def onSuccess(self, funcName: str, qe, durationNs: int) -> None:
"""
Вызывается после успешного завершения любого Action.
qe - QueryExecution объект с доталями физического плана и метриками.
durationNs - длительность выполнения в наносекундах.
"""
duration_ms = durationNs // 1_000_000
# Извлекаем метрики из физического плана
metrics = self._extract_write_metrics(qe)
log_record = {
"event": "query_execution_success",
"pipeline": self.pipeline_name,
"run_id": self.run_id,
"trace_id": self.trace_id,
"func_name": funcName,
"duration_ms": duration_ms,
**metrics,
}
logger.info(json.dumps(log_record))
def onFailure(self, funcName: str, qe, exception: Exception) -> None:
"""
Вызывается при ошибке любого Action.
Логируем детали ошибки с полным контекстом.
"""
log_record = {
"event": "query_execution_failure",
"pipeline": self.pipeline_name,
"run_id": self.run_id,
"trace_id": self.trace_id,
"func_name": funcName,
"error_type": type(exception).__name__,
"error_message": str(exception)[:500],
}
logger.error(json.dumps(log_record))
def _extract_write_metrics(self, qe) -> dict:
"""
Извлекаем метрики из выполненного плана.
Ищем WriteFilesExec или InsertIntoHadoopFsRelation для операций записи.
"""
metrics = {}
try:
# Получаем все метрики из Spark Plan
executed_plan = qe.executedPlan()
for node in executed_plan:
node_name = node.__class__.__name__
# Метрики для операций записи в файлы
if "WriteFiles" in node_name or "InsertInto" in node_name:
spark_metrics = node.metrics()
if spark_metrics:
metrics["num_output_rows"] = (
spark_metrics.get("numOutputRows", {}).get("value", 0)
)
metrics["num_output_bytes"] = (
spark_metrics.get("numOutputBytes", {}).get("value", 0)
)
metrics["num_parts"] = (
spark_metrics.get("numParts", {}).get("value", 0)
)
except Exception:
pass # Метрики опциональны - не ломаем пайплайн при их отсутствии
return metrics
def register_pipeline_listener(
spark: SparkSession,
pipeline_name: str,
run_id: str,
trace_id: str = None,
) -> DataPipelineQueryListener:
"""
Регистрируем listener в SparkSession.
После этого все Actions автоматически логируются.
"""
listener = DataPipelineQueryListener(pipeline_name, run_id, trace_id)
spark.streams.addListener(listener) # для Structured Streaming
spark._jvm.org.apache.spark.sql.execution.QueryExecution.listeners().add(
listener._jlistener # для Batch DataFrame
)
return listener
4.3 SparkListener для Task-уровневых метрик¶
Для мониторинга data skew, spill и GC нужен SparkListener на уровне JVM:
def register_task_metrics_listener(spark: SparkSession, log: SparkPipelineLogger):
"""
Регистрируем SparkListener через Py4J для мониторинга Task-метрик.
Позволяет детектировать data skew и GC pressure без Spark UI.
"""
from py4j.java_gateway import java_import
sc = spark.sparkContext
jvm = sc._jvm
jsc = sc._jsc
java_import(jvm, "org.apache.spark.scheduler.*")
class TaskMetricsListener:
"""Python-обёртка над Java SparkListener."""
class Java:
implements = ["org.apache.spark.scheduler.SparkListener"]
def onTaskEnd(self, task_end_event):
task_info = task_end_event.taskInfo()
task_metrics = task_end_event.taskMetrics()
if task_metrics is None:
return
# Ключевые метрики каждой Task
shuffle_read = task_metrics.shuffleReadMetrics().totalBytesRead()
shuffle_write = task_metrics.shuffleWriteMetrics().bytesWritten()
spill_disk = task_metrics.diskBytesSpilled()
spill_memory = task_metrics.memoryBytesSpilled()
gc_time = task_metrics.jvmGCTime()
duration = task_info.duration()
# Детектируем аномалии
if spill_disk > 1024 * 1024 * 1024: # spill > 1GB
log.warning(
"Large disk spill detected",
event="disk_spill_alert",
task_id=task_info.taskId(),
stage_id=task_end_event.stageId(),
spill_disk_bytes=spill_disk,
spill_gb=round(spill_disk / (1024**3), 2),
)
if gc_time > 30_000: # GC > 30 секунд
log.warning(
"Excessive GC time",
event="gc_pressure_alert",
task_id=task_info.taskId(),
gc_time_ms=gc_time,
task_duration_ms=duration,
gc_ratio=round(gc_time / duration, 2) if duration > 0 else 0,
)
listener = TaskMetricsListener()
jsc.sc().addSparkListener(listener)
return listener
4.4 Детекция Data Skew через Task-метрики¶
Data skew - одна из самых частых причин деградации производительности в Spark. Один Executor обрабатывает миллион строк, остальные - по тысяче. Весь Job ждёт самого медленного.
def detect_data_skew(
spark: SparkSession,
app_id: str,
stage_id: int,
skew_ratio_threshold: float = 5.0,
) -> dict:
"""
Анализируем Task-метрики конкретного Stage для детекции data skew.
Skew определяем как: max_task_duration / avg_task_duration > threshold.
В production используем Spark History Server REST API или SparkListener.
Здесь показана логика обнаружения.
"""
# В реальной системе Task-метрики получаем из Spark History Server API:
# GET /api/v1/applications/{appId}/stages/{stageId}/taskList
# Для демонстрации - псевдокод с реальными именами метрик
task_durations = [] # заполняется из Task End Events
if not task_durations:
return {"status": "no_data"}
max_duration = max(task_durations)
avg_duration = sum(task_durations) / len(task_durations)
min_duration = min(task_durations)
skew_ratio = max_duration / avg_duration if avg_duration > 0 else 0
result = {
"stage_id": stage_id,
"task_count": len(task_durations),
"max_duration_ms": max_duration,
"avg_duration_ms": round(avg_duration),
"min_duration_ms": min_duration,
"skew_ratio": round(skew_ratio, 2),
"skew_detected": skew_ratio > skew_ratio_threshold,
}
if result["skew_detected"]:
print(
f"[SKEW ALERT] Stage {stage_id}: max={max_duration}ms, "
f"avg={avg_duration:.0f}ms, ratio={skew_ratio:.1f}x. "
f"Consider salting or repartitioning."
)
return result
Часть 5. Метрики и Prometheus¶
5.1 Встроенный PrometheusServlet в Spark 3.x¶
Начиная со Spark 3.0, встроен HTTP endpoint для экспорта метрик в формате Prometheus. Активируется через конфигурацию:
# spark-defaults.conf или --conf при spark-submit
# Включаем Prometheus endpoint на Driver
spark.ui.prometheus.enabled=true
# Spark metrics system
spark.metrics.conf.*.sink.prometheus.class=org.apache.spark.metrics.sink.PrometheusServlet
spark.metrics.conf.*.sink.prometheus.path=/metrics/prometheus
# Метрики JVM (GC, heap, threads)
spark.metrics.conf.driver.source.jvm.class=org.apache.spark.metrics.source.JvmSource
spark.metrics.conf.executor.source.jvm.class=org.apache.spark.metrics.source.JvmSource
После этого http://<driver-host>:4040/metrics/prometheus возвращает:
# HELP spark_driver_jvm_heap_used_bytes JVM heap used
# TYPE spark_driver_jvm_heap_used_bytes gauge
spark_driver_jvm_heap_used_bytes{app_id="application_001"} 2147483648
# HELP spark_executor_shuffleRead_totalBytesRead Total shuffle read bytes
# TYPE spark_executor_shuffleRead_totalBytesRead counter
spark_executor_shuffleRead_totalBytesRead{executor_id="1"} 53687091200
5.2 Кастомные application-метрики через Dropwizard¶
Для business-level метрик (количество обработанных строк, quarantine rate, SLA) используем Dropwizard Metrics, который Spark использует внутренне:
from pyspark import SparkContext
class SparkCustomMetrics:
"""
Регистрируем кастомные метрики в Spark Metrics Registry.
Они автоматически экспортируются через PrometheusServlet.
"""
def __init__(self, sc: SparkContext, namespace: str):
self.sc = sc
self.namespace = namespace
# Получаем Java Metrics Registry через Py4J
self._registry = sc._jvm.org.apache.spark.metrics.MetricsSystem
def register_gauge(self, name: str, value_func):
"""Регистрируем Gauge - метрику, читаемую в момент опроса."""
# Через Py4J создаём Java Gauge
jvm = self.sc._jvm
gauge = jvm.com.codahale.metrics.Gauge.of(value_func)
# self._registry.registerGauge(f"{self.namespace}.{name}", gauge)
def increment_counter(self, name: str, amount: int = 1):
"""Инкремент счётчика - монотонно возрастающая метрика."""
pass # В реальной реализации через Py4J
def report_pipeline_metrics(
spark: SparkSession,
pipeline_name: str,
stats: dict,
) -> None:
"""
Отправляем кастомные метрики через statsd или напрямую в Prometheus
(через pushgateway для batch jobs).
"""
import socket
import time
# Формат statsd: metric_name:value|type
metrics_lines = [
f"spark.pipeline.rows_processed:{stats['total_rows']}|g",
f"spark.pipeline.rows_clean:{stats['clean_rows']}|g",
f"spark.pipeline.quarantine_rate:{stats['quarantine_rate']:.4f}|g",
f"spark.pipeline.duration_seconds:{stats['duration_seconds']:.1f}|g",
]
# Тэги для идентификации метрик
tags = f",pipeline={pipeline_name},env=prod"
try:
sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
for line in metrics_lines:
message = f"{line}{tags}".encode()
sock.sendto(message, ("statsd-host", 8125))
sock.close()
except Exception as e:
print(f"[Metrics] Failed to send to statsd: {e}")
5.3 Grafana Dashboard для Spark ETL¶
Ключевые дашборды для инженера данных (запросы на PromQL):
# Время выполнения последних 10 запусков пайплайна
avg_over_time(
spark_pipeline_duration_seconds{pipeline="bronze_to_silver_orders"}[24h]
)
# Quarantine rate - тренд за 7 дней (алерт при росте)
spark_pipeline_quarantine_rate{pipeline="bronze_to_silver_orders"}
> 0.05 # alert threshold
# Shuffle bytes как индикатор data skew
sum by (stage_id) (
spark_executor_shuffleRead_totalBytesRead
)
# GC pressure на Executor'ах (алерт > 20% времени в GC)
sum(rate(spark_executor_jvm_gc_time_total[5m])) /
sum(rate(spark_executor_duration_total[5m])) > 0.20
# Количество spill to disk (должно быть 0 в нормальной ситуации)
sum(spark_executor_diskBytesSpilled_total)
Часть 6. Observability в Medallion Architecture¶
6.1 Разные метрики для разных слоёв¶
Каждый слой Medallion требует собственного набора observability signals:
from dataclasses import dataclass, field
from typing import Optional
@dataclass
class LayerObservabilityEvent:
"""Структурированное событие observability для конкретного слоя."""
pipeline: str
run_id: str
trace_id: str
layer: str # bronze | silver | gold
table: str
processing_date: str
timestamp: str
# Bronze-специфичные метрики
rows_ingested: Optional[int] = None
bytes_ingested: Optional[int] = None
parse_errors: Optional[int] = None
schema_drift_detected: Optional[bool] = None
source_lag_seconds: Optional[int] = None # задержка от источника
# Silver-специфичные метрики
rows_before_dedup: Optional[int] = None
rows_after_dedup: Optional[int] = None
duplicates_removed: Optional[int] = None
quarantine_rows: Optional[int] = None
quarantine_rate: Optional[float] = None
null_violations: Optional[dict] = None
# Gold-специфичные метрики
aggregation_groups: Optional[int] = None
metric_delta_pct: Optional[float] = None # % изменения ключевой метрики
freshness_lag_minutes: Optional[int] = None # задержка данных
6.2 Data Freshness Monitoring¶
Свежесть данных - критическая метрика для Gold слоя. Аналитики должны знать: данные на дашборде за какой момент?
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from datetime import datetime, timedelta
def check_data_freshness(
spark: SparkSession,
table_path: str,
event_time_col: str,
max_lag_minutes: int = 30,
log: SparkPipelineLogger = None,
) -> dict:
"""
Проверяем свежесть данных в таблице.
Сравниваем MAX(event_time) с текущим временем.
Если разница > max_lag_minutes - алерт.
"""
result = (
spark.read.format("delta").load(table_path)
.agg(
F.max(event_time_col).alias("latest_event"),
F.min(event_time_col).alias("earliest_event"),
F.count("*").alias("total_rows"),
)
.collect()[0]
)
latest_event = result["latest_event"]
now = datetime.utcnow()
lag_minutes = None
if latest_event:
lag_seconds = (now - latest_event.replace(tzinfo=None)).total_seconds()
lag_minutes = int(lag_seconds / 60)
freshness_ok = lag_minutes is not None and lag_minutes <= max_lag_minutes
if log:
log.info(
"Data freshness check",
event="freshness_check",
table=table_path,
latest_event=str(latest_event),
lag_minutes=lag_minutes,
threshold_minutes=max_lag_minutes,
freshness_ok=freshness_ok,
)
if not freshness_ok:
log.warning(
"Data freshness SLA violated",
event="freshness_sla_violation",
table=table_path,
lag_minutes=lag_minutes,
threshold_minutes=max_lag_minutes,
)
return {
"latest_event": str(latest_event),
"lag_minutes": lag_minutes,
"freshness_ok": freshness_ok,
"total_rows": result["total_rows"],
}
Часть 7. Debugging с Observability: реальный сценарий¶
7.1 Инцидент: OOM на Silver transform¶
Рассмотрим реальный сценарий расследования инцидента с готовой системой observability.
Симптом: задача transform_silver_orders упала в 02:17. SLA - 03:00.
Шаг 1: фильтрация логов по trace_id
В Kibana/Grafana Loki фильтруем:
{pipeline="bronze_to_silver_orders"} |= "2024-01-15"
Находим события:
{"event": "stage_end", "stage": "bronze_read", "rows_read": 52348910, "duration_ms": 82341}
{"event": "rows_processed", "stage": "deduplication", "rows_in": 52348910, ...}
{"event": "query_execution_failure", "error_type": "OutOfMemoryError", "func_name": "save"}
rows_read: 52348910 - в 10 раз больше обычного! Bronze ingestion загрузил в 10 раз больше строк. Подозрение: дублирование или ошибка в фильтрации партиций.
Шаг 2: Spark History Server
По spark_app_id из лога открываем History Server. В SQL tab находим Query план:
Filter (isnotnull(processing_date#12))
Scan parquet datalake/bronze/orders
PartitionFilters: []
PushedFilters: [IsNotNull(processing_date)]
PartitionFilters: [] - Spark не использовал partition pruning! Прочитал всю таблицу, включая исторические данные.
Шаг 3: корень проблемы
Находим коммит: вчера изменили тип колонки processing_date с date на string. Spark потерял возможность делать partition pruning по строковому предикату с DateType-сравнением. Читал всю Bronze-таблицу за два года.
Без observability: 4+ часов расследования, несколько рестартов, ручной grep по логам. С observability: 15 минут от алерта до корневой причины.
7.2 Полный observability-стек в одной функции¶
def run_observable_pipeline(
spark: SparkSession,
batch_date: str,
run_id: str,
) -> None:
"""
Демонстрация полного observability-стека в одном пайплайне:
- Structured logging с JSON
- Correlation ID / trace propagation
- SparkContext local properties
- Автоматические метрики через QueryExecutionListener
- Data freshness check
- Business metrics logging
"""
import time
from datetime import datetime
# 1. Инициализация trace context
ctx = setup_trace_context(spark)
log = SparkPipelineLogger(
pipeline_name="observable_orders_pipeline",
run_id=run_id,
trace_id=ctx["trace_id"],
layer="silver",
)
# 2. Регистрируем QueryExecutionListener
register_pipeline_listener(
spark,
pipeline_name="observable_orders_pipeline",
run_id=run_id,
trace_id=ctx["trace_id"],
)
pipeline_start = time.time()
log.info(
"Observable pipeline started",
event="pipeline_start",
spark_app_id=ctx["spark_app_id"],
batch_date=batch_date,
)
# 3. Bronze read с логированием
t0 = log.stage_start("bronze_read", batch_date=batch_date)
bronze_df = spark.read.format("delta").load(
f"s3a://datalake/bronze/orders/date={batch_date}/"
)
bronze_count = bronze_df.count()
log.stage_end("bronze_read", t0, rows_read=bronze_count)
# 4. Quality check с логированием
log.quality_gate(
gate_name="bronze_row_count",
passed=(bronze_count > 0),
metric_value=bronze_count,
threshold=1,
)
# 5. Transform с логированием каждого шага
t1 = log.stage_start("deduplication")
silver_df = bronze_df.dropDuplicates(["order_id"])
silver_count = silver_df.count()
log.stage_end("deduplication", t1, rows_deduplicated=bronze_count - silver_count)
log.rows_processed("deduplication", bronze_count, silver_count, table="orders")
# 6. Write - QueryExecutionListener логирует автоматически
t2 = log.stage_start("silver_write")
(
silver_df
.write.mode("overwrite")
.format("delta")
.partitionBy("processing_date")
.save("s3a://datalake/silver/orders/")
)
log.stage_end("silver_write", t2, rows_written=silver_count)
# 7. Freshness check
freshness = check_data_freshness(
spark,
"s3a://datalake/silver/orders/",
event_time_col="created_at",
max_lag_minutes=120,
log=log,
)
# 8. Итоговая метрика пайплайна
pipeline_duration = time.time() - pipeline_start
log.info(
"Pipeline complete",
event="pipeline_complete",
duration_seconds=round(pipeline_duration, 1),
bronze_rows=bronze_count,
silver_rows=silver_count,
dedup_ratio=round(1 - silver_count / bronze_count, 4) if bronze_count else 0,
data_freshness_ok=freshness["freshness_ok"],
freshness_lag_minutes=freshness["lag_minutes"],
)
# 9. Отправка business-метрик в Prometheus/StatsD
report_pipeline_metrics(
spark,
pipeline_name="observable_orders_pipeline",
stats={
"total_rows": bronze_count,
"clean_rows": silver_count,
"quarantine_rate": 0.0,
"duration_seconds": pipeline_duration,
},
)
Часть 8. Anti-Patterns в Observability¶
8.1 Logging Everything: избыточность убивает сигнал¶
# АНТИПАТТЕРН: логируем каждую строку в цикле
for row in df.collect(): # уже антипаттерн - collect на большом DF
logger.debug(f"Processing row: {row}") # миллионы строк в логах
# ПРАВИЛЬНО: агрегированные метрики, не отдельные строки
logger.info("Batch processed", rows=df.count(), stage="transform")
Логи - это не замена данным. Задача лога - ответить на вопрос «что произошло и когда», а не воспроизвести каждую строку датасета.
8.2 Missing Correlation ID¶
# АНТИПАТТЕРН: каждый шаг логирует независимо
# В ELK невозможно связать logs одного pipeline run
logger.info("Bronze read complete: 5M rows") # нет run_id, нет trace_id
logger.info("Silver write complete: 4.9M rows") # это тот же запуск? другой?
# ПРАВИЛЬНО: единый trace_id через всю цепочку
log.info("Bronze read complete", trace_id=trace_id, rows=5_000_000)
log.info("Silver write complete", trace_id=trace_id, rows=4_900_000)
8.3 Noisy Alerts: alert fatigue¶
# АНТИПАТТЕРН: алерт на каждое небольшое отклонение
if current_rows < previous_rows: # срабатывает при любом снижении, даже на 1 строку
send_pagerduty_alert("Row count decreased!")
# ПРАВИЛЬНО: осмысленный порог с контекстом
drop_ratio = (previous_rows - current_rows) / previous_rows
if drop_ratio > 0.30: # > 30% снижение - реальная проблема
send_alert(
level="critical",
message=f"Row count dropped by {drop_ratio:.1%}",
context={"previous": previous_rows, "current": current_rows},
)
8.4 Observability Theater: метрик много, пользы нет¶
Самый опасный антипаттерн: команда создаёт десятки дашбордов, сотни алертов, гигабайты логов - но при инциденте всё равно не может быстро найти проблему.
Признаки «observability theater»:
- Дашборды показывают CPU и Memory, но не бизнес-метрики (rows_processed, quarantine_rate)
- Алерты срабатывают > 50% ложных тревог → команда их игнорирует
- Нет correlation ID → логи разных слоёв невозможно связать
- Метрики есть, но нет контекста (в каком pipeline, в каком run?)
# ПРАВИЛЬНЫЙ подход: минимальный, но полезный observability stack
ESSENTIAL_SIGNALS = {
# Один структурированный лог с trace_id - лучше тысячи разрозненных
"logging": "structured JSON with trace_id",
# Три ключевые метрики на пайплайн - лучше 50 технических
"metrics": [
"pipeline_duration_seconds", # как долго?
"rows_quarantine_rate", # насколько чисты данные?
"data_freshness_lag_minutes", # насколько свежи данные?
],
# Один алерт на pipeline - лучше 20 шумных
"alerting": "quarantine_rate > 5% OR duration > 2x_sla OR freshness_lag > 2h",
}
8.5 Observability как afterthought¶
Самая системная ошибка: observability добавляется «потом», после того как пайплайн уже работает в production. В итоге архитектура не предусматривает correlation ID, функции возвращают void вместо статистики, нет стандартного формата логов.
# Признак правильной архитектуры: каждая функция возвращает observable результат
# ПЛОХО: функция ничего не возвращает, observability невозможна
def process_bronze(df):
df.write.parquet(silver_path)
# void - не знаем, сколько строк записано
# ХОРОШО: функция возвращает статистику, observability встроена
def process_bronze(df, log: SparkPipelineLogger) -> dict:
t = log.stage_start("write_silver")
(df.write.format("delta").save(silver_path))
count = df.count()
log.stage_end("write_silver", t, rows_written=count)
return {"rows_written": count, "path": silver_path}
Часть 9. Чек-лист готовности пайплайна к production¶
9.1 Observability Readiness Checklist¶
Перед деплоем нового Spark-пайплайна в production:
OBSERVABILITY_CHECKLIST = {
"structured_logging": [
"JSON-формат для всех логов (не plain text)",
"Каждое событие содержит: pipeline, run_id, trace_id, timestamp",
"Числовые метрики в полях (rows_read, duration_ms) - не в строках",
"log.stage_start()/stage_end() на каждую значимую операцию",
],
"tracing": [
"trace_id передаётся из Airflow через --conf",
"sc.setLocalProperty('trace_id', ...) вызывается при старте",
"Python UDF имеют доступ к trace_id через TaskContext",
"Все Spark jobs описаны: sc.setJobDescription()",
],
"metrics": [
"Итоговые метрики пайплайна логируются в pipeline_complete event",
"Quarantine rate присутствует если есть DLQ",
"Data freshness проверяется на Gold слое",
"Метрики отправляются в Prometheus/StatsD",
],
"alerting": [
"Алерт при quarantine_rate > 5%",
"Алерт при duration > 2x исторического среднего",
"Алерт при data freshness lag > SLA",
"Алерт при пустом батче (rows == 0)",
"Алерт не чаще 1 раза в 15 минут (группировка)",
],
"debugging": [
"Spark History Server настроен и хранит event logs > 30 дней",
"По trace_id можно найти все логи запуска в ELK/Loki",
"По application_id можно найти Stage/Task метрики в History Server",
"QueryExecutionListener зарегистрирован для автоматических write-метрик",
],
}
9.2 Минимальный рабочий observability-стек¶
Для команды, которая только начинает:
Фаза 1 (MVP): Structured JSON logs → Grafana Loki (Docker Compose) → один Grafana Dashboard с 5 метриками → Slack-алерты.
Фаза 2: Prometheus + PrometheusServlet → расширенные дашборды → alertmanager.
Фаза 3: Distributed tracing (OpenTelemetry/Jaeger) → полный lineage tracking → SLO dashboards.
Итоги¶
Observability - это не набор инструментов, а инженерная дисциплина. Grafana без структурированных логов - бесполезна. Логи без correlation ID - разрозненные тексты. Алерты без контекста - шум.
Structured Logging - фундамент. Каждое событие - JSON с pipeline, run_id, trace_id, числовыми метриками. Без этого невозможны автоматические дашборды и фильтрация в ELK.
Correlation ID / trace_id - связующее звено. Передаётся из Airflow через --conf, устанавливается в SparkContext через setLocalProperty(), автоматически копируется на все Executor'ы. По одному ID можно найти все события запуска от Airflow до S3.
SparkContext.setLocalProperty() - встроенный механизм распространения контекста из Driver в Task'и. Не нужно изобретать велосипед - это официальный API Spark для correlation ID.
QueryExecutionListener - автоматический сбор метрик без изменения кода пайплайна. Регистрируем один раз - получаем rows_written, bytes_written, duration_ms для каждого Action автоматически.
Data freshness - ключевая бизнес-метрика Gold слоя. MAX(event_time) vs NOW() должны быть в пределах SLA. Аналитики должны видеть не только «данные есть», но и «данные за какой момент».
Минимальный стек для начала: python-json-logger + trace_id в каждом log record + PrometheusServlet + Grafana Loki. Это даёт 80% observability за 20% усилий.