Настройка SparkSession

Создание SparkSession через builder API: управление памятью и ресурсами, настройки shuffle, логирования, JDBC-подключений и прочие опции конфигурации.

core

🔧 Настройка 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 крайне гибким инструментом для различных сценариев работы с большими данными.