Блок 5 - Практические кейсы

20 сценарных вопросов: top-N, incremental load, OOM, big joins, Delta Lake, TF-IDF и разбор реальных ситуаций.

core optimization

Как найти топ-3 зарплаты в каждом отделе?

Классический вопрос на Window Functions. Два подхода:

from pyspark.sql.window import Window
from pyspark.sql.functions import dense_rank, col

# Подход 1: Window + dense_rank (правильный - сохраняет связи)
w = Window.partitionBy("dept").orderBy(col("salary").desc())
df.withColumn("rnk", dense_rank().over(w)) \
  .filter(col("rnk") <= 3) \
  .drop("rnk")

# Подход 2: через SQL (то же самое)
df.createOrReplaceTempView("employees")
spark.sql("""
    SELECT dept, name, salary FROM (
        SELECT dept, name, salary,
               DENSE_RANK() OVER (PARTITION BY dept ORDER BY salary DESC) AS rnk
        FROM employees
    ) WHERE rnk <= 3
""")

RANK vs DENSE_RANK: при одинаковых зарплатах RANK пропускает номера (1,1,3), DENSE_RANK - нет (1,1,2). Для «топ-N» обычно нужен DENSE_RANK.


Опишите процесс инкрементальной загрузки данных

Инкрементальная загрузка = загружаем только новые/изменённые записи, не перезаписывая всё.

Ключевые элементы:

  1. Водяной знак (watermark) - последнее значение updated_at из предыдущего запуска (хранить в Airflow variable или отдельной таблице)
  2. Дедупликация - left_anti join новых против существующих по business key
  3. Upsert - MERGE INTO Iceberg/Delta для обновления существующих + вставки новых
  4. Идемпотентность - повторный запуск не должен дублировать данные

Как прочитать только новые файлы из директории?

# Подход 1: Auto Loader (Databricks/OSS)
df = spark.readStream \
    .format("cloudFiles") \
    .option("cloudFiles.format", "parquet") \
    .load("s3://bucket/incoming/")

# Подход 2: Structured Streaming с maxFilesPerTrigger
df = spark.readStream \
    .format("parquet") \
    .option("maxFilesPerTrigger", 100) \
    .load("s3://bucket/incoming/")

# Подход 3: Batch - читать только файлы новее определённой даты
from datetime import datetime, timedelta
cutoff = datetime.now() - timedelta(hours=1)

import boto3
s3 = boto3.client("s3")
new_files = [f for f in list_files("bucket", "prefix")
             if f["LastModified"] > cutoff]
df = spark.read.parquet(*new_files)

Streaming-подход (Auto Loader) предпочтителен для production - Spark отслеживает обработанные файлы через checkpoint.


Что делать, если задача падает с OutOfMemoryError?

Диагностика по месту ошибки:

Место Причина Решение
Driver OOM collect(), broadcast слишком большой Избегать collect; поднять driver.memory; проверить broadcast
Executor OOM Skew - одна партиция огромная AQE Skew Join; salting; увеличить executor.memory
Executor OOM Слишком мало партиций Увеличить shuffle.partitions; repartition()
Executor OOM UDF утечка памяти Профилировать Python-код; Pandas UDF вместо row-by-row

Практический чеклист:

  1. spark.executor.memory + spark.executor.memoryOverhead - суммарно ≥ container size
  2. spark.sql.shuffle.partitions - поднять пропорционально объёму данных
  3. spark.sql.adaptive.enabled=true - включить AQE
  4. Проверить Spark UI → Stages → Task max duration → skew?

Как соединить таблицу 1 ТБ и таблицу 100 МБ?

100 МБ - кандидат для Broadcast Join:

from pyspark.sql.functions import broadcast

# Явный hint
result = big_1tb.join(broadcast(small_100mb), "key")

# Поднять порог автоматического broadcast
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "200m")

# AQE сделает это автоматически если small_100mb окажется маленькой после фильтра
spark.conf.set("spark.sql.adaptive.enabled", "true")

