Блок 5 - Практические кейсы
20 сценарных вопросов: top-N, incremental load, OOM, big joins, Delta Lake, TF-IDF и разбор реальных ситуаций.
Как найти топ-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.
Опишите процесс инкрементальной загрузки данных¶
Инкрементальная загрузка = загружаем только новые/изменённые записи, не перезаписывая всё.
Ключевые элементы:
- Водяной знак (watermark) - последнее значение
updated_atиз предыдущего запуска (хранить в Airflow variable или отдельной таблице) - Дедупликация -
left_anti joinновых против существующих по business key - Upsert -
MERGE INTOIceberg/Delta для обновления существующих + вставки новых - Идемпотентность - повторный запуск не должен дублировать данные
Как прочитать только новые файлы из директории?¶
# Подход 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 |
Практический чеклист:
spark.executor.memory+spark.executor.memoryOverhead- суммарно ≥ container sizespark.sql.shuffle.partitions- поднять пропорционально объёму данныхspark.sql.adaptive.enabled=true- включить AQE- Проверить 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
Стратегия:
- Ingestion: Auto Loader / Kafka → Bronze (raw, без трансформации)
- Cleansing: Bronze → Silver (дедупликация, типы, null-handling) - partitionBy(date)
- Aggregation: Silver → Gold (бизнес-метрики, joins) - bucketing для часто join'ируемых ключей
- Оптимизация: ежедневный
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¶
Стратегия поэтапной миграции:
- Аудит кода - найти Pandas-антипаттерны:
iterrows,applyпо строкам,to_dict, in-place mutation - Переписать на встроенные функции - они выполняются в JVM без Python overhead
- Сложная логика →
@pandas_udf- получить векторизацию без полного рефактора - Тестировать на 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
Как работает:
spark-submitподключается к K8s API- Создаётся Driver Pod, который запускает SparkContext
- Driver запрашивает K8s создать Executor Pod'ы
- При завершении 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 - версионировать не только код, но и модели, схемы