KryoSerializer: регистрация классов и измеримый выигрыш

Полный разбор KryoSerializer в Spark: проблема JavaSerializer и раздутые метаданные, физика сериализации в Shuffle/Broadcast/Cache, включение Kryo, обязательная регистрация классов и spark.kryo.registrationRequired, KryoRegistrator, настройка буферов, измерение выигрыша в Spark UI, область применения в PySpark и Tungsten контекст.

optimization

1. Проблема дефолтной сериализации: JavaSerializer и его цена

Сериализация — это преобразование объектов из памяти JVM в бинарный поток байт для передачи по сети или записи на диск. В Spark сериализация происходит постоянно: при каждом Shuffle, Broadcast, RDD Cache и Checkpoint.

Когда Spark сериализует данные

Схема показывает критически важное разграничение: Kryo влияет на передачу Java-объектов (Shuffle, Broadcast, RDD Cache, Tasks), но не влияет на DataFrame API и работу с Parquet/ORC/Delta — там используется внутренний формат Tungsten/UnsafeRow, который вообще не использует Java-сериализацию.

Что записывает JavaSerializer в бинарный поток

JavaSerializer (стандартный java.io.ObjectOutputStream) при сериализации каждого объекта записывает:

  • Полное имя класса включая пакет: org.apache.spark.mllib.regression.LabeledPoint — это 44 байта только для имени!
  • Версию класса (serialVersionUID)
  • Метаданные всех полей: имена, типы
  • Хэдеры объектов: заголовки для каждого вложенного объекта

Для простого объекта LabeledPoint(label=1.0, features=[0.1, 0.2, 0.3]) размер в JavaSerializer может быть 300-400 байт. В Kryo — 25-30 байт. Разница 10-15x в размере.

# Демонстрация разницы в размерах (псевдокод для понимания)
import sys

# Java Serialization overhead:
# "org.apache.spark.mllib.regression.LabeledPoint" = 47 bytes
# + serialVersionUID = 8 bytes
# + field names: "label" = 5, "features" = 8, "indices" = 7...
# + values: 1 double + array_header + N doubles
# ИТОГО для LabeledPoint(1.0, [0.1, 0.2, 0.3]): ~350 bytes!

# Kryo с регистрацией:
# class_id: 2 bytes (зарегистрированный числовой ID)
# + values: 8 bytes (double) + 12 bytes (3 doubles array)
# ИТОГО: ~22 bytes! В 15 раз меньше!

Практический эффект на Shuffle

Представьте Stage с 10 000 000 объектов LabeledPoint в Shuffle:

  • JavaSerializer: 10M × 350 bytes = 3.5 GB Shuffle Write
  • Kryo с регистрацией: 10M × 22 bytes = 220 MB Shuffle Write

Разница: в 16 раз меньше сетевого трафика, в 16 раз меньше диска, значительно меньше времени на сериализацию.


2. Физика работы KryoSerializer

Kryo — библиотека сериализации для JVM, созданная специально для высокопроизводительных сред. Её автор — Натан Свит (Nathan Sweet), сейчас библиотека используется в Google, Twitter, Uber и в самом Apache Spark.

Принцип компактного кодирования

Схема показывает три уровня компактности. Разница между Java (200 байт) и Kryo с регистрацией (5 байт) — 40x! Это достигается двумя ключевыми приёмами:

1. Числовые идентификаторы классов. Вместо строки "org.example.Point" (17 байт) записывается число 10 (1-2 байта). Это возможно только при предварительной регистрации.

2. Variable-length encoding чисел. Числа 0-127 кодируются 1 байтом вместо 8 (double) или 4 (int). Для небольших значений это огромная экономия.

Место Kryo в архитектуре Spark

Ключевое понимание: В современном Spark большинство данных обрабатывается через Tungsten без каких-либо Java-объектов — это бинарные колоночные данные в памяти. Kryo влияет на то, что ещё остаётся в виде Java-объектов: Shuffle, Broadcast, RDD Cache, Task closures.


3. Включение KryoSerializer

Конфигурация при создании SparkSession

from pyspark.sql import SparkSession

