Жизненный цикл Spark-приложения: от SparkContext до Block Manager
Как Spark запускается изнутри: SparkSession, DAGScheduler, TaskScheduler, CoarseGrainedExecutorBackend, Block Manager и полный путь от Action до результата.
Spark - распределённое Java-приложение¶
Когда вы запускаете spark-submit myjob.py, стартует обычный Java-процесс. Вся магия - это оркестрация JVM-процессов по кластеру через RPC. Понимание этой оркестрации объясняет большинство ошибок: долгий холодный старт, FetchFailed, heartbeat timeout, No alive nodes.
Жизненный цикл делится на три крупных фазы:
Фаза 1: Инициализация¶
Entry Point: эволюция от SparkContext к SparkSession¶
До Spark 2.0 было несколько entry points:
SparkContext- основной, для RDD APISQLContext- для DataFrame/SQLHiveContext- для Hive metastore
Начиная с Spark 2.0, всё объединено в SparkSession:
spark = SparkSession.builder \
.appName("MyJob") \
.master("yarn") \
.config("spark.executor.memory", "4g") \
.getOrCreate()
# SparkSession → SparkContext (доступен как spark.sparkContext)
sc = spark.sparkContext
SparkSession.builder.getOrCreate() - идемпотентен: если сессия уже существует в JVM, возвращает её. Полезно для библиотек, которые не знают, был ли уже создан SparkContext.
Что происходит внутри Driver JVM при старте¶
DAGScheduler - строит и управляет DAG стадий. При каждом Action разбивает логический план на Stage → TaskSet.
TaskScheduler - получает TaskSet от DAGScheduler и назначает Task'и на конкретные Executor'ы с учётом Data Locality.
SchedulerBackend - RPC-клиент к Cluster Manager (YARN ResourceManager, Kubernetes API, Standalone Master). Через него происходит запрос ресурсов.
BlockManagerMaster - центральный реестр: "блок B001 находится на Executor 3". Executor'ы спрашивают его, чтобы найти чужие shuffle-блоки.
MapOutputTracker - реестр shuffle output location. После завершения Map-фазы каждый Executor сообщает Driver'у, куда записал свои shuffle-блоки.
SparkUI - HTTP-сервер на порту 4040. Работает всё время жизни приложения.
Запрос ресурсов у Cluster Manager¶
| Cluster Manager | Как запускается Executor |
|---|---|
| YARN | ApplicationMaster запрашивает Container у ResourceManager; NodeManager стартует CoarseGrainedExecutorBackend |
| Kubernetes | Driver создаёт Pod через K8s API; Pod содержит образ Spark с CoarseGrainedExecutorBackend |
| Standalone | Driver регистрируется у Master; Master командует Worker-демону стартовать Executor |
CoarseGrainedExecutorBackend - org.apache.spark.executor.CoarseGrainedExecutorBackend. Это JVM-процесс, который:
- Живёт всё время приложения (отсюда "Coarse-grained" - противоположность fine-grained, где процесс на каждую Task)
- Поднимает внутри себя
Executorс пулом Task-потоков - Регистрируется у Driver'а по Netty RPC
- Принимает Task'и, выполняет, отчитывается через heartbeat
Структура процессов на кластере после инициализации¶
Фаза 2: Выполнение - от Action до результата¶
Путь от Python-кода до физического плана¶
Иерархия: Job → Stage → Task¶
Job - всё вычисление, запущенное одним Action. Один .write() = один Job. df.cache(); df.count() = два Job'а.
Stage - группа Task'ей, которые можно выполнить без сетевой передачи данных между собой. Граница Stage = граница Shuffle (Wide Transformation).
Task - минимальная единица: один поток на одном Executor'е обрабатывает одну партицию. Число Task'ей = число партиций.
DAGScheduler: построение графа стадий¶
DAGScheduler смотрит на Physical Plan и находит ShuffleExchange узлы - это границы Stage. Всё до ShuffleExchange - один Stage, всё после - следующий.
DAGScheduler также управляет повторными запусками: если Stage падает (Executor умер), DAGScheduler перезапускает только упавшие Task'и (или весь Stage при потере shuffle-данных).
TaskScheduler: назначение на Executor'ы¶
TaskScheduler получает TaskSet от DAGScheduler и:
- Запрашивает у SchedulerBackend список живых Executor'ов с их местоположением
- Для каждой Task определяет предпочтительную локальность (PROCESS_LOCAL → ANY)
- Применяет Delay Scheduling: ждёт нужный Executor, если он занят
- Сериализует Task-замыкание (closure) и отправляет на выбранный Executor
Speculative Execution - особая стратегия: если Task работает значительно дольше медианы своего Stage (>75% задач уже завершены, а она всё ещё идёт), TaskScheduler запускает дубликат Task'и на другом Executor'е. Кто первый завершит - тот и победил.
# Включить speculative execution (по умолчанию: false)
spark.conf.set("spark.speculation", "true")
spark.conf.set("spark.speculation.multiplier", "1.5") # задача медленнее в 1.5× медианы
spark.conf.set("spark.speculation.quantile", "0.75") # ждём пока 75% задач завершатся
Выполнение Task на Executor¶
Heartbeat: каждый Executor отправляет Driver'у сердцебиение каждые spark.executor.heartbeatInterval (default: 10 s). Если Driver не получает heartbeat дольше spark.network.timeout (default: 120 s) - считает Executor мёртвым и перезапускает его Task'и.
Block Manager: сердце хранения данных¶
Архитектура Block Manager¶
Block Manager - сервис на каждом Executor'е (и на Driver'е для broadcast). Он управляет всеми блоками данных:
Типы блоков, которыми управляет Block Manager¶
| Тип блока | Префикс | Где хранится | Кто создаёт |
|---|---|---|---|
| RDD-партиция (кэш) | rdd_N_P |
Memory + Disk | df.cache() / df.persist() |
| Shuffle Map Output | shuffle_S_M_R |
Disk | Map-фаза (shuffle write) |
| Broadcast | broadcast_N_piece0 |
Memory | spark.sparkContext.broadcast() |
| Temporary shuffle | temp_shuffle_* |
Disk | Spill во время sort |
| Stream chunk | input_N_X |
Memory | Streaming receiver |
Жизненный цикл кэш-блока¶
Shuffle lifecycle через Block Manager¶
ExternalShuffleService: shuffle без привязки к Executor'у¶
По умолчанию shuffle-файлы хранятся в JVM-процессе Executor'а. Если Executor упал (OOM, nodemanager убил контейнер) - его shuffle-данные теряются, и Stage придётся перезапускать с нуля.
ExternalShuffleService (ESS) - отдельный daemon-процесс на каждом Worker-узле:
- Хранит shuffle-файлы независимо от Executor JVM
- Executor может быть убит (при Dynamic Allocation) - файлы остаются
- Обязателен при
spark.dynamicAllocation.enabled = true
spark.conf.set("spark.shuffle.service.enabled", "true") # включить ESS
spark.conf.set("spark.dynamicAllocation.enabled", "true") # требует ESS
Коммуникация внутри Spark¶
Netty RPC: основной транспорт¶
Spark использует Netty для всех внутренних коммуникаций. Каждый компонент - это RpcEndpoint с уникальным адресом:
Heartbeat и обнаружение сбоев¶
# Ключевые таймауты
spark.conf.set("spark.executor.heartbeatInterval", "10s") # как часто Executor пингует Driver
spark.conf.set("spark.network.timeout", "120s") # таймаут соединения (должен > heartbeatInterval)
spark.conf.set("spark.task.maxFailures", "4") # число попыток Task перед отказом Stage
Если Driver не получает heartbeat дольше network.timeout:
- Executor помечается как мёртвый
- Его Task'и повторно ставятся в очередь
- При Dynamic Allocation - заказывается новый Executor
- Если потеряны shuffle-данные - перезапускается Map Stage целиком
Fault Tolerance: как Spark восстанавливается¶
Lineage (цепочка RDD-преобразований) - основа fault tolerance. Spark не копирует данные для резервирования. Вместо этого он знает, как пересчитать любой блок данных, повторив преобразования над исходными данными. Поэтому восстановление после сбоя Task = просто повторить вычисление.
Dynamic Allocation: масштабирование в runtime¶
Dynamic Allocation позволяет добавлять и удалять Executor'ы по мере необходимости:
spark.conf.set("spark.dynamicAllocation.enabled", "true")
spark.conf.set("spark.dynamicAllocation.minExecutors", "2")
spark.conf.set("spark.dynamicAllocation.maxExecutors", "50")
spark.conf.set("spark.dynamicAllocation.initialExecutors", "5")
spark.conf.set("spark.dynamicAllocation.executorIdleTimeout", "60s")
spark.conf.set("spark.dynamicAllocation.schedulerBacklogTimeout", "1s")
Влияние на локальность: новые Executor'ы не знают о расположении данных → уровень ANY до первого кэша.
Фаза 3: Завершение¶
Что происходит при spark.stop()¶
Shuffle-файлы удаляются автоматически при нормальном завершении. При аварийном падении Driver'а (kill -9) временные файлы в spark.local.dir могут остаться - их нужно чистить вручную или через cron.
# Явный вызов - лучшая практика в долгоживущих приложениях
try:
spark.sql("SELECT ...").write.parquet("/output/")
finally:
spark.stop() # гарантированное завершение и cleanup
Практика: трейсинг жизненного цикла через Spark UI и логи¶
1. Наблюдение за инициализацией¶
import logging
logging.basicConfig(level=logging.INFO)
spark = SparkSession.builder \
.appName("LifecycleDemo") \
.config("spark.executor.instances", "3") \
.config("spark.executor.cores", "2") \
.config("spark.executor.memory", "2g") \
.getOrCreate()
# После создания SparkSession:
print(f"App ID: {spark.sparkContext.applicationId}")
print(f"UI URL: {spark.sparkContext.uiWebUrl}")
print(f"Executors: {len(spark.sparkContext._jsc.sc().statusTracker().getExecutorInfos())}")
В логах Driver'а ищем:
INFO SparkContext: Submitted application: LifecycleDemo
INFO StandaloneSchedulerBackend: Registered executor...
INFO BlockManagerMaster: Registering block manager...
2. Наблюдение за Job → Stage → Task¶
df = spark.range(10_000_000) \
.selectExpr("id % 100 AS region", "id * 1.5 AS value") \
.groupBy("region") \
.agg({"value": "sum"})
# Перед Action открыть Spark UI: http://localhost:4040
# SQL tab → посмотреть план
df.explain(mode="extended")
# После Action
df.collect()
# Jobs tab: один Job
# Stages tab: Stage 0 (Map) + Stage 1 (Reduce)
# Tasks tab: Locality Level, Duration, Shuffle Write/Read
3. Наблюдение за Block Manager через Spark UI¶
large_df = spark.read.parquet("/data/events/")
large_df.cache()
large_df.count() # первый Action - заполнение кэша
# Storage tab: смотреть закэшированные RDD
# - Partitions Cached
# - Size in Memory
# - Size on Disk (если spill)
# - Locality (PROCESS_LOCAL / NODE_LOCAL)
4. Программный доступ к метрикам¶
# Статус Executor'ов
exec_infos = spark.sparkContext._jsc.sc().statusTracker().getExecutorInfos()
for info in exec_infos:
print(f"Host: {info.host()}, Cores: {info.totalCores()}")
# Активные стадии
active_stages = spark.sparkContext._jsc.sc().statusTracker().getActiveStageIds()
for stage_id in active_stages:
stage_info = spark.sparkContext._jsc.sc().statusTracker().getStageInfo(stage_id)
if stage_info.isDefined():
s = stage_info.get()
print(f"Stage {stage_id}: {s.numActiveTasks()} active, {s.numCompletedTasks()} done")
Типичные проблемы и их причины¶
| Симптом | Компонент | Причина | Что делать |
|---|---|---|---|
No alive nodes при старте |
SchedulerBackend | Executor'ы не зарегистрировались вовремя | Увеличить spark.network.timeout, проверить YARN/K8s ресурсы |
FetchFailed при shuffle read |
Block Manager | Executor с shuffle-файлами упал | Включить ESS, увеличить spark.task.maxFailures |
heartbeat timeout |
CoarseGrainedExecutorBackend | GC STW или перегрузка Executor'а | G1GC tuning, уменьшить heap, проверить OOM |
| Долгий холодный старт | Cluster Manager | Медленное выделение контейнеров/Pod'ов | Prewarmed pool, spark.executor.instances static allocation |
Could not find block |
BlockManagerMaster | Кэш был вытеснен (eviction) | Увеличить spark.executor.memory или spark.memory.fraction |
| Stage перезапускается бесконечно | DAGScheduler | Потеря shuffle-данных при каждом retry | Включить ESS, проверить стабильность Executor'ов |
Резюме: полная карта жизненного цикла¶
Ключевые компоненты и их роли:
- SparkContext / SparkSession - точка входа, создаёт все сервисы Driver'а
- DAGScheduler - переводит физический план в граф Stage'ей, управляет повторными запусками
- TaskScheduler - назначает Task'и на Executor'ы с учётом локальности, запускает speculative execution
- CoarseGrainedExecutorBackend - JVM-процесс на Worker-узле, содержит Executor и Block Manager
- BlockManagerMaster - реестр всех блоков кластера; без него невозможен shuffle read и кэш
- MapOutputTracker - знает, куда записал shuffle-выходы каждый Map Task
- ExternalShuffleService - отвязывает shuffle-файлы от жизни Executor'а