При Broadcast Join: 100 МБ × N Executors = сетевой трафик. При 100 Executor'ах - 10 ГБ трафика. Если Executor'ов очень много - проверить хватает ли памяти на каждом.


Сценарий: join двух таблиц по 1 ТБ каждая - ваши действия?

1. EXPLAIN → посмотреть план: SortMergeJoin или ShuffleHashJoin?
2. ANALYZE TABLE both_tables COMPUTE STATISTICS → дать CBO статистику
3. Проверить Data Skew: есть ли hot keys?
   → Если есть: salting или AQE SkewJoin

4. Проверить партиционирование:
   → Если таблицы bucketed по join-key с одинаковым N → shuffle не нужен

5. Настроить ресурсы:
   → shuffle.partitions: для 1 ТБ × 2 = ~2000 партиций (по ~1 ГБ каждая)
   → executor.memory: достаточно для sort буферов

6. Включить AQE:
   → spark.sql.adaptive.enabled=true
   → spark.sql.adaptive.skewJoin.enabled=true

7. Рассмотреть предварительный фильтр:
   → Уменьшить одну таблицу до join? (predicate pushdown)

Как посчитать количество уникальных пользователей за каждый час?

from pyspark.sql.functions import date_trunc, countDistinct, window, col

# Вариант 1: truncate до часа + groupBy
df.withColumn("hour", date_trunc("hour", col("event_time"))) \
  .groupBy("hour") \
  .agg(countDistinct("user_id").alias("unique_users")) \
  .orderBy("hour")

# Вариант 2: через window (для streaming или batch)
df.groupBy(window("event_time", "1 hour")) \
  .agg(countDistinct("user_id").alias("unique_users"))

# Вариант 3: approx (быстрее при огромных данных, ~5% погрешность)
from pyspark.sql.functions import approx_count_distinct
df.groupBy("hour") \
  .agg(approx_count_distinct("user_id", rsd=0.05).alias("approx_unique"))

approx_count_distinct использует алгоритм HyperLogLog - в 10–100× быстрее точного countDistinct на больших данных.


Как развернуть (pivot) таблицу?

# Исходная таблица: dept, month, amount
# Цель: dept | jan | feb | mar

df.groupBy("dept") \
  .pivot("month", ["jan", "feb", "mar"]) \
  .agg(sum("amount"))

# Без явного списка значений (Spark сам определит, но медленнее - 2 прохода)
df.groupBy("dept").pivot("month").agg(sum("amount"))

# Обратный pivot (unpivot) - через stack
from pyspark.sql.functions import expr
df.select("dept",
    expr("stack(3, 'jan', jan, 'feb', feb, 'mar', mar) as (month, amount)"))

Как нормализовать данные перед обучением модели?

from pyspark.ml.feature import StandardScaler, MinMaxScaler, VectorAssembler

# Шаг 1: собрать features в вектор
assembler = VectorAssembler(inputCols=["age", "salary", "score"],
                            outputCol="features_raw")

# Шаг 2а: StandardScaler - (x - mean) / stddev
scaler = StandardScaler(inputCol="features_raw", outputCol="features",
                        withMean=True, withStd=True)

# Шаг 2б: MinMaxScaler - [0, 1]
scaler = MinMaxScaler(inputCol="features_raw", outputCol="features",
                      min=0.0, max=1.0)

from pyspark.ml import Pipeline
pipeline = Pipeline(stages=[assembler, scaler])
model = pipeline.fit(train_df)
scaled = model.transform(train_df)

Scaler обязательно обучается (fit) только на train, применяется (transform) на train и test - иначе data leakage.


Как обработать JSON с динамической схемой?

from pyspark.sql.functions import from_json, schema_of_json, col

# Вариант 1: schema_of_json - определить схему из одной записи (если схема стабильна)
sample_json = '{"id": 1, "tags": ["a", "b"], "meta": {"source": "web"}}'
schema = schema_of_json(sample_json)
df.withColumn("parsed", from_json(col("json_str"), schema))

# Вариант 2: get_json_object - извлечь конкретные поля без полной схемы
df.withColumn("id",     get_json_object(col("json_str"), "$.id")) \
  .withColumn("source", get_json_object(col("json_str"), "$.meta.source"))