# ── Способ 1: В SparkSession.builder ─────────────────────────────────
spark = SparkSession.builder \
    .appName("kryo-optimized-job") \

    # Главный параметр: заменяем сериализатор
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \

    # КРИТИЧНО: явная регистрация классов (разбираем в следующем разделе)
    # Без этого Kryo работает, но не максимально эффективно
    .config("spark.kryo.registrationRequired", "false") \  # дефолт, без ошибок

    # Буферы Kryo (разбираем в разделе 7)
    .config("spark.kryoserializer.buffer", "128k") \       # начальный буфер
    .config("spark.kryoserializer.buffer.max", "128m") \   # максимальный буфер

    .getOrCreate()

# ── Способ 2: Через spark-submit ────────────────────────────────────
# spark-submit \
#   --conf spark.serializer=org.apache.spark.serializer.KryoSerializer \
#   --conf spark.kryoserializer.buffer.max=128m \
#   your_job.py

# ── Способ 3: В spark-defaults.conf (кластерная настройка) ───────────
# В $SPARK_HOME/conf/spark-defaults.conf:
# spark.serializer   org.apache.spark.serializer.KryoSerializer
# spark.kryoserializer.buffer.max   128m

Проверка что Kryo активирован

# Проверяем активную конфигурацию
serializer = spark.conf.get("spark.serializer")
print(f"Активный сериализатор: {serializer}")
# org.apache.spark.serializer.KryoSerializer

# Также проверяем через SparkContext
sc = spark.sparkContext
print(f"Kryo активен: {sc._conf.get('spark.serializer', 'default')}")

4. Критическая важность регистрации классов

Это самый важный раздел урока. Просто включить Kryo без регистрации классов — значит получить ~50% выигрыш вместо возможных 100%.

Что происходит без регистрации

Kryo без регистрации:
→ Встречает объект типа com.company.model.Transaction
→ Записывает: длина_имени + "com.company.model.Transaction" + данные
→ Экономия vs Java: ~30-50% (только за счёт variable-length encoding)
→ Строка "com.company.model.Transaction" = 34 байта в каждую запись!

Kryo с регистрацией:
→ Transaction зарегистрирован как ID=42
→ Записывает: 42 (1 байт) + данные
→ Экономия vs Java: ~70-90%
→ Вместо 34 байт имени класса — 1 байт числа!

Параметр spark.kryo.registrationRequired

spark = SparkSession.builder \
    # РЕЖИМ РАЗРАБОТКИ: Kryo брасает исключение если класс не зарегистрирован
    # Используйте это чтобы найти ВСЕ незарегистрированные классы
    .config("spark.kryo.registrationRequired", "true") \

    .getOrCreate()

# При незарегистрированном классе:
# com.esotericsoftware.kryo.KryoException:
#   Class is not registered: com.company.model.Transaction
# Note: To register this class use: kryo.register(com.company.model.Transaction.class);

# РЕЖИМ PRODUCTION (можно оставить false для безопасности):
# .config("spark.kryo.registrationRequired", "false") \
# Незарегистрированные классы сериализуются менее эффективно,
# но не вызывают ошибок

Рекомендованный workflow:

  1. В dev/test среде включить registrationRequired=true
  2. Запустить пайплайн и собрать все KryoException — это список классов для регистрации
  3. Зарегистрировать все найденные классы
  4. В prod оставить registrationRequired=false как страховку (на случай если появятся новые классы)
  5. При крупных новых разработках периодически повторять шаги 1-3

Способы регистрации классов в Spark

Способ 1: spark.kryo.classesToRegister (самый простой)

spark = SparkSession.builder \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \

    # Список классов через запятую
    # Каждый класс получает автоматический числовой ID
    .config("spark.kryo.classesToRegister",
            "com.company.model.Transaction,"
            "com.company.model.Customer,"
            "org.apache.spark.mllib.regression.LabeledPoint,"
            "scala.collection.mutable.ArrayBuffer") \

    .getOrCreate()

Способ 2: KryoRegistrator (для сложных случаев)

KryoRegistrator позволяет написать кастомные сериализаторы для специфических классов. Это нужно когда:

  • Класс не имеет no-arg constructor (нужен для обычной Kryo десериализации)
  • Класс из сторонней библиотеки и его нельзя модифицировать
  • Нужна специальная логика сериализации (например, игнорировать некоторые поля)
// KryoRegistrator.scala
// ЭТОТ КОД НА SCALA/JAVA — компилируется в JAR и передаётся Spark'у
import com.esotericsoftware.kryo.Kryo
import org.apache.spark.serializer.KryoRegistrator

