Tungsten: UnsafeRow, off-heap память и Whole-Stage CodeGen

Как Project Tungsten переписал движок Spark на уровне CPU и памяти: бинарное представление данных, обход GC и генерация Java bytecode прямо во время выполнения

core internals optimization

Почему классическая 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 работает целиком в памяти:

  1. WholeStageCodegen обходит дерево операторов, вызывая produce()/consume() на каждом
  2. Каждый оператор добавляет свой код в StringBuilder
  3. Итоговый Java-исходник (строка) передаётся в Janino.compile()
  4. Janino возвращает Class<?>, который Spark загружает через ClassLoader
  5. 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), и что происходит при нехватке каждого вида.