Observability: структурированное логирование и трассировка Spark job

Три столпа observability в data engineering: structured logging (JSON logs, python-json-logger), distributed tracing (correlation IDs, SparkContext local properties), Spark Listeners для автоматических метрик, интеграция с Prometheus и Grafana

streaming

Введение: проблема «чёрного ящика» в распределённых вычислениях

Представьте: в 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/Grafana
  • trace_id - сквозной идентификатор, связывающий все логи одного запуска через Airflow, Spark Driver и Executor'ы
  • rows_dropped + drop_ratio - числовые метрики, которые можно агрегировать в dashboard
  • spark_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% усилий.