Жизненный цикл Spark-приложения: от SparkContext до Block Manager

Как Spark запускается изнутри: SparkSession, DAGScheduler, TaskScheduler, CoarseGrainedExecutorBackend, Block Manager и полный путь от Action до результата.

core

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 API
  • SQLContext - для DataFrame/SQL
  • HiveContext - для 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 и:

  1. Запрашивает у SchedulerBackend список живых Executor'ов с их местоположением
  2. Для каждой Task определяет предпочтительную локальность (PROCESS_LOCAL → ANY)
  3. Применяет Delay Scheduling: ждёт нужный Executor, если он занят
  4. Сериализует 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:

  1. Executor помечается как мёртвый
  2. Его Task'и повторно ставятся в очередь
  3. При Dynamic Allocation - заказывается новый Executor
  4. Если потеряны 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'а