class MySparkKryoRegistrator extends KryoRegistrator {
    override def registerClasses(kryo: Kryo): Unit = {
        // Простая регистрация: Kryo использует FieldSerializer по умолчанию
        kryo.register(classOf[com.company.model.Transaction])
        kryo.register(classOf[com.company.model.Customer])

        // Регистрация с кастомным сериализатором:
        // Когда Transaction имеет поле которое не нужно сериализовать
        kryo.register(classOf[com.company.model.BigTransaction],
                      new BigTransactionSerializer())

        // Массивы и коллекции тоже нужно регистрировать!
        kryo.register(classOf[Array[com.company.model.Transaction]])
        kryo.register(classOf[Array[Double]])

        // MLlib классы
        kryo.register(classOf[org.apache.spark.mllib.regression.LabeledPoint])
        kryo.register(classOf[org.apache.spark.mllib.linalg.DenseVector])
        kryo.register(classOf[org.apache.spark.mllib.linalg.SparseVector])
    }
}

Подключение KryoRegistrator в PySpark:

# В PySpark: регистратор написан на Java/Scala и передаётся через конфиг
spark = SparkSession.builder \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \

    # Путь к JVM-классу регистратора (должен быть в classpath через --jars)
    .config("spark.kryo.registrator", "com.company.spark.MySparkKryoRegistrator") \

    .getOrCreate()

# Альтернатива для pure PySpark без Scala/Java классов:
# Используем spark.kryo.classesToRegister для стандартных JVM-классов
spark = SparkSession.builder \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.kryo.classesToRegister",
            "org.apache.spark.mllib.regression.LabeledPoint,"
            "org.apache.spark.mllib.linalg.DenseVector,"
            "org.apache.spark.mllib.linalg.SparseVector") \
    .getOrCreate()

5. Режим spark.kryo.registrationRequired и отладка

Диагностика незарегистрированных классов

# Включаем strict режим для поиска всех незарегистрированных классов
spark_strict = SparkSession.builder \
    .config("spark.serializer", "org.apache.spark.serializer.KryoSerializer") \
    .config("spark.kryo.registrationRequired", "true") \  # строгий режим
    .getOrCreate()

# Запускаем наш пайплайн
try:
    result = run_pipeline(spark_strict)  # наш основной код
    result.count()
except Exception as e:
    if "KryoException" in str(e) or "is not registered" in str(e):
        print(f"Незарегистрированный класс обнаружен:")
        print(f"  {e}")
        # Вывод: Class is not registered: com.company.model.Transaction
        # Добавляем этот класс в список регистрации и повторяем

Автоматический сбор незарегистрированных классов

import re


def find_unregistered_classes(spark, pipeline_func) -> set[str]:
    """
    Запускает пайплайн с strict Kryo и собирает все незарегистрированные классы.
    Полезно при первичной настройке KryoSerializer.

    ВНИМАНИЕ: pipeline_func должен запускать все характерные операции
    пайплайна чтобы обнаружить все классы.
    """
    spark.conf.set("spark.kryo.registrationRequired", "true")
    unregistered = set()

    # Итерируемся пока не соберём все классы
    max_iterations = 20
    for i in range(max_iterations):
        try:
            pipeline_func()
            print(f"Итерация {i+1}: все классы зарегистрированы!")
            break
        except Exception as e:
            error_str = str(e)
            # Парсим имя незарегистрированного класса из ошибки
            match = re.search(r'Class is not registered: ([^\n]+)', error_str)
            if match:
                class_name = match.group(1).strip()
                if class_name not in unregistered:
                    unregistered.add(class_name)
                    print(f"Итерация {i+1}: найден класс: {class_name}")

                    # Регистрируем и пробуем снова
                    current_classes = spark.conf.get(
                        "spark.kryo.classesToRegister", ""
                    )
                    new_classes = f"{current_classes},{class_name}".lstrip(",")
                    spark.conf.set("spark.kryo.classesToRegister", new_classes)
            else:
                print(f"Неожиданная ошибка: {error_str}")
                break

    return unregistered

6. Kryo в PySpark: особенности и ограничения

Это раздел, который часто упускают из виду. PySpark — это два слоя: JVM (Spark Core) и Python (пользовательский код). Kryo влияет только на JVM-слой.

Архитектура PySpark и сериализация

Что это означает на практике:

