SparkSession: точка входа, конфигурация и master URL
Как устроен SparkSession изнутри: SparkContext, SparkEnv, DAGScheduler, BlockManager, SharedState, SessionState. Полная цепочка запуска приложения Spark - от builder до первого executor.
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 - slaveMapOutputTracker- реестр расположений shuffle-блоков. Driver держитMapOutputTrackerMaster, Executor -MapOutputTrackerWorkerUnifiedMemoryManager- управляет разделением памяти между 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.
Полная цепочка запуска¶
Ключевые моменты последовательности:
- SparkEnv стартует первым - BlockManager и MapOutputTracker должны быть готовы до любых операций с данными
- SparkUI стартует на :4040 - если порт занят, Spark автоматически пробует :4041, :4042, ...
- HTTP FileServer - раздаёт jar/py файлы Executor-ам. Без него Executor не сможет загрузить пользовательский код
- SharedState и SessionState инициализируются после SparkContext - они зависят от уже запущенного SC
- 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() запускает сложную цепочку:
SparkConfсобирает конфигурацию из всех источников (код → spark-submit → defaults → env)SparkContextсоздаётSparkEnv: BlockManager, MapOutputTracker, UnifiedMemoryManagerDAGSchedulerиTaskSchedulerготовы принимать задачиSchedulerBackendрегистрируется в кластер-менеджере и запрашивает Executor-ыSharedStateинициализирует Catalog и CacheManagerSessionStateсоздаёт Catalyst pipeline (Parser → Analyzer → Optimizer → Planner)- Executor-ы регистрируются асинхронно - Driver готов раньше них
Понимание этой цепочки объясняет большинство runtime-проблем: почему нельзя создать два SparkContext, почему spark.stop() важен, почему конфиг не меняется после старта, и почему первый job запускается медленнее последующих.