Tungsten: UnsafeRow, off-heap память и Whole-Stage CodeGen
Как Project Tungsten переписал движок Spark на уровне CPU и памяти: бинарное представление данных, обход GC и генерация Java bytecode прямо во время выполнения
Почему классическая JVM медленна для Big Data¶
До Spark 1.4 и проекта Tungsten Spark работал как типичное JVM-приложение: хранил данные в виде Java-объектов на heap, итерировался по строкам через цепочку вызовов Iterator.next(), сериализовал объекты при shuffle через Java Serialization или Kryo. Всё это - общие идиомы Java, разумные для OLTP-нагрузок, но катастрофические для обработки миллиардов строк.
Оверхед Java-объектов¶
Простое целое число 42 в Java не занимает 4 байта. Оно занимает минимум 16 байт как объект Integer (object header - 12 байт + 4 байта данных), а при хранении в коллекциях - ещё указатель на него (8 байт). Строка "hello" - это объект String с полем char[], что добавляет ещё один объект и ещё один header.
Integer(42):
┌──────────────┬──────────────┬──────────┐
│ mark word │ klass ptr │ value │
│ 8 bytes │ 4 bytes │ 4 bytes │
└──────────────┴──────────────┴──────────┘
= 16 байт для хранения числа 42
Tungsten UnsafeRow: 4 байта
Экономия: 4x
На таблице с 1 млрд строк и 10 колонками разница между Java-объектами и плотным бинарным форматом - десятки гигабайт дополнительной памяти.
GC pressure: Stop-the-World паузы¶
Garbage Collector Java работает с heap. Чем больше живых объектов - тем дольше каждая фаза GC. При обработке больших датасетов heap заполняется миллиардами объектов-строк. GC запускает Stop-the-World (STW) пауза: все потоки Spark-executor останавливаются на десятки миллисекунд или даже секунды.
Типичная картина в Spark UI (Task Metrics):
┌─────────────────────────────────────────┐
│ Duration: 120 s │
│ GC Time: 35 s ← 29%! │
│ Executor CPU Time: 85 s │
└─────────────────────────────────────────┘
Треть времени задача стоит - JVM убирает за собой мусор. Project Tungsten решил эту проблему радикально: вынести данные из-под контроля GC.
Volcano Model: виртуальные вызовы на каждую строку¶
Классический подход в реляционных движках - Volcano Model (iterator model): каждый оператор реализует метод next(), который возвращает одну строку, взяв её у нижестоящего оператора.
result = Filter.next()
↓ вызов
Filter → Project.next()
↓ вызов
Project → Scan.next()
↓ вызов
Scan → возвращает строку
Для 1 млрд строк с 5 операторами это 5 млрд виртуальных вызовов next(). Виртуальный вызов (через interface/abstract) дорог: JVM не может заинлайнить его статически, CPU branch predictor не может предсказать цель, instruction cache постоянно промахивается. Результат - CPU большую часть времени ждёт, а не считает.
Project Tungsten: три кита оптимизации¶
В 2015 году команда Databricks представила Project Tungsten - набор оптимизаций движка исполнения Spark, нацеленных на эффективность на уровне железа: CPU, кэши L1/L2/L3, RAM.
UnsafeRow: бинарное представление строки¶
Устройство UnsafeRow¶
UnsafeRow - внутренний класс Spark, представляющий строку DataFrame не как Java-объект, а как непрерывный массив байт в памяти. Структура:
UnsafeRow (5 колонок, row = [42, null, "hello", 3.14, true])
Byte offset:
┌──────────┬──────────────────────────────────────────┐
│ 0-7 │ Null Bitmap (1 бит на колонку) │
│ │ Bit 1 = 1 (колонка 1 null) │
├──────────┼──────────────────────────────────────────┤
│ 8-15 │ Поле 0: 42 (Long, 8 байт) │
│ 16-23 │ Поле 1: 0 (null, место зарезервировано) │
│ 24-31 │ Поле 2: offset+length "hello" (8 байт) │
│ 32-39 │ Поле 3: 3.14 (Double, 8 байт) │
│ 40-47 │ Поле 4: true (Boolean в 8 байтах) │
├──────────┼──────────────────────────────────────────┤
│ 48-52 │ Variable-length data: "hello" (5 байт) │
└──────────┴──────────────────────────────────────────┘
Ключевые свойства:
- Fixed-size поля (Int, Long, Double, Boolean) хранятся напрямую в фиксированной секции - 8 байт каждое
- Variable-size поля (String, Array, Map) хранятся в variable-length секции; в фиксированной секции - 8-байтный указатель (4 байта offset + 4 байта length)
- NULL кодируется через bitmap (0 = not null, 1 = null), данные не хранятся
- Весь ряд - один непрерывный блок памяти: можно скопировать через
System.arraycopyбез сериализации
sun.misc.Unsafe: прямой доступ к памяти¶
Spark использует sun.misc.Unsafe - недокументированный API JVM для прямой работы с памятью в обход обычных гарантий JVM:
// Примерная механика (упрощённо)
Unsafe unsafe = Unsafe.getUnsafe();
// Выделить off-heap память
long address = unsafe.allocateMemory(1024 * 1024); // 1 МБ вне heap
// Записать long по адресу
unsafe.putLong(address + offset, 42L);
// Прочитать long
long value = unsafe.getLong(address + offset);
// Освободить - явно, не через GC
unsafe.freeMemory(address);
Память, выделенная через Unsafe.allocateMemory, не управляется GC. Spark сам отслеживает выделение и освобождение через TaskMemoryManager и MemoryPool.
On-heap vs Off-heap хранение¶
# Управление off-heap в конфиге
spark.conf.set("spark.memory.offHeap.enabled", "true")
spark.conf.set("spark.memory.offHeap.size", "4g") # 4 ГБ off-heap на executor
По умолчанию Spark хранит UnsafeRow on-heap (в byte[]). Off-heap включается явно и полезен когда heap GC паузы критичны (например, большой unified memory pool).
Whole-Stage Code Generation¶
Whole-Stage Code Generation (WSCG) - вторая половина Tungsten. Вместо того чтобы интерпретировать физический план оператор за оператором (Volcano Model), Spark генерирует Java-код для всего пайплайна операторов, компилирует его и выполняет как единый метод.
От плана к коду: шаги¶
Что генерирует Spark: пример¶
Рассмотрим запрос:
df.filter(col("amount") > 100).select("customer_id", "amount")
Volcano Model (интерпретация):
FilterExec.next() вызывается 1 млрд раз
└─ ProjectExec.next() вызывается 1 млрд раз
└─ ScanExec.next() вызывается 1 млрд раз
Whole-Stage CodeGen генерирует примерно такой Java-код:
// Это упрощённая иллюстрация сгенерированного кода
final class GeneratedIterator extends BufferedRowIterator {
private UnsafeRowWriter rowWriter = new UnsafeRowWriter(2);
protected void processNext() throws IOException {
while (scan_batchIdx < scan_numBatch) {
// Scan: читаем напрямую из Parquet column vector
double amount = scan_doubleVector.getDouble(scan_batchIdx);
long customer_id = scan_longVector.getLong(scan_batchIdx);
scan_batchIdx++;
// Filter: встроен прямо в цикл - нет виртуального вызова!
if (!(amount > 100.0)) continue;
// Project: формируем UnsafeRow
rowWriter.reset();
rowWriter.write(0, customer_id); // customer_id
rowWriter.write(1, amount); // amount
append(rowWriter.getRow());
}
}
}
Один while-цикл заменил три уровня Iterator.next(). JIT-компилятор HotSpot видит простой цикл с предсказуемой ветвью if (!(amount > 100.0)) continue - и применяет:
- Inlining - убирает все вызовы методов
- Loop unrolling - разворачивает цикл в несколько итераций подряд
- Branch prediction - CPU предсказывает
continueи подгружает данные заранее - Vectorization (SIMD) - если архитектура позволяет, обрабатывает несколько значений за одну инструкцию
Janino: компилятор прямо в runtime¶
Spark использует Janino - лёгкий Java-компилятор, который можно встроить в приложение. Не javac (тот требует файловой системы и JDK), а Janino работает целиком в памяти:
WholeStageCodegenобходит дерево операторов, вызываяproduce()/consume()на каждом- Каждый оператор добавляет свой код в
StringBuilder - Итоговый Java-исходник (строка) передаётся в
Janino.compile() - Janino возвращает
Class<?>, который Spark загружает черезClassLoader - Spark создаёт экземпляр и вызывает
processNext()в цикле
Время компиляции - 5–50 мс на запрос. Это пренебрежимо мало по сравнению с минутами выполнения, но для очень коротких запросов overhead заметен.
Cache-Aware Computation¶
Иерархия памяти и cache locality¶
Современный CPU имеет несколько уровней кэша с радикально разным временем доступа:
| Кэш | Размер | Latency | Bandwidth |
|---|---|---|---|
| L1 data cache | 32–64 КБ | ~1 нс | ~1 ТБ/с |
| L2 cache | 256 КБ – 1 МБ | ~5 нс | ~400 ГБ/с |
| L3 cache | 8–64 МБ | ~20 нс | ~200 ГБ/с |
| DRAM (RAM) | ГБ–ТБ | ~100 нс | ~50 ГБ/с |
Cache line - минимальная единица загрузки из RAM в кэш - 64 байта. Если данные расположены последовательно в памяти (Sequential Access), каждый cache miss подгружает 64 байта полезных данных. Если данные разбросаны (Random Access через указатели), каждый miss даёт только 8 байт полезного и 56 байт мусора.
Сортировка с 8-байтными ключами¶
При sort Spark не перемещает UnsafeRow-объекты (они могут быть большими). Вместо этого строится массив 8-байтных записей: 4 байта - sort key (или хеш ключа), 4 байта - указатель на строку. Сортируется только этот компактный массив:
Sort Key Array (8 байт × N):
┌──────────┬──────────┐
│ key: 042 │ ptr: 0x0 │ → Row 1
│ key: 017 │ ptr: 0x1 │ → Row 2
│ key: 099 │ ptr: 0x2 │ → Row 3
└──────────┴──────────┘
↓ sort по key (компактный массив влезает в L2 cache)
┌──────────┬──────────┐
│ key: 017 │ ptr: 0x1 │ → Row 2
│ key: 042 │ ptr: 0x0 │ → Row 1
│ key: 099 │ ptr: 0x2 │ → Row 3
└──────────┴──────────┘
Весь массив ключей для 10 млн строк = 80 МБ - может целиком поместиться в L3 кэш типичного server CPU (8–32 МБ). Сравнение: массив указателей на Java-объекты того же размера вызовет 10 млн cache miss при каждом compare.
WSCG в Spark UI и explain()¶
Чтение физического плана¶
df.filter(col("amount") > 100) \
.groupBy("country") \
.agg({"amount": "sum"}) \
.explain("formatted")
== Physical Plan ==
*(2) HashAggregate(keys=[country#1], functions=[sum(amount#2)]) ← stage 2
+- Exchange hashpartitioning(country#1, 200)
+- *(1) HashAggregate(keys=[country#1], functions=[partial_sum(amount#2)]) ← stage 1
+- *(1) Filter (isnotnull(amount#2) AND (amount#2 > 100.0))
+- *(1) FileScan parquet [country#1,amount#2]
Расшифровка:
*(1)- операторы в WholeStageCodegen Stage 1: Scan + Filter + partial HashAggregate скомпилированы в один Java-класс*(2)- операторы в WholeStageCodegen Stage 2: финальный HashAggregate скомпилирован отдельноExchange- shuffle, разрывает WSCG (данные уходят по сети → новый stage)- Операторы без
*- выполняются в интерпретированном режиме
Просмотр сгенерированного кода¶
# Способ 1: через queryExecution
print(df.queryExecution.debug.codegen())
# Способ 2: через SQL
spark.conf.set("spark.sql.codegen.comments", "true") # добавляет комментарии в код
print(df.queryExecution.debug.codegen())
# Способ 3: для конкретного плана
from pyspark.sql.functions import col
df2 = df.filter(col("amount") > 100).select("customer_id", "amount")
print(df2.queryExecution.debug.codegen())
Пример вывода для простого filter + project:
/* 001 */ public Object generate(Object[] references) {
/* 002 */ return new GeneratedIteratorForCodegenStage1(references);
/* 003 */ }
/* 004 */
/* 005 */ final class GeneratedIteratorForCodegenStage1 extends BufferedRowIterator {
/* 006 */ private Object[] references;
/* 007 */ private scala.collection.Iterator[] inputs;
/* 008 */ private boolean scan_onHeap_0;
/* 009 */ // ... поля для векторов колонок ...
/* 010 */
/* 011 */ protected void processNext() throws java.io.IOException {
/* 012 */ while (scan_batchIdx < scan_numBatch) {
/* 013 */ // Filter: amount > 100
/* 014 */ double scan_value_1 = scan_doubleVector.getDouble(scan_batchIdx);
/* 015 */ if (!(scan_value_1 > 100.0)) { scan_batchIdx++; continue; }
/* 016 */ // Project: customer_id, amount
/* 017 */ long scan_value_0 = scan_longVector.getLong(scan_batchIdx);
/* 018 */ scan_batchIdx++;
/* 019 */ result.setLong(0, scan_value_0);
/* 020 */ result.setDouble(1, scan_value_1);
/* 021 */ append(result);
/* 022 */ }
/* 023 */ }
/* 024 */ }
Обратите внимание: никаких виртуальных вызовов, никакой абстракции - прямой цикл с прямым доступом к column vectors.
Когда WSCG не работает¶
Не все операторы поддерживают WSCG. Если в пайплайн попадает оператор без поддержки, он разрывает stage на две части:
Операторы, нарушающие WSCG:
| Оператор | Причина отключения WSCG |
|---|---|
Python UDF (udf) |
Выполняется в Python-процессе через Py4J, данные сериализуются |
Pandas UDF (pandas_udf) |
Лучше: батчи через Arrow, но всё равно разрыв codegen |
rdd.map(), rdd.filter() |
RDD API вне Catalyst полностью |
BroadcastNestedLoopJoin |
Не поддерживается в codegen |
| Очень длинные выражения | Компилятор JVM ограничивает размер метода (64 KB bytecode) |
Window с complex frame |
Частичная поддержка |
| Некоторые агрегатные функции | Зависит от реализации |
Диагностика разрыва:
# В explain() ищем операторы без *
df.explain("formatted")
# *(1) Filter ← с codegen
# BatchEvalPython [my_udf] ← без * → разрыв!
# *(2) Project ← новый codegen stage
# Включить подробные метрики
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
Ограничение 64 KB метода¶
JVM не позволяет создать bytecode-метод размером более 64 KB. При очень широких таблицах (200+ колонок) или сложных CASE WHEN выражениях сгенерированный код превышает лимит:
org.apache.spark.SparkException: Failed to compile:
org.codehaus.commons.compiler.CompileException: Code of method exceeds 64 KB limit
# Обходные пути:
# 1. Разбить запрос на несколько шагов
df = df.select("col1", "col2", ..., "col50") # сначала выбрать нужные
df = df.withColumn(...) # потом трансформировать
# 2. Отключить codegen для конкретного запроса
spark.conf.set("spark.sql.codegen.wholeStage", "false")
df.filter(...).show()
spark.conf.set("spark.sql.codegen.wholeStage", "true") # вернуть
# 3. Уменьшить число колонок через Column Pruning (SELECT нужных)
Columnar Vectorized Execution¶
Tungsten работает не только со строками (UnsafeRow), но и с батчами колонок при чтении Parquet. Это Vectorized Parquet Reader:
Батч из 4096 значений типа double занимает 32 КБ - помещается в L1 кэш. Filter над этим батчем - это SIMD-операция над непрерывным массивом. Spark codegen работает с ColumnVector напрямую, без создания промежуточных UnsafeRow.
# Vectorized reader включён по умолчанию для Parquet и ORC
spark.conf.get("spark.sql.parquet.enableVectorizedReader") # true
spark.conf.get("spark.sql.orc.enableVectorizedReader") # true
# Размер батча
spark.conf.set("spark.sql.parquet.columnarReaderBatchSize", "4096")
DataFrame vs RDD: где теряются оптимизации¶
# Потеря Tungsten при переходе к RDD
df = spark.table("orders")
# DataFrame: Tungsten всё время
result1 = df.filter("amount > 100").groupBy("country").sum("amount")
# Переход к RDD → потеря всех оптимизаций
result2 = df.rdd \
.filter(lambda row: row.amount > 100) \ # Python UDF + RDD
.map(lambda row: (row.country, row.amount)) \
.reduceByKey(lambda a, b: a + b)
# Возврат к DataFrame - Tungsten снова работает
result3 = spark.createDataFrame(result2, ["country", "total"])
# Но: round-trip через RDD = 2x сериализация + потеря статистики
Практика: диагностика Tungsten¶
Шаг 1: Убедиться что WSCG включён¶
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName("Tungsten-Demo") \
.getOrCreate()
# Проверить конфиги
print(spark.conf.get("spark.sql.codegen.wholeStage")) # true
print(spark.conf.get("spark.sql.codegen.fallback")) # true (fallback при ошибке)
print(spark.conf.get("spark.sql.parquet.enableVectorizedReader")) # true
Шаг 2: Сравнить codegen vs interpreted¶
import time
df = spark.range(100_000_000).withColumn("v", col("id") * 1.5)
# С codegen (по умолчанию)
spark.conf.set("spark.sql.codegen.wholeStage", "true")
t0 = time.time()
df.filter(col("v") > 1000).agg({"v": "sum"}).show()
print(f"With codegen: {time.time() - t0:.2f}s")
# Без codegen
spark.conf.set("spark.sql.codegen.wholeStage", "false")
t0 = time.time()
df.filter(col("v") > 1000).agg({"v": "sum"}).show()
print(f"Without codegen: {time.time() - t0:.2f}s")
# Типичный результат:
# With codegen: 2.3s
# Without codegen: 8.7s ← 3.8x медленнее
Шаг 3: Найти bottleneck через explain и UI¶
# Видим план с WSCG stages
df.filter(col("amount") > 100) \
.groupBy("country") \
.sum("amount") \
.explain("formatted")
# В Spark UI → SQL → выбрать query → DAG:
# WholeStageCodegen (1): зелёный прямоугольник, несколько операторов внутри
# Exchange: серый прямоугольник (shuffle, разрыв)
# WholeStageCodegen (2): следующий зелёный прямоугольник
Шаг 4: Диагностика GC через Task Metrics¶
# В Spark UI → Stages → выбрать stage → Task Metrics table:
# Смотреть колонки:
# GC Time - если > 10% от Task Duration → GC pressure
# Peak Execution Memory - пиковая память executor
# Spill (Memory) / Spill (Disk) - если есть → нехватка памяти
# Программно через SparkListener (advanced)
Шаг 5: Просмотр сгенерированного кода¶
# Включить комментарии в codegen (полезно для понимания)
spark.conf.set("spark.sql.codegen.comments", "true")
query = df.filter(col("amount") > 100).select("customer_id", "amount")
# Получить сгенерированный код
codegen_str = query.queryExecution.debug.codegen()
print(codegen_str)
# Сохранить в файл для изучения
with open("/tmp/generated_code.java", "w") as f:
f.write(codegen_str)
Best Practices для Tungsten-friendly кода¶
1. Использовать встроенные функции вместо Python UDF:
from pyspark.sql import functions as F
# Плохо: Python UDF разрывает codegen
@udf("double")
def calc_tax(amount):
return amount * 0.2
df.withColumn("tax", calc_tax(col("amount")))
# Хорошо: встроенная функция → WSCG не разрывается
df.withColumn("tax", col("amount") * 0.2)
df.withColumn("tax", F.round(col("amount") * 0.2, 2))
2. Избегать лишних rdd.* операций:
# Плохо: два round-trip через RDD
result = df.rdd.filter(lambda r: r.amount > 100).toDF()
# Хорошо: остаёмся в DataFrame
result = df.filter(col("amount") > 100)
3. Columnar форматы для максимальной пользы от Vectorized Reader:
# Parquet и ORC дают Vectorized Read → SIMD обработка батчей
df.write.parquet("/data/orders") # хорошо
df.write.csv("/data/orders") # csv не векторизован
4. Контролировать размер выражений:
# Плохо при 100+ условиях - превысит 64 KB bytecode limit
condition = " OR ".join([f"category = '{c}'" for c in categories]) # 200 категорий
df.filter(condition)
# Хорошо: join с маленькой таблицей-фильтром
categories_df = spark.createDataFrame([(c,) for c in categories], ["category"])
df.join(categories_df, "category", "semi")
5. Off-heap при GC pressure:
# Если в Spark UI видите GC Time > 15% - включить off-heap
spark = SparkSession.builder \
.config("spark.memory.offHeap.enabled", "true") \
.config("spark.memory.offHeap.size", "8g") \
.getOrCreate()
Итого¶
| Характеристика | RDD (без Tungsten) | DataFrame (с Tungsten) |
|---|---|---|
| Хранение данных | Java-объекты в Heap | UnsafeRow (бинарный формат) |
| Garbage Collection | Высокое давление, STW паузы | Минимальное (или нулевое off-heap) |
| Выполнение | Volcano: N×M виртуальных вызовов | WSCG: один цикл, нет виртуальных вызовов |
| Использование CPU кэша | Cache miss при каждом pointer chasing | Sequential access, prefetching, SIMD |
| Сериализация | Java Serialization / Kryo | Нет: UnsafeRow = готовый бинарный буфер |
| Чтение колонок | Строка за строкой | Vectorized: батч 4096 значений |
| Ускорение vs RDD | 1x (baseline) | 5–20x типично |
Следующий урок рассматривает Unified Memory Model - как Spark делит память executor между execution (shuffle, join, sort) и storage (кэш RDD/DataFrame), и что происходит при нехватке каждого вида.