Kryo в PySpark даёт выигрыш только для:

  • Shuffle JVM-уровня (до 2-3x меньше Shuffle объём)
  • Broadcast JVM-объектов (lookup таблицы из DataFrame)
  • JVM Task Descriptions
  • Spark MLlib (работает на JVM, активно использует сериализацию)

Kryo не влияет на:

  • Python UDF (они используют pickle/cloudpickle)
  • Pandas UDF (используют Apache Arrow)
  • Python objects в mapPartitions
  • PySpark-специфичные операции через Py4J
# Проверим: когда Kryo реально помогает в PySpark

# СЛУЧАЙ 1: Broadcast с DataFrame-based lookup — Kryo ПОМОГАЕТ
# JVM объект (HashedRelation) передаётся через Kryo
lookup_df = spark.table("dims.products")  # DataFrame → JVM объект
from pyspark.sql import functions as F
result = fact_df.join(F.broadcast(lookup_df), "product_id")
# Kryo сериализует HashedRelation при broadcast → меньший объём

# СЛУЧАЙ 2: Python UDF — Kryo НЕ ПОМОГАЕТ
@F.udf(returnType=StringType())
def classify_amount(amount: float) -> str:
    # Эта функция сериализуется через cloudpickle, не через Kryo
    if amount > 1000:
        return "high"
    return "low"
result = df.withColumn("category", classify_amount(F.col("amount")))

# СЛУЧАЙ 3: Spark MLlib — Kryo ПОМОГАЕТ СИЛЬНО
from pyspark.mllib.regression import LabeledPoint
from pyspark.mllib.tree import RandomForest

# MLlib работает с JVM объектами (LabeledPoint) через RDD API
rdd = sc.parallelize([LabeledPoint(1.0, [0.1, 0.2, 0.3])])
# Каждый LabeledPoint сериализуется при Shuffle — Kryo даёт большой выигрыш

# СЛУЧАЙ 4: Spark GraphX (Scala API через PySpark GraphFrames) — Kryo ПОМОГАЕТ
# GraphX передаёт вершины и рёбра как Java объекты через Shuffle

7. Настройка размеров буфера Kryo

Ошибка Buffer Overflow

# Типичная ошибка при слишком маленьком буфере:
com.esotericsoftware.kryo.KryoException:
  Buffer overflow. Available: 0, required: 1024
    at com.esotericsoftware.kryo.io.Output.require(Output.java:138)
    at com.esotericsoftware.kryo.io.Output.writeBytes(Output.java:220)
    ...
# Причина: один объект > spark.kryoserializer.buffer.max

Параметры буфера и когда их увеличивать

spark = SparkSession.builder \
    # Начальный размер буфера Kryo
    # Kryo автоматически расширяет буфер при необходимости
    # (до buffer.max) — это стартовый размер
    # Дефолт: 64 KB. Увеличьте до 128-512 KB если объекты в среднем больше
    .config("spark.kryoserializer.buffer", "256k") \

    # МАКСИМАЛЬНЫЙ размер буфера Kryo
    # Один объект НЕ МОЖЕТ быть больше этого значения!
    # Дефолт: 64 MB. Для ML с большими моделями/эмбеддингами:
    .config("spark.kryoserializer.buffer.max", "512m") \

    .getOrCreate()

# Когда увеличивать buffer.max:
# 1. ML модели (weights vectors): RandomForest с 100 деревьями = 500+ MB
# 2. Эмбеддинги (embedding vectors): 10K dim vector × 4 bytes × 1M строк
# 3. Большие broadcast переменные (lookup dictionaries с миллионами записей)
# 4. Graph algorithms (вершины с большими payload'ами)

# Расчёт нужного buffer.max:
# max_object_size_bytes = max(sizeof(all_objects_to_serialize))
# buffer.max = max_object_size_bytes × 1.5 (запас 50%)

def estimate_buffer_max(df_sample, n_rows=1000) -> int:
    """
    Оценивает нужный buffer.max через сэмплирование.
    Сериализует N объектов и берёт максимальный размер.
    """
    import pickle

    rows = df_sample.limit(n_rows).collect()
    max_size = max(len(pickle.dumps(row)) for row in rows)
    recommended = int(max_size * 2.0)  # 2x запас для Kryo overhead

    return recommended

8. Измерение выигрыша в Spark UI

Правильное измерение эффекта Kryo — это не просто «было 5 минут, стало 4 минуты». Нужно смотреть конкретные метрики в Spark UI.