# Вариант 3: PERMISSIVE + _corrupt_record для аномалий
df = spark.read \
    .option("mode", "PERMISSIVE") \
    .option("columnNameOfCorruptRecord", "_bad") \
    .json("data.json")
bad_records = df.filter(col("_bad").isNotNull())

Как реализовать поиск по тексту (TF-IDF)?

from pyspark.ml.feature import Tokenizer, StopWordsRemover, HashingTF, IDF

# Пайплайн TF-IDF
tokenizer = Tokenizer(inputCol="text", outputCol="words")
remover   = StopWordsRemover(inputCol="words", outputCol="filtered",
                             stopWords=StopWordsRemover.loadDefaultStopWords("russian"))
hashing   = HashingTF(inputCol="filtered", outputCol="raw_features", numFeatures=10000)
idf       = IDF(inputCol="raw_features", outputCol="features")

from pyspark.ml import Pipeline
pipeline = Pipeline(stages=[tokenizer, remover, hashing, idf])
model = pipeline.fit(df)
result = model.transform(df)
# result.features - разреженный вектор TF-IDF весов

Для сходства документов - косинусное расстояние через BucketedRandomProjectionLSH (Approximate Nearest Neighbors).


Как переложить данные из HDFS в S3?

# Вариант 1: Spark job (рекомендуется для больших данных - трансформация при переносе)
df = spark.read.parquet("hdfs://namenode/data/table/")
df.write.parquet("s3a://bucket/data/table/")

# Вариант 2: hadoop distcp (инструмент Hadoop - копирует без трансформации)
# hadoop distcp hdfs://namenode/data/ s3a://bucket/data/

# Конфигурация S3A
spark.conf.set("fs.s3a.endpoint",            "https://s3.amazonaws.com")
spark.conf.set("fs.s3a.access.key",          "ACCESS_KEY")
spark.conf.set("fs.s3a.secret.key",          "SECRET_KEY")
spark.conf.set("fs.s3a.impl",                "org.apache.hadoop.fs.s3a.S3AFileSystem")
spark.conf.set("fs.s3a.committer.magic.enabled", "true")  # Magic Committer

При переносе рекомендуется rewrite to optimal format: parquet + zstd-сжатие + правильное партиционирование.


Разница Spark на On-premise vs Облаке (AWS EMR/Databricks)

Аспект On-premise AWS EMR Databricks
Управление кластером Полностью вручную Managed YARN/K8s Полностью managed
Масштабирование Статично (долго) Auto-scaling Auto-scaling (быстрый)
Стоимость CAPEX (железо) OPEX (Spot instance) OPEX (DBU)
Формат хранения HDFS S3 + Delta Lake Delta Lake (нативно)
Движок Open-source Spark Open-source Spark Photon (C++ SIMD)
SQL интерфейс Spark SQL Athena / Spark SQL Databricks SQL
Self-host Да Нет Нет

Наш курс ориентирован на self-host: MinIO вместо S3, Iceberg вместо Delta, K8s вместо managed YARN.


Что такое Delta Lake и какие проблемы решает?

Delta Lake - open-source table format (от Databricks), добавляющий ACID-транзакции поверх Parquet на S3/HDFS. Решаемые проблемы:

  • Атомарность - запись либо целиком применяется, либо откатывается (transaction log)
  • Нет «dirty reads» - читатели видят только закоммиченные данные
  • Schema enforcement - нельзя записать данные с несовместимой схемой
  • Time Travel - запрос к историческим версиям (VERSION AS OF N, TIMESTAMP AS OF)
  • DML операции - UPDATE, DELETE, MERGE поверх Parquet

Аналог: Apache Iceberg (более открытый стандарт, поддерживается Spark, Flink, Trino).


Как сделать upsert данных в Delta-таблицу?

from delta.tables import DeltaTable

delta_table = DeltaTable.forPath(spark, "s3://bucket/delta/users/")

