SparkSession: точка входа, конфигурация и master URL

Как устроен SparkSession изнутри: SparkContext, SparkEnv, DAGScheduler, BlockManager, SharedState, SessionState. Полная цепочка запуска приложения Spark - от builder до первого executor.

core internals

SparkSession - единственный объект, с которого начинается любое PySpark-приложение. За одной строкой SparkSession.builder.getOrCreate() скрывается каскад из десятков компонентов. Этот урок - разбор каждого из них: зачем существует, что инициализирует и почему именно в таком порядке.


История: от SparkContext к SparkSession

До Spark 2.0 существовало несколько точек входа:

Объект Spark Назначение
SparkContext 1.x RDD API, управление кластером
SQLContext 1.x DataFrame API поверх SparkContext
HiveContext 1.x HiveQL, метастор, UDF Hive
StreamingContext 1.x DStream streaming
SparkSession 2.0+ Единая точка входа для всего

SparkSession объединил все предыдущие контексты. Внутри он по-прежнему создаёт SparkContext - ничего не исчезло, просто скрылось за унифицированным API.


SparkContext: фундамент всего

SparkContext - это ядро Spark. Он существует в одном экземпляре на JVM-процесс и управляет всем жизненным циклом приложения.

Что создаёт SparkContext при инициализации

SparkConf - иммутабельный объект конфигурации. После передачи в SparkContext изменить конфигурацию нельзя.

SparkEnv - контейнер всех runtime-сервисов. Существует в двух экземплярах: один на Driver, один на каждом Executor. Содержит:

  • BlockManager - локальное хранилище блоков данных (RDD кэш, shuffle файлы, broadcast переменные). На Driver - master, на Executor - slave
  • MapOutputTracker - реестр расположений shuffle-блоков. Driver держит MapOutputTrackerMaster, Executor - MapOutputTrackerWorker
  • UnifiedMemoryManager - управляет разделением памяти между execution (shuffle, sort, join) и storage (RDD cache). Доля задаётся через spark.memory.fraction (default 0.6)
  • SerializerManager - выбирает между Java serializer и Kryo для сериализации shuffle-данных

DAGScheduler - получает RDD граф, делит его на стадии по shuffle-границам (wide transformations), строит DAG стадий и передаёт в TaskScheduler.

TaskScheduler - принимает TaskSet от DAGScheduler, ищет доступные Executor, назначает задачи с учётом data locality (PROCESS_LOCAL → NODE_LOCAL → RACK_LOCAL → ANY).

SchedulerBackend - адаптер между TaskScheduler и конкретным кластер-менеджером (Standalone, YARN, Kubernetes). Отвечает за регистрацию Executor и отправку LaunchTask.


SharedState и SessionState

SparkSession добавляет поверх SparkContext два объекта:

SharedState - разделяется между всеми SparkSession в одном SparkContext:

  • InMemoryCatalog / HiveCatalog - реестр баз данных, таблиц, вьюшек
  • warehouse.dir - корневая директория для managed таблиц
  • CacheManager - хранит кэшированные DataFrame (.cache())

SessionState - изолирован для каждой SparkSession. Это позволяет создавать несколько сессий с разными конфигами в рамках одного SparkContext. Содержит полный Catalyst pipeline: Parser → Analyzer → Optimizer → Planner.


Полная цепочка запуска

Ключевые моменты последовательности:

  1. SparkEnv стартует первым - BlockManager и MapOutputTracker должны быть готовы до любых операций с данными
  2. SparkUI стартует на :4040 - если порт занят, Spark автоматически пробует :4041, :4042, ...
  3. HTTP FileServer - раздаёт jar/py файлы Executor-ам. Без него Executor не сможет загрузить пользовательский код
  4. SharedState и SessionState инициализируются после SparkContext - они зависят от уже запущенного SC
  5. Executor регистрируется асинхронно - getOrCreate() возвращается до появления первого Executor

Client-режим vs Cluster-режим

Параметр Client Mode Cluster Mode
Где работает Driver На клиентской машине Внутри кластера
Сетевой трафик Через клиент Внутри кластера
Закрытие терминала Убивает job Job продолжается
Логи Driver Прямо в консоль Через yarn logs / Spark History
Использование Разработка, отладка Продакшн

Builder Pattern и getOrCreate()

spark = (
    SparkSession.builder
    .master("local[*]")
    .appName("MyApp")
    .config("spark.executor.memory", "4g")
    .config("spark.sql.shuffle.partitions", "200")
    .enableHiveSupport()   # HiveCatalog + HiveQL
    .getOrCreate()
)

getOrCreate() работает по следующей логике:

Важно: getOrCreate() возвращает существующую сессию если SparkContext уже запущен. Новый SparkContext в одной JVM создать нельзя - нужно остановить предыдущий через spark.stop().


Master URL: где запускать задачи