Что измерять в Spark UI

Вкладка Stages → конкретный Stage:

МЕТРИКИ ШАФФЛА (Shuffle Metrics):
  Shuffle Write Size: 3.2 GB  ← С JavaSerializer
  Shuffle Write Size: 0.4 GB  ← С Kryo + Registration (8x меньше!)

  Shuffle Read Size: 3.1 GB   ← С JavaSerializer
  Shuffle Read Size: 0.4 GB   ← С Kryo

МЕТРИКИ SPILL:
  Spill (Memory): 4.5 GB     ← С JavaSerializer (объекты большие, не помещаются)
  Spill (Memory): 0 GB       ← С Kryo (объекты меньше, помещаются в Execution Memory!)

МЕТРИКИ ЗАДАЧ:
  Task Deserialization Time: 45s avg  ← С JavaSerializer
  Task Deserialization Time: 8s avg   ← С Kryo
import requests
import json


def compare_serializer_impact(
    history_server_url: str,
    app_id_java: str,
    app_id_kryo: str,
    stage_id: int = 0,
) -> dict:
    """
    Сравнивает метрики Shuffle для двух запусков:
    одного с JavaSerializer и одного с KryoSerializer.
    """

    def get_stage_metrics(app_id: str, stage: int) -> dict:
        url = (f"{history_server_url}/api/v1/applications/{app_id}"
               f"/stages/{stage}/0/taskList")
        tasks = requests.get(url, timeout=10).json()

        total_shuffle_write = sum(
            t.get("taskMetrics", {}).get("shuffleWriteMetrics", {})
             .get("bytesWritten", 0) for t in tasks
        )
        total_shuffle_read = sum(
            t.get("taskMetrics", {}).get("shuffleReadMetrics", {})
             .get("totalBytesRead", 0) for t in tasks
        )
        total_spill_memory = sum(
            t.get("taskMetrics", {}).get("memoryBytesSpilled", 0) for t in tasks
        )
        total_deser_time = sum(
            t.get("taskMetrics", {}).get("executorDeserializeTime", 0) for t in tasks
        )

        return {
            "shuffle_write_gb": total_shuffle_write / 1024**3,
            "shuffle_read_gb": total_shuffle_read / 1024**3,
            "spill_memory_gb": total_spill_memory / 1024**3,
            "deser_time_ms": total_deser_time,
        }

    java_metrics = get_stage_metrics(app_id_java, stage_id)
    kryo_metrics = get_stage_metrics(app_id_kryo, stage_id)

    def pct_improvement(before, after):
        if before == 0:
            return 0
        return (before - after) / before * 100

    return {
        "shuffle_write_reduction_pct":
            pct_improvement(java_metrics["shuffle_write_gb"],
                           kryo_metrics["shuffle_write_gb"]),
        "shuffle_read_reduction_pct":
            pct_improvement(java_metrics["shuffle_read_gb"],
                           kryo_metrics["shuffle_read_gb"]),
        "spill_reduction_pct":
            pct_improvement(java_metrics["spill_memory_gb"],
                           kryo_metrics["spill_memory_gb"]),
        "deser_time_reduction_pct":
            pct_improvement(java_metrics["deser_time_ms"],
                           kryo_metrics["deser_time_ms"]),
        "java_metrics": java_metrics,
        "kryo_metrics": kryo_metrics,
    }

9. Kryo и RDD Cache: драматический эффект

Один из самых заметных эффектов Kryo — улучшение эффективности RDD Cache с StorageLevel.MEMORY_ONLY_SER.

Почему MEMORY_ONLY_SER + Kryo = мощная комбинация

from pyspark import StorageLevel

# Вариант А: MEMORY_ONLY (дефолт для .cache())
# Каждый объект хранится как "живой" Java-объект в Heap
# → Много места в RAM (Java object overhead)
# → GC давление: JVM должен отслеживать все объекты
# → Full GC на большом Heap может занимать 30-90 секунд!
rdd.persist(StorageLevel.MEMORY_ONLY)

# Вариант Б: MEMORY_ONLY_SER без Kryo
# Объекты сериализуются через JavaSerializer → byte[]
# → Байтовые массивы в Heap (меньше overhead vs живых объектов)
# → Но JavaSerializer неэффективен: раздутые byte[]
# → GC лучше (byte[] проще для GC чем граф объектов)
rdd.persist(StorageLevel.MEMORY_ONLY_SER)

