Настройка SparkSession
Создание SparkSession через builder API: управление памятью и ресурсами, настройки shuffle, логирования, JDBC-подключений и прочие опции конфигурации.
🔧 Настройка SparkSession¶
SparkSession служит точкой входа для работы со Spark и предлагает обширные опции конфигурации для управления ресурсами, оптимизации обработки и настройки логирования.
Шаг 1: Базовая настройка SparkSession¶
Начнём с создания SparkSession через builder API:
from pyspark.sql import SparkSession
# Инициализация базовой Spark-сессии
spark = SparkSession.builder \
.appName("My Spark Application") \
.getOrCreate()
print("Spark session is ready! 🚀")
Параметр appName помогает идентифицировать приложение в Spark UI.
Шаг 2: Опции конфигурации SparkSession¶
Управление памятью и ресурсами¶
Управление выделением памяти и использованием CPU:
- Память executor'а (
spark.executor.memory): устанавливает память для каждого executor'а - Память driver'а (
spark.driver.memory): выделяет память для driver'а - Ядра (
spark.executor.cores): задаёт количество CPU-ядер на executor
spark = SparkSession.builder \
.appName("Optimized App") \
.config("spark.executor.memory", "2g") \
.config("spark.driver.memory", "1g") \
.config("spark.executor.cores", "2") \
.getOrCreate()
Настройки shuffle и хранения¶
- Shuffle-партиции (
spark.sql.shuffle.partitions): настраивает количество партиций для shuffle-операций в join'ах и агрегациях - Уровень хранения (
spark.storage.level): управляет кешированием для эффективного повторного использования DataFrame
spark = SparkSession.builder \
.appName("Shuffle Optimization") \
.config("spark.sql.shuffle.partitions", "100") \
.getOrCreate()
Логирование и отладка¶
- Уровень логов (
spark.eventLog.enabled): включает логирование событий для мониторинга задач - Директория checkpoint'ов (
spark.checkpoint.dir): указывает директорию для данных checkpoint'ов
spark = SparkSession.builder \
.appName("Logging App") \
.config("spark.eventLog.enabled", "true") \
.config("spark.eventLog.dir", "/path/to/logs") \
.getOrCreate()
Настройка JDBC-подключения и базы данных¶
- JDBC URL (
spark.sql.sources.jdbc.url): задаёт URL для подключения к базе данных - Тайм-аут соединения (
spark.network.timeout): устанавливает длительность тайм-аута сети
spark = SparkSession.builder \
.appName("Database App") \
.config("spark.sql.sources.jdbc.url", "jdbc:postgresql://localhost:5432/mydb") \
.config("spark.network.timeout", "120s") \
.getOrCreate()
Прочие опции¶
- Порог broadcast-join'а (
spark.sql.autoBroadcastJoinThreshold): управляет порогом размера данных для broadcast-джоинов
spark = SparkSession.builder \
.appName("Misc Configurations") \
.config("spark.local.dir", "/path/to/temp") \
.config("spark.sql.autoBroadcastJoinThreshold", "10MB") \
.getOrCreate()
Шаг 3: Финализация настройки SparkSession¶
Вызовите .getOrCreate() для инициализации сессии с заданными настройками.
Каждая опция конфигурации адаптирует поведение Spark под конкретные потребности в ресурсах, предпочтения по логированию и форматы хранения, делая SparkSession крайне гибким инструментом для различных сценариев работы с большими данными.