# MERGE - upsert: обновить существующие, вставить новые
delta_table.alias("target").merge(
    new_data.alias("source"),
    "target.user_id = source.user_id"
).whenMatchedUpdate(set={
    "name":       "source.name",
    "email":      "source.email",
    "updated_at": "source.updated_at"
}).whenNotMatchedInsert(values={
    "user_id":    "source.user_id",
    "name":       "source.name",
    "email":      "source.email",
    "updated_at": "source.updated_at"
}).execute()

Эквивалент для Iceberg: spark.sql("MERGE INTO target USING source ON ...")


Как работает Time Travel в Delta Lake?

Delta Lake хранит transaction log - JSON-файлы с описанием каждой операции. Time Travel позволяет запрашивать данные на момент конкретной версии или времени.

# По версии
df = spark.read.format("delta") \
    .option("versionAsOf", 5) \
    .load("s3://bucket/delta/users/")

# По времени
df = spark.read.format("delta") \
    .option("timestampAsOf", "2024-01-15 10:00:00") \
    .load("s3://bucket/delta/users/")

# SQL синтаксис
spark.sql("SELECT * FROM users VERSION AS OF 5")
spark.sql("SELECT * FROM users TIMESTAMP AS OF '2024-01-15'")

# Откат таблицы к версии
delta_table.restoreToVersion(3)

История хранится пока не запущен VACUUM (удаляет файлы старше retention периода, default 7 дней).


Зачем использовать формат Avro в стриминге?

Avro - row-based бинарный формат со встроенной schema. Преимущества для стриминга:

  • Schema Registry интеграция - Kafka producer/consumer автоматически проверяют совместимость схем; schema хранится в Confluent Schema Registry
  • Compact - хорошее сжатие для небольших событий
  • Schema evolution - поддерживает добавление полей с default-значениями без поломки потребителей
# Kafka + Avro через Schema Registry (with confluent_kafka)
df = spark.readStream.format("kafka") \
    .option("subscribe", "events") \
    .load()

# Десериализация Avro (with spark-avro)
from pyspark.sql.functions import from_avro
avro_schema = "..."  # JSON-строка схемы
df.select(from_avro(col("value"), avro_schema).alias("data"))

Как автоматизировать запуск Spark-задач (Airflow)?

from airflow import DAG
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator
from airflow.providers.cncf.kubernetes.operators.pod import KubernetesPodOperator
from datetime import datetime

with DAG("daily_etl", schedule_interval="0 6 * * *",
         start_date=datetime(2024, 1, 1), catchup=False) as dag:

    # Вариант 1: SparkSubmitOperator (YARN/Standalone)
    run_job = SparkSubmitOperator(
        task_id="run_spark_job",
        application="s3://bucket/jobs/etl.py",
        conf={"spark.executor.memory": "4g",
              "spark.executor.cores": "2"},
    )

    # Вариант 2: KubernetesPodOperator (K8s)
    run_k8s = KubernetesPodOperator(
        task_id="spark_k8s",
        image="myregistry/spark-job:latest",
        cmds=["python", "/app/etl.py"],
        namespace="spark",
    )

Best practice: разделять ETL на атомарные Spark-задачи (по одной трансформации), выстраивать зависимости в Airflow DAG, обеспечивать идемпотентность каждого шага.


Как спроектировать пайплайн для обработки 10 ТБ данных ежедневно?

Ключевые решения:

Хранение: Iceberg/Delta Lake на S3 (Parquet + zstd, партиционирование по дате)
Чтение: только нужные партиции (predicate pushdown на дату)
Обработка:
  - shuffle.partitions: ~2000 (10 ТБ / 5 МБ per partition)
  - executor.memory: 8–16g с достаточным overhead
  - AQE enabled: динамическая корректировка партиций и skew join

Стратегия:

  1. Ingestion: Auto Loader / Kafka → Bronze (raw, без трансформации)
  2. Cleansing: Bronze → Silver (дедупликация, типы, null-handling) - partitionBy(date)
  3. Aggregation: Silver → Gold (бизнес-метрики, joins) - bucketing для часто join'ируемых ключей
  4. Оптимизация: ежедневный OPTIMIZE / rewrite_data_files для компакции мелких файлов