Master URL Режим Описание
local Локальный 1 поток, Driver = Executor
local[N] Локальный N потоков
local[*] Локальный Все доступные ядра
local[N,M] Локальный N потоков, M ретраев задач
spark://host:7077 Standalone Spark Standalone кластер
spark://h1:7077,h2:7077 Standalone HA с несколькими Master
yarn YARN Hadoop YARN
k8s://https://host:port Kubernetes K8s кластер
mesos://host:5050 Mesos Apache Mesos (устаревший)
# Локально для разработки
spark = SparkSession.builder.master("local[*]").getOrCreate()

# Standalone кластер
spark = SparkSession.builder.master("spark://master:7077").getOrCreate()

# YARN - master URL не нужен, берётся из HADOOP_CONF_DIR
spark = SparkSession.builder.master("yarn").getOrCreate()

# Kubernetes
spark = SparkSession.builder \
    .master("k8s://https://k8s-api:6443") \
    .config("spark.kubernetes.container.image", "my-spark:3.5") \
    .getOrCreate()

Иерархия приоритетов конфигурации

Spark применяет настройки в следующем порядке (выше = приоритетнее):

Правило: если видите неожиданное поведение - сначала проверьте spark.sparkContext.getConf().getAll().


Критически важные параметры конфигурации

Ресурсы

Параметр Default Описание
spark.executor.memory 1g Память JVM Executor
spark.executor.memoryOverhead max(10%, 384MB) Off-heap (UDF, native libs)
spark.executor.cores 1 (YARN/K8s), все ядра (Standalone) vCPU на Executor
spark.driver.memory 1g Память Driver
spark.driver.maxResultSize 1g Макс. размер collect()

Shuffle и партиции

Параметр Default Описание
spark.sql.shuffle.partitions 200 Партиции после shuffle
spark.default.parallelism 2× CPU Партиции RDD
spark.sql.files.maxPartitionBytes 128MB Макс. размер partition при чтении

AQE (Adaptive Query Execution)

Параметр Default Описание
spark.sql.adaptive.enabled true (Spark 3.2+) Включить AQE
spark.sql.adaptive.coalescePartitions.enabled true Слияние мелких партиций
spark.sql.adaptive.skewJoin.enabled true Автоматическое лечение skew

Оптимизация

Параметр Default Описание
spark.sql.autoBroadcastJoinThreshold 10MB Порог Broadcast Join
spark.serializer JavaSerializer KryoSerializer быстрее
spark.sql.parquet.compression.codec snappy Сжатие Parquet

S3/MinIO (важно для хранилищ)

spark = SparkSession.builder \
    .config("spark.hadoop.fs.s3a.endpoint", "http://minio:9000") \
    .config("spark.hadoop.fs.s3a.access.key", "minioadmin") \
    .config("spark.hadoop.fs.s3a.secret.key", "minioadmin") \
    .config("spark.hadoop.fs.s3a.path.style.access", "true") \
    .config("spark.hadoop.fs.s3a.impl", "org.apache.hadoop.fs.s3a.S3AFileSystem") \
    .getOrCreate()

SparkSession API: основные операции

Чтение и запись данных

# Чтение
df = spark.read.parquet("s3a://bucket/data/")
df = spark.read.csv("data.csv", header=True, inferSchema=True)
df = spark.read.json("data.json")
df = spark.read.format("iceberg").load("catalog.db.table")

# Запись
df.write.mode("overwrite").parquet("s3a://bucket/output/")
df.write.mode("append").partitionBy("date").parquet("s3a://bucket/partitioned/")

SQL и временные вьюшки

# Создать временную вьюшку (живёт в рамках SparkSession)
df.createOrReplaceTempView("events")

# Глобальная вьюшка (видна в других SparkSession)
df.createOrReplaceGlobalTempView("global_events")

# SQL-запрос
result = spark.sql("""
    SELECT user_id, COUNT(*) as cnt
    FROM events
    WHERE date >= '2024-01-01'
    GROUP BY user_id
""")

# Читать из глобальной вьюшки
result2 = spark.sql("SELECT * FROM global_temp.global_events")

Catalog API

# Список баз данных
spark.catalog.listDatabases()

# Список таблиц
spark.catalog.listTables("my_database")

# Проверка существования таблицы
spark.catalog.tableExists("my_database.my_table")

# Очистка кэша
spark.catalog.clearCache()
spark.catalog.uncacheTable("events")

# Смена базы данных
spark.catalog.setCurrentDatabase("my_database")

Регистрация UDF

from pyspark.sql.functions import udf
from pyspark.sql.types import StringType

# Регистрация через SparkSession (доступно в SQL)
def clean_phone(phone):
    return phone.replace("-", "").replace(" ", "") if phone else None

spark.udf.register("clean_phone", clean_phone, StringType())