# Вариант В: MEMORY_ONLY_SER + KryoSerializer ← ЛУЧШИЙ ВЫБОР
# byte[] компактны благодаря Kryo
# → В 5-10x меньше памяти чем MEMORY_ONLY
# → Практически нет GC давления (byte[] = simple memory blocks)
# → Помещается больше данных в RAM
# → Меньше Spill на диск
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
rdd.persist(StorageLevel.MEMORY_ONLY_SER)

Практическое сравнение для RDD из 10M LabeledPoint объектов:

Стратегия RAM на объект Суммарный RAM GC паузы
MEMORY_ONLY (Java objects) ~350 bytes 3.5 GB Частые, длинные
MEMORY_ONLY_SER (Java Serializer) ~200 bytes 2.0 GB Умеренные
MEMORY_ONLY_SER (Kryo + reg.) ~25 bytes 250 MB Минимальные

DataFrame Cache vs RDD Cache: важное отличие

# DataFrame .cache() использует Tungsten columnar format
# Это совершенно другой механизм — не через Kryo!
df.cache()  # → StorageLevel.MEMORY_AND_DISK
# Tungsten хранит данные в компактном columnar формате
# Kryo здесь не участвует

# RDD .cache() или .persist(MEMORY_ONLY_SER) — вот где Kryo работает
rdd = df.rdd  # конвертируем в RDD (каждая строка = Java Row object)
rdd.persist(StorageLevel.MEMORY_ONLY_SER)  # теперь Kryo сериализует
# Kryo кодирует каждый Row объект компактно

# Вывод: для DataFrame API Kryo почти не влияет на Cache
# Для RDD API Kryo + MEMORY_ONLY_SER = значительная экономия RAM

10. Best Practices, антипаттерны и когда Kryo применять в 2026

Где Kryo даёт ощутимый выигрыш сегодня

# ✅ СЛУЧАЙ 1: Spark MLlib / Spark GraphX (активные пользователи JVM объектов)
from pyspark.mllib.regression import LabeledPoint, LinearRegressionWithSGD
from pyspark.mllib.linalg import Vectors

# Каждый LabeledPoint сериализуется при каждом Shuffle
# На 100M объектов: JavaSerializer = 35 GB, Kryo+reg = 2.5 GB
training_data = sc.parallelize([
    LabeledPoint(1.0, Vectors.dense([0.1, 0.2, 0.3])),
    ...
])
# Регистрация классов MLlib:
# spark.kryo.classesToRegister = org.apache.spark.mllib.linalg.DenseVector,...

# ✅ СЛУЧАЙ 2: Сложные кастомные JVM объекты в RDD трансформациях
# Если вы вынуждены использовать RDD API с кастомными Java/Scala объектами
custom_rdd = sc.parallelize(list_of_java_objects)
custom_rdd.persist(StorageLevel.MEMORY_ONLY_SER)  # Kryo сжимает

# ✅ СЛУЧАЙ 3: Broadcast переменных (JVM objects)
# Broadcast Table переданная Driver'ом → Executor'ам
broadcast_lookup = sc.broadcast(large_dict)  # Python dict → pickle, не Kryo
# НО: внутренние JVM Broadcast (HashedRelation) используют Kryo

# ✅ СЛУЧАЙ 4: Streaming (legacy DStream API)
# DStream активно использует JVM сериализацию
from pyspark.streaming import StreamingContext
ssc = StreamingContext(sc, 1)
# Kryo ускоряет DStream checkpoint и window операции

# ⚠️ ОГРАНИЧЕННЫЙ СЛУЧАЙ: современный DataFrame/SQL
spark.conf.set("spark.serializer", "org.apache.spark.serializer.KryoSerializer")
df = spark.read.parquet("hdfs://cluster/data/")
result = df.groupBy("category").count()  # Tungsten format, не через Kryo!
# Эффект: минимальный. Task closures меньше, но данные обрабатывает Tungsten

Антипаттерны при работе с Kryo

# ❌ АНТИПАТТЕРН 1: Включить Kryo без registrationRequired для аудита
# Kryo тихо сериализует незарегистрированные классы с именами
# → Эффект есть, но не максимальный
# ✅ Правильно: включить registrationRequired=true в dev, найти все классы