Мониторинг: Spark UI → Stages, проверять spill, skew, GC time.


Вы столкнулись с ошибкой ExecutorLostFailure. Каковы ваши действия?

ExecutorLostFailure = Executor неожиданно завершился. Диагностика по причине:

Причина Признаки Решение
OOM Executor Container killed, heap dump Увеличить executor.memory; проверить data skew
GC overhead Spark UI: GC Time > 10% Task Time G1GC; уменьшить memory.storageFraction
Preemption (YARN) Executor lost после долгого ожидания Повысить приоритет приложения или static allocation
Hardware failure Несколько Executor'ов на одном узле Проверить node logs; Spark перезапустит Tasks автоматически
# Посмотреть stderr Executor'а (YARN)
yarn logs -applicationId application_xxx -containerId container_yyy

# Признаки OOM в stderr:
# java.lang.OutOfMemoryError: Java heap space
# GC overhead limit exceeded

Spark автоматически перезапустит упавшие Tasks на других Executor'ах (до spark.task.maxFailures раз, default=4). Если Task падает больше 4 раз - Job завершается с ошибкой.


Опишите процесс миграции пайплайна с Pandas на PySpark

Стратегия поэтапной миграции:

  1. Аудит кода - найти Pandas-антипаттерны: iterrows, apply по строкам, to_dict, in-place mutation
  2. Переписать на встроенные функции - они выполняются в JVM без Python overhead
  3. Сложная логика@pandas_udf - получить векторизацию без полного рефактора
  4. Тестировать на subset данных в local[2]
# Pandas → PySpark: основные паттерны
# Добавить колонку
df["new"] = df["a"] + df["b"]
# → df.withColumn("new", col("a") + col("b"))

# Фильтр
df[df["age"] > 18]
# → df.filter(col("age") > 18)

# GroupBy агрегация
df.groupby("dept")["salary"].mean()
# → df.groupBy("dept").agg(avg("salary"))

# Merge
df.merge(other, on="id", how="left")
# → df.join(other, "id", "left")

# Apply по строкам (сложная логика)
df["result"] = df.apply(complex_func, axis=1)
# → @pandas_udf("double")
#    def complex_func(s: pd.Series) -> pd.Series: ...
#    df.withColumn("result", complex_func(col("value")))

Что не переносится 1-в-1: df.iloc, df.loc, библиотеки типа scipy/sklearn внутри apply - для них используют Pandas UDF с батчевой обработкой или выносят на Ray/MLflow.


Как Spark работает с Kubernetes в качестве Cluster Manager?

В K8s-режиме Driver и Executor'ы - отдельные Pod'ы в Kubernetes кластере.

spark-submit \
  --master k8s://https://k8s-api:6443 \
  --deploy-mode cluster \
  --conf spark.kubernetes.container.image=myregistry/spark:3.5 \
  --conf spark.kubernetes.namespace=spark-jobs \
  --conf spark.executor.instances=10 \
  --conf spark.kubernetes.executor.request.cores=2 \
  s3a://bucket/jobs/etl.py

Как работает:

  1. spark-submit подключается к K8s API
  2. Создаётся Driver Pod, который запускает SparkContext
  3. Driver запрашивает K8s создать Executor Pod'ы
  4. При завершении Job - все Pod'ы уничтожаются

Преимущества перед YARN:

  • Полная изоляция ресурсов (namespace + resource quotas)
  • Docker-образ = воспроизводимое окружение
  • Auto-scaling с KEDA/Karpenter
  • Нативная интеграция с CI/CD (ArgoCD, Flux)

Как реализовать Data Quality Checks внутри Spark-процесса?

Два подхода:

1. Custom assertions - простые проверки через assert/Exception:

from pyspark.sql.functions import col, count, when