# Использование в SQL
result = spark.sql("SELECT clean_phone(phone_number) FROM customers")

# Использование в DataFrame API
clean_phone_udf = udf(clean_phone, StringType())
df.withColumn("clean_phone", clean_phone_udf("phone_number"))

Несколько SparkSession в одном приложении

# Основная сессия
spark1 = SparkSession.builder.master("local[*]").appName("App").getOrCreate()

# Вторая сессия (отдельный SessionState, общий SparkContext)
spark2 = spark1.newSession()

# Разные конфигурации на уровне сессии
spark1.conf.set("spark.sql.shuffle.partitions", "200")
spark2.conf.set("spark.sql.shuffle.partitions", "50")

# Разные temp views
spark1.createDataFrame([(1, "a")], ["id", "val"]).createTempView("t1")
spark2.sql("SHOW TABLES").show()  # t1 не видна в spark2

# Общий кэш через SharedState
df = spark1.read.parquet("data.parquet").cache()
# spark2 тоже видит кэш - CacheManager общий

Это полезно при тестировании: каждый тест создаёт newSession() с чистым каталогом temp views.


Жизненный цикл и корректная остановка

spark = SparkSession.builder.getOrCreate()

try:
    # работа с данными
    df = spark.read.parquet("data/")
    df.write.parquet("output/")
finally:
    # Корректная остановка:
    # 1. Завершает все running jobs
    # 2. Останавливает Executor-ы (через SchedulerBackend)
    # 3. Останавливает SparkUI
    # 4. Сохраняет Event Log для History Server
    spark.stop()

Если не вызвать spark.stop():

  • Executor-ы зависнут до таймаута (YARN: spark.yarn.executor.failuresValidityInterval)
  • Event Log не будет закрыт - History Server не покажет приложение
  • В Jupyter / PySpark shell это обычно нормально - процесс завершается и так

Паттерн: фабрика SparkSession для разных окружений

import os
from pyspark.sql import SparkSession


def create_spark_session(
    app_name: str,
    env: str = "local",
) -> SparkSession:
    builder = SparkSession.builder.appName(app_name)

    if env == "local":
        builder = (
            builder
            .master("local[*]")
            .config("spark.sql.shuffle.partitions", "4")
            .config("spark.driver.memory", "2g")
        )
    elif env == "dev":
        builder = (
            builder
            .master("yarn")
            .config("spark.executor.instances", "2")
            .config("spark.executor.memory", "4g")
            .config("spark.executor.cores", "2")
            .config("spark.sql.shuffle.partitions", "50")
        )
    elif env == "prod":
        builder = (
            builder
            .master("yarn")
            .config("spark.executor.instances", "20")
            .config("spark.executor.memory", "16g")
            .config("spark.executor.cores", "4")
            .config("spark.sql.shuffle.partitions", "400")
            .config("spark.sql.adaptive.enabled", "true")
            .config("spark.serializer",
                    "org.apache.spark.serializer.KryoSerializer")
        )

    return builder.getOrCreate()


# Использование
env = os.getenv("SPARK_ENV", "local")
spark = create_spark_session("etl-pipeline", env=env)

Диагностика: что проверить если Spark ведёт себя странно

# Текущая конфигурация
for key, val in spark.sparkContext.getConf().getAll():
    print(f"{key} = {val}")

# Версия Spark
print(spark.version)

# Приложение и UI
print(spark.sparkContext.applicationId)
print(spark.sparkContext.uiWebUrl)

# Статус Executor-ов
print(spark.sparkContext.statusTracker().getExecutorInfos())

# Сколько памяти выделено
print(spark.sparkContext.getConf().get("spark.executor.memory"))

Spark UI на :4040 показывает:

  • Jobs - все запущенные и завершённые job
  • Stages - DAG стадий с timing и размером shuffle
  • Storage - кэшированные RDD/DataFrame
  • Environment - все активные параметры конфигурации
  • Executors - список Executor с памятью, GC time, shuffle read/write

Итог

SparkSession.builder.getOrCreate() запускает сложную цепочку:

  1. SparkConf собирает конфигурацию из всех источников (код → spark-submit → defaults → env)
  2. SparkContext создаёт SparkEnv: BlockManager, MapOutputTracker, UnifiedMemoryManager
  3. DAGScheduler и TaskScheduler готовы принимать задачи
  4. SchedulerBackend регистрируется в кластер-менеджере и запрашивает Executor-ы
  5. SharedState инициализирует Catalog и CacheManager
  6. SessionState создаёт Catalyst pipeline (Parser → Analyzer → Optimizer → Planner)
  7. Executor-ы регистрируются асинхронно - Driver готов раньше них

Понимание этой цепочки объясняет большинство runtime-проблем: почему нельзя создать два SparkContext, почему spark.stop() важен, почему конфиг не меняется после старта, и почему первый job запускается медленнее последующих.