KryoSerializer: регистрация классов и измеримый выигрыш
Полный разбор KryoSerializer в Spark: проблема JavaSerializer и раздутые метаданные, физика сериализации в Shuffle/Broadcast/Cache, включение Kryo, обязательная регистрация классов и spark.kryo.registrationRequired, KryoRegistrator, настройка буферов, измерение выигрыша в Spark UI, область применения в PySpark и Tungsten контекст.
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:
- В dev/test среде включить
registrationRequired=true - Запустить пайплайн и собрать все
KryoException— это список классов для регистрации - Зарегистрировать все найденные классы
- В prod оставить
registrationRequired=falseкак страховку (на случай если появятся новые классы) - При крупных новых разработках периодически повторять шаги 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 с примитивными типами
Три вещи для немедленного внедрения:
- Включить
spark.serializer=org.apache.spark.serializer.KryoSerializer— это практически бесплатно - Для ML пайплайнов зарегистрировать MLlib классы и поднять
buffer.maxдо 512m - При использовании RDD Cache перейти на
MEMORY_ONLY_SER— это особенно ощутимо