# ❌ АНТИПАТТЕРН 2: Ожидать ускорения от Kryo для DataFrame SQL операций
# "Включу Kryo и мои Parquet запросы станут быстрее"
# ✅ Реальность: Tungsten format уже оптимален для DataFrame
# Kryo тут почти ничего не даёт

# ❌ АНТИПАТТЕРН 3: Игнорировать buffer.max при работе с ML
# MLlib с Random Forest: каждое дерево может весить 100+ MB
# Default buffer.max = 64 MB → Buffer overflow!
spark.conf.set("spark.kryoserializer.buffer.max", "512m")  # для ML обязательно!

# ❌ АНТИПАТТЕРН 4: Использовать MEMORY_ONLY без SER даже с Kryo
# MEMORY_ONLY всегда хранит живые Java объекты (не байты)
# Kryo для этого режима не используется!
rdd.persist(StorageLevel.MEMORY_ONLY)     # Kryo НЕ влияет
rdd.persist(StorageLevel.MEMORY_ONLY_SER) # Kryo ВЛИЯЕТ ✅

Полная конфигурация для production с комментариями

from pyspark.sql import SparkSession

def create_kryo_optimized_spark(
    app_name: str,
    ml_workload: bool = False,
    custom_classes: list[str] = None,
) -> SparkSession:
    """
    Создаёт SparkSession с оптимизированным KryoSerializer.

    ml_workload: True если пайплайн использует MLlib/GraphX
    custom_classes: список JVM-классов для регистрации
    """
    builder = SparkSession.builder.appName(app_name) \

        # Активируем KryoSerializer
        .config("spark.serializer",
                "org.apache.spark.serializer.KryoSerializer") \

        # В prod: false (не ломать при незарегистрированных классах)
        # В dev: true (найти все незарегистрированные классы)
        .config("spark.kryo.registrationRequired", "false") \

        # Начальный буфер: 256 KB для типичных объектов
        .config("spark.kryoserializer.buffer", "256k")

    # Для ML: большой буфер под векторы и модели
    if ml_workload:
        builder = builder \
            .config("spark.kryoserializer.buffer.max", "512m") \
            .config("spark.kryo.classesToRegister",
                    "org.apache.spark.mllib.regression.LabeledPoint,"
                    "org.apache.spark.mllib.linalg.DenseVector,"
                    "org.apache.spark.mllib.linalg.SparseVector,"
                    "org.apache.spark.mllib.linalg.Vectors,"
                    "org.apache.spark.mllib.feature.Word2VecModel")
    else:
        builder = builder \
            .config("spark.kryoserializer.buffer.max", "128m")

    # Добавляем пользовательские классы
    if custom_classes:
        existing = builder.getOrCreate().conf.get(
            "spark.kryo.classesToRegister", ""
        )
        all_classes = f"{existing},{','.join(custom_classes)}".lstrip(",")
        builder = builder.config("spark.kryo.classesToRegister", all_classes)

    return builder.getOrCreate()


# Использование:
spark_ml = create_kryo_optimized_spark(
    "ml-training-job",
    ml_workload=True,
    custom_classes=["com.company.model.FeatureVector"]
)

spark_etl = create_kryo_optimized_spark(
    "etl-job",
    ml_workload=False,
)

Итоги: место Kryo в современном Spark

KryoSerializer в 2026 году — это стандартная базовая настройка для любого production Spark-кластера. Включать его нужно, но с реалистичными ожиданиями.

Где даёт ощутимый выигрыш:

  • Spark MLlib (RDD-based): 5-15x компактнее данные → кратно меньше Shuffle трафик
  • Сложные пользовательские JVM типы в RDD API
  • RDD Cache с MEMORY_ONLY_SER: в 5-10x меньше RAM, нет GC давления
  • Spark GraphX и GraphFrames: меньше Shuffle на больших графах

Где эффект минимален (до 10%):

  • Чистые DataFrame/SQL операции с Parquet/Delta/Iceberg
  • Python UDF и Pandas UDF (они через pickle/Arrow, не Kryo)
  • Простые ETL с примитивными типами

Три вещи для немедленного внедрения:

  1. Включить spark.serializer=org.apache.spark.serializer.KryoSerializer — это практически бесплатно
  2. Для ML пайплайнов зарегистрировать MLlib классы и поднять buffer.max до 512m
  3. При использовании RDD Cache перейти на MEMORY_ONLY_SER — это особенно ощутимо