def check_quality(df, name):
    total = df.count()
    assert total > 0, f"{name}: пустой DataFrame"

    null_ids = df.filter(col("user_id").isNull()).count()
    assert null_ids == 0, f"{name}: {null_ids} NULL user_id"

    # Сводная статистика по всем колонкам сразу
    null_stats = df.select([
        count(when(col(c).isNull(), c)).alias(c)
        for c in df.columns
    ]).collect()[0].asDict()

    for col_name, null_count in null_stats.items():
        if null_count > total * 0.05:  # > 5% null
            raise ValueError(f"{name}.{col_name}: {null_count} nulls ({null_count/total:.1%})")

2. Great Expectations - декларативный фреймворк:

import great_expectations as gx

context = gx.get_context()
batch = context.sources.spark_datasource.get_batch(df)

# Правила проверки
batch.expect_column_values_to_not_be_null("user_id")
batch.expect_column_values_to_be_between("age", min_value=0, max_value=150)
batch.expect_column_values_to_be_unique("order_id")

results = context.run_checkpoint("daily_etl_checkpoint")
if not results["success"]:
    raise RuntimeError(f"DQ check failed: {results}")

Как эффективно пересчитать исторические данные за 2 года (Backfilling)?

Стратегия: батчевый backfill по партициям с поддержкой idempotentности.

from datetime import date, timedelta
from dateutil.relativedelta import relativedelta

start = date(2022, 1, 1)
end   = date(2023, 12, 31)

# Batching по месяцам - меньший риск при сбое
month = start
while month <= end:
    next_month = month + relativedelta(months=1)

    df = spark.read.parquet("s3://bucket/events/") \
        .filter(col("date").between(str(month), str(next_month - timedelta(days=1))))

    result = transform(df)

    # Overwrite только конкретной партиции (idempotent)
    result.write \
        .partitionBy("date") \
        .mode("overwrite") \
        .parquet("s3://bucket/output/")

    month = next_month

Ключевые принципы:

  • Партиционировать по дате → overwrite конкретного месяца безопасен
  • Idempotent: повторный запуск не дублирует данные
  • Небольшие батчи: при сбое перезапускаем только упавший месяц
  • Iceberg Time Travel: можно запустить пайплайн на историческое состояние таблицы

Что такое Data Lineage и как его реализовать для Spark-задач?

Data Lineage - отслеживание происхождения данных: откуда пришли, через какие трансформации прошли, где сохранены.

Подходы:

1. Apache Atlas - метаданные + lineage из коробки для Spark через Spark Atlas Connector:

# При регистрации SparkSession
spark = SparkSession.builder \
    .config("spark.extraListeners", "com.hortonworks.spark.atlas.SparkAtlasEventTracker") \
    .getOrCreate()
# Atlas автоматически записывает: какие таблицы читались, какие создавались

2. OpenLineage - open-standard для lineage событий (используется в Airflow, Spark, dbt):

spark-submit \
  --conf spark.extraListeners=io.openlineage.spark.agent.OpenLineageSparkListener \
  --conf spark.openlineage.transport.url=http://marquez:5000

3. Ручное логирование - минимальный вариант:

lineage = {
    "job": "daily_etl",
    "inputs":  ["s3://bucket/events/", "s3://bucket/users/"],
    "outputs": ["s3://bucket/output/summary/"],
    "timestamp": datetime.now().isoformat()
}
spark.createDataFrame([lineage]).write.json("s3://bucket/lineage/")

Какую роль играет Spark в современной архитектуре Data Lakehouse?

Lakehouse - архитектура, объединяющая гибкость Data Lake (дешёвое объектное хранилище) с возможностями Data Warehouse (ACID, схемы, SQL).

Spark в Lakehouse - центральный compute-движок:

Spark отвечает за:

  • Ingestion: чтение из Kafka, JDBC, файловых источников
  • Transformation: ETL, cleansing, aggregation (Bronze→Silver→Gold)
  • Writes: ACID-запись в Iceberg/Delta (MERGE INTO, schema evolution)
  • Maintenance: компакция, expire snapshots, Z-ordering

Spark не единственный движок в Lakehouse - Trino/Athena используется для ad-hoc SQL, Flink для low-latency streaming. Iceberg как open table format обеспечивает единый доступ для всех движков без дублирования данных.


Как отлаживать задачи, которые зависают на 99% (Stuck Tasks)?

"Задание на 99%" - Stage завершён почти полностью, но одна-две Task работают непропорционально долго. Причины и диагностика:

Причина Признак Действие
Data Skew Одна Task обрабатывает гигантскую партицию Spark UI → Tasks → Sort by Duration; смотреть Shuffle Read Size одной Task
GC pressure GC Time > 50% Task Time Spark UI → Executors → GC Time / Task Time
Stragglers (медленный узел) Task на конкретном Executor всегда медленнее Включить Speculative Execution
External call в UDF UDF делает HTTP/DB-запрос на каждую строку Искать в коде UDF; использовать mapPartitions
Shuffle не завершён Task в статусе "Running", Shuffle Read = 0 Executor не получает данные - проверить сеть, shuffle service
# Диагностика data skew
df.groupBy("join_key").count().orderBy(col("count").desc()).show(20)

# Включить Speculative Execution для stragglers
spark.conf.set("spark.speculation",           "true")
spark.conf.set("spark.speculation.multiplier", "3")   # Task в 3× медленнее медианы
spark.conf.set("spark.speculation.quantile",   "0.9") # 90% Tasks завершились

Самая частая причина в production - data skew. Проверяйте её первой.


Инструменты для мониторинга Spark-кластера в продакшене

Инструмент Что мониторит Как подключить
Spark UI Real-time: DAG, Stages, Tasks, GC, Shuffle, Storage Встроен в Spark (порт 4040)
Spark History Server Завершённые Job'ы spark.eventLog.enabled=true + History Server
Prometheus + Grafana Метрики кластера, memory, CPU, GC JMX Exporter на Driver/Executor
SparkMeasure Programmatic сбор Stage/Task метрик StageMetrics(spark).runAndMeasure(...)
Ganglia / Datadog Node-level: CPU, RAM, network Агент на каждом узле
# Включить event log для History Server
spark.conf.set("spark.eventLog.enabled",  "true")
spark.conf.set("spark.eventLog.dir",      "s3://bucket/spark-logs/")
spark.conf.set("spark.history.fs.logDirectory", "s3://bucket/spark-logs/")

# SparkMeasure - программный сбор метрик
from sparkmeasure import StageMetrics
stage_metrics = StageMetrics(spark)
stage_metrics.begin()

result = df.groupBy("key").count().collect()

stage_metrics.end()
stage_metrics.print_report()
# → executorRunTime: 45231 ms, shuffleReadBytes: 2.3 GB, jvmGCTime: 1203 ms

Ключевые метрики для контроля в production: jvmGCTime, shuffleReadBytes, executorRunTime, количество Spill. Настроить алерты на аномалии этих метрик.


Как реализовать CI/CD для Spark-приложений?

Этапы CI pipeline:

# .github/workflows/spark-ci.yml (пример)
jobs:
  test:
    steps:
      - name: Unit tests
        run: pytest tests/ -v --tb=short

      - name: Type check
        run: mypy src/

      - name: Integration tests (local Spark)
        run: pytest tests/integration/ -v
        # SparkSession.builder.master("local[2]") - не нужен кластер

      - name: Build Docker image
        run: docker build -t myregistry/spark-job:${{ github.sha }} .

      - name: Push to registry
        run: docker push myregistry/spark-job:${{ github.sha }}

Этапы CD:

# Деплой через spark-submit или K8s Job
spark-submit \
  --master k8s://https://k8s-api:6443 \
  --conf spark.kubernetes.container.image=myregistry/spark-job:abc1234 \
  s3a://bucket/jobs/etl.py

Ключевые практики:

  • Версионировать Docker образ по git SHA - не использовать latest
  • Юнит-тесты в local[2] - быстро, не требует кластера
  • Интеграционные тесты - против реального MinIO/Kafka в docker-compose
  • Идемпотентность - каждый job безопасен при повторном запуске
  • Artifact registry - версионировать не только код, но и модели, схемы