Проблема JVM: почему Tungsten не может использовать SIMD и при чём тут колоночное хранение

Проблема JVM: почему Tungsten не может использовать SIMD и при чём тут колоночное хранение

optimization

Исторический контекст: эволюция узких мест в Spark

От I/O-bound к CPU-bound

Когда Apache Spark появился в 2014 году, главным узким местом аналитических систем была сеть и диски. Жёсткие диски давали 100-200 MB/s, сеть в датацентре - 1-10 Gbps. В такой среде shuffle-операция (передача данных между executor'ами) могла занимать 80% времени выполнения джоба. Алгоритмы оптимизировались прежде всего для минимизации I/O: именно поэтому появились predicate pushdown, column pruning, in-memory caching.

Прошло десять лет. Современный датацентр - это NVMe-диски со скоростью 5-7 GB/s, 100-гигабитная Ethernet-сеть, S3 с параллельной пропускной способностью в десятки GB/s. Данные перемещаются быстро. Узкое место переместилось: теперь на большинстве тяжёлых аналитических джобов CPU и память - вычисление GROUP BY по терабайту данных, хэш-джоин двух больших таблиц, агрегация с оконными функциями.

Это называется переход от I/O-bound к CPU-bound. Spark был спроектирован под один профиль узких мест, а эксплуатируется теперь в другом. Понимание этого перехода объясняет, почему сегодня появляются принципиально новые архитектуры вычислительных движков.

Что значит CPU-bound в аналитике

Когда мы говорим «джоб уперся в CPU», это не всегда означает, что процессор загружен на 100%. Иногда процессор загружен на 100%, но при этом почти ничего не вычисляет - он ждёт данных из памяти. Это называется memory bandwidth bottleneck: CPU может обрабатывать данные быстрее, чем оперативная память успевает их доставить.

Современный процессор Intel/AMD выполняет миллиарды операций в секунду. Оперативная память DDR5 имеет пропускную способность 50-80 GB/s. При обработке терабайтной таблицы процессор может буквально простаивать, ожидая, пока RAM загрузит следующую порцию данных. Это ещё один вид CPU-bound - не вычислительный, а коммуникационный.

Оба вида CPU-bound напрямую связаны с тем, как данные организованы в памяти: как часто процессор может найти нужные данные в быстром кэше (L1/L2/L3), а когда вынужден обращаться к медленной RAM.


Project Tungsten: попытка выжать максимум из JVM

Что было до Tungsten

До Spark 2.x (до Project Tungsten) Spark работал с данными как с обычными Java-объектами. Каждая строка DataSet была Java-объектом в heap-памяти, каждый Integer - объектом java.lang.Integer (а не примитивным int). Это имело серьёзные последствия:

Огромный overhead памяти. Java-объект занимает минимум 16 байт заголовка, плюс данные. Примитивный int - 4 байта, объект Integer - 16 байт заголовка + 4 байта данных + 4 байта выравнивания = 24 байта. В 6 раз больше. Строка с 10 числовыми полями в виде Java-объектов занимала в 5-10 раз больше памяти, чем те же данные в компактном бинарном представлении.

Давление на Garbage Collector. При обработке миллиарда строк создаётся миллиард Java-объектов. GC должен периодически останавливать всё приложение (Stop-The-World pause) чтобы освободить память. В аналитических джобах GC-паузы могут занимать от секунд до минут - это прямые потери производительности.

Pointer chasing. Java-объекты в heap-памяти разбросаны хаотично. Ссылка на объект - это адрес в памяти, по которому нужно «прыгнуть». При обходе коллекции Java-объектов процессор делает сотни тысяч «прыжков» по памяти, каждый из которых - потенциальный cache miss. Кэш L1 имеет время доступа ~4 такта, оперативная память - ~200 тактов. Каждый cache miss - потеря 200 тактов вхолостую.

Три столпа Project Tungsten

Project Tungsten, появившийся в Spark 2.x, атаковал каждую из этих проблем:

1. Off-Heap Memory Management (управление памятью вне кучи). Spark начал выделять память напрямую через sun.misc.Unsafe - это API, позволяющее работать с raw-памятью, минуя JVM heap. Данные, хранящиеся в off-heap-памяти, не видны GC и не создают давления на него. Это устранило GC-паузы для больших датасетов.

2. Binary Row Format - UnsafeRow. Вместо Java-объектов Spark начал хранить строки данных как компактные бинарные blob'ы. Структура UnsafeRow:

[Null Bitmap (8 bytes)] [Fixed-Length Fields] [Variable-Length Fields]
                        [int32][int64][double] [string offset+length][string data]
  • Null bitmap: битовая маска, где каждый бит указывает, является ли поле NULL
  • Fixed-length fields: числа хранятся непосредственно как 4 или 8 байт
  • Variable-length fields: строки хранятся как offset+length в начале и данные в конце

Вместо объекта Integer (24 байта) - просто 4 байта в нужном месте. Вместо объекта String с internal char array - компактное представление с offset. Размер строки сократился в 5-10 раз, pointer chasing исчез.

3. Whole-Stage Code Generation (WSCG). Вместо цепочки операторов (Filter → Project → Aggregate), каждый из которых вызывает следующий через виртуальный вызов, Spark начал генерировать монолитный Java-код, объединяющий все операторы в одну функцию:

// Вместо: строка → Filter.process(строка) → Project.process(строка) → Aggregate.process(строка)
// Генерируется единый Java-метод:
void processRow(UnsafeRow row) {
    // Inline: фильтр
    int amount = row.getInt(2);
    if (amount <= 0) return;
    // Inline: проекция
    long userId = row.getLong(0);
    // Inline: агрегация
    hashMap.update(userId, amount);
}

Это устранило накладные расходы на виртуальные вызовы методов и улучшило использование I-cache процессора (instruction cache): один большой горячий метод вместо цепочки маленьких.

Успех и потолок Tungsten

Tungsten дал значительный прирост производительности - в 2-5 раз на типичных рабочих нагрузках. Это был огромный шаг. Spark из «медленного MapReduce» превратился в серьёзный аналитический движок.

Но Tungsten уткнулся в потолок, который нельзя преодолеть, оставаясь в рамках JVM. Этот потолок - принципиальная невозможность эффективного использования SIMD-инструкций процессора.


Анатомия SIMD: параллелизм на уровне процессора

Что такое SIMD

SIMD (Single Instruction, Multiple Data) - это архитектурный принцип, при котором одна процессорная инструкция применяется одновременно к нескольким элементам данных. Это фундаментально иной вид параллелизма: не параллелизм потоков (несколько ядер), не параллелизм задач (несколько машин) - а параллелизм внутри одного ядра, в одном такте процессора.

Современные процессоры Intel/AMD имеют специализированные SIMD-регистры:

  • SSE (Streaming SIMD Extensions): 128-битные регистры (4 × int32 или 2 × int64 за одну операцию)
  • AVX2 (Advanced Vector Extensions 2): 256-битные регистры (8 × int32 или 4 × int64)
  • AVX-512: 512-битные регистры (16 × int32 или 8 × int64)

Что это означает на практике: если у вас есть массив из 16 целых чисел, и вы хотите прибавить к каждому константу 100:

Скалярное (обычное) выполнение:

add eax, 100    ; такт 1: прибавляем к элементу [0]
add eax, 100    ; такт 2: прибавляем к элементу [1]
add eax, 100    ; такт 3: прибавляем к элементу [2]
...
add eax, 100    ; такт 16: прибавляем к элементу [15]

16 тактов для 16 элементов.

SIMD (AVX-512) выполнение:

vmovdqa32 zmm0, [data]        ; загрузить 16 × int32 в регистр zmm0
vpbroadcastd zmm1, 100         ; 16 раз скопировать константу 100 в zmm1
vpaddd zmm2, zmm0, zmm1        ; сложить векторно: zmm2 = zmm0 + zmm1 (все 16 за один такт)
vmovdqa32 [result], zmm2       ; сохранить 16 результатов

4 инструкции, примерно 4 такта - для всех 16 элементов. Ускорение в 4× по сравнению со скалярным кодом, при той же тактовой частоте.

Для агрегации по миллиарду строк - это разница между 10 секундами и 2.5 секундами только от применения AVX-512.

SIMD в аналитических операциях

Аналитические операции особенно хорошо поддаются SIMD-ускорению:

Фильтрация: WHERE amount > 0 - сравниваем 16 значений amount с нулём за одну AVX-512 инструкцию, получаем 16-битную маску результатов.

Агрегация: SUM(amount) - суммируем 16 значений за одну инструкцию, затем горизонтально складываем регистр.

Хэширование: CRC32-хэш для GROUP BY - аппаратные инструкции CRC32 на последовательном массиве в 16 раз быстрее, чем поэлементно.

Сортировка: Bitonic sort, vectorized merge - стандартные алгоритмы ускоряются в разы при колоночном хранении.

Ключевое требование для SIMD: данные должны лежать в памяти последовательно (contiguous array). Если нужно просуммировать колонку amount - все значения amount должны идти подряд в памяти: amount[0], amount[1], amount[2], .... Именно тогда можно загрузить 16 значений одной инструкцией vmovdqa32.

Vectorized Execution: обработка батчами

Концепция, делающая SIMD применимой к аналитике, называется vectorized execution (не путать с SIMD - это более широкая концепция):

Вместо обработки строк одну за другой:

Строка 1 → Filter → Project → Aggregate
Строка 2 → Filter → Project → Aggregate
...
Строка 1 000 000 → Filter → Project → Aggregate

Обработка происходит батчами (обычно 1024-4096 строк за раз):

Батч строк 1-1024:
  Колонка amount[0..1023] → Filter (SIMD: 64 итерации по 16 элементов)
  Колонка user_id[0..1023] → Project (копирование массива)
  Колонка amount[0..1023] → Aggregate (SIMD: vectorized sum)

Это принципиально другая архитектура исполнения. Каждый оператор работает с целым батчем данных, представленным как массив значений (колоночный батч). Именно так работают DuckDB, ClickHouse, Apache Velox - нативные аналитические движки.


Проблема JVM: почему Java не умеет в SIMD

Абстракция против производительности

JVM (Java Virtual Machine) проектировалась с фундаментальным приоритетом: Write Once, Run Anywhere. Один и тот же Java-байткод должен выполняться на x86_64, ARM, RISC-V, старых Sparc-процессорах. Это великолепная инженерная цель - но она несовместима с аппаратно-специфичными оптимизациями.

SIMD-инструкции (AVX2, AVX-512) специфичны для конкретной микроархитектуры. Код, написанный для Intel AVX-512, не запустится на старом процессоре без этих инструкций. Для JVM это неприемлемо: программа должна работать везде.

Выход, который JVM пытается использовать - автовекторизация JIT-компилятора: HotSpot JIT анализирует байткод и при определённых условиях может заменить скалярный цикл на SIMD-инструкции. Это происходит автоматически, без участия программиста.

Но автовекторизация в JVM работает крайне нестабильно и ограниченно.

Почему автовекторизация JIT ненадёжна

Условия для автовекторизации очень строгие. JIT может векторизовать только простейшие циклы:

// ЭТО JIT может векторизовать:
double[] a = new double[n];
double[] b = new double[n];
double[] c = new double[n];
for (int i = 0; i < n; i++) {
    c[i] = a[i] + b[i];  // Простое поэлементное сложение
}

// ЭТО JIT не может векторизовать:
for (int i = 0; i < n; i++) {
    if (a[i] > 0) {          // Условная ветка - блокирует векторизацию
        c[i] = a[i] + b[i];
    }
}

В реальном Spark-коде таких «идеальных» циклов почти нет. GROUP BY с хэш-таблицей, JOIN с probing, фильтры с условиями - все они содержат ветки, обращения к хэш-таблицам, динамические адреса. JIT не сможет их векторизовать.

Проверки безопасности блокируют SIMD. JVM гарантирует: выход за границы массива → ArrayIndexOutOfBoundsException, разыменование null → NullPointerException. Для этого JIT вставляет скрытые проверки в каждый доступ к массиву:

; Что видит программист:
c[i] = a[i] + b[i];

; Что генерирует JIT:
cmp i, a.length      ; проверка границ для a[i]
jge throw_aiobe
cmp i, c.length      ; проверка границ для c[i]
jge throw_aiobe
null_check(a)        ; проверка на null
null_check(c)        ; проверка на null
mov eax, [a + i*4]   ; загрузка a[i]
mov ebx, [b + i*4]   ; загрузка b[i]
add eax, ebx         ; сложение
mov [c + i*4], eax   ; сохранение c[i]

Эти скрытые проверки происходят в каждой итерации. Они нарушают паттерн «чистая последовательная обработка массива», который AVX-компиляторы ищут для векторизации. SIMD-инструкция не может проверить границы для 16 элементов сразу так, как это ожидает JVM.

Иногда JIT способен доказать, что проверки избыточны (например, цикл for (int i = 0; i < a.length; i++) - явно в пределах), и убрать их. Но это работает только в простейших сценариях.

Whole-Stage Code Generation мешает векторизации. Это парадоксально, но инновация Tungsten - WSCG - сама мешает векторизации JIT. WSCG генерирует огромные функции, объединяя 5-10 операторов:

// Сгенерированный WSCG код (упрощённо):
void processWholeStage(long[] colA, long[] colB, HashMap hashMap) {
    for (int i = 0; i < numRows; i++) {
        long a = colA[i];          // Filter check
        if (a > 0) {               // Branch!
            long b = colB[i];
            long hash = hash(a);   // Вызов функции!
            HashEntry entry = hashMap.getOrCreate(hash);  // Pointer access!
            entry.sum += b;
        }
    }
}

Такой код содержит if-ветку, вызов функции (hash()), обращение к хэш-таблице через указатель. JIT видит этот цикл как слишком сложный для векторизации - и отказывается от SIMD.

Маленькие простые циклы JIT мог бы векторизовать. Один большой сложный цикл (WSCG) - нет.

Project Panama: медленный ответ JVM

Java 16+ представил Vector API (Project Panama) - возможность явно использовать SIMD из Java-кода:

// Java Vector API (явный SIMD):
import jdk.incubator.vector.*;

static final VectorSpecies<Integer> SPECIES = IntVector.SPECIES_512;

void vectorAdd(int[] a, int[] b, int[] c) {
    int i = 0;
    for (; i < SPECIES.loopBound(a.length); i += SPECIES.length()) {
        IntVector va = IntVector.fromArray(SPECIES, a, i);
        IntVector vb = IntVector.fromArray(SPECIES, b, i);
        va.add(vb).intoArray(c, i);
    }
    // Скалярный хвост для остатка
    for (; i < a.length; i++) c[i] = a[i] + b[i];
}

Это решает проблему явного использования SIMD из Java. Но для Spark это не помогает по нескольким причинам:

  • Vector API только в инкубаторном статусе с Java 16 и всё ещё нестабилен (Java 21+)
  • Переписать Spark Tungsten под Vector API - это многолетний проект
  • Накладные расходы JVM (GC, object overhead) остаются
  • Даже с явным SIMD JVM не конкурирует с нативным C++ в аналитических задачах

При чём тут колоночное хранение

Конфликт форматов: Parquet vs UnsafeRow

Apache Parquet - колоночный формат. Данные хранятся по колонкам: сначала все значения user_id, затем все значения amount, затем все значения event_date. Это оптимально для аналитических запросов: SELECT SUM(amount) FROM orders WHERE event_date > '2024-01-01' требует только две колонки - Parquet читает только их.

UnsafeRow (формат Tungsten) - строчный формат. Данные хранятся по строкам: все поля одной строки рядом, затем все поля следующей строки. Это оптимально для операций, которые обрабатывают строку целиком: JOIN (нужны все поля строки для сравнения), UDF (получает строку и возвращает строку).

Parquet (колоночный):
user_id: [1, 2, 3, 4, 5, ...]
amount:  [100, 200, 150, 300, 50, ...]
date:    [2024-01-01, 2024-01-01, 2024-01-02, ...]

UnsafeRow (строчный, Tungsten):
[user_id=1 | amount=100 | date=2024-01-01]
[user_id=2 | amount=200 | date=2024-01-01]
[user_id=3 | amount=150 | date=2024-01-02]

Парадокс: Parquet → Tungsten → потеря потенциала

Вот что происходит при выполнении SELECT SUM(amount) FROM orders:

  1. Spark читает Parquet: получает только колонку amount как непрерывный массив [100, 200, 150, 300, 50, ...] (отлично для SIMD!)
  2. Spark конвертирует в UnsafeRow: каждое значение amount упаковывается в строку [null_bitmap | amount_value] (теряем непрерывность колонки!)
  3. Tungsten обрабатывает строки последовательно: row[0].getInt(1) → +, row[1].getInt(1) → +, ...

После шага 2 данные amount больше не лежат непрерывно - они «разбросаны» по UnsafeRow-блобам, разделённые null bitmap и другими полями. Загрузить их в SIMD-регистр одной инструкцией невозможно.

Мы прочитали данные в колоночном формате (хорошо!) и немедленно переложили их в строчный (плохо!). Колоночный потенциал уничтожен на этапе конвертации.

Влияние CPU-кэша на производительность

Чтобы понять, почему это важно, нужно понять иерархию кэша процессора:

Регистры:     ~0.1 нс,  ~1 KB     (прямо в ALU)
L1 кэш:       ~1 нс,    ~32 KB    (на каждом ядре)
L2 кэш:       ~5 нс,    ~256 KB   (на каждом ядре)
L3 кэш:       ~20 нс,   ~8-32 MB  (общий для всех ядер)
DRAM (RAM):   ~60 нс,   ~64-512 GB

Производительность критически зависит от того, находятся ли нужные данные в L1/L2, или процессор вынужден идти в RAM. Разница в скорости - 60×.

Колоночное хранение дружит с кэшем. При сканировании колонки amount данные читаются последовательно: amount[0], amount[1], amount[2], .... Процессор заранее (prefetch) загружает следующие элементы в кэш. L1 кэш на 32 KB вмещает 8192 значений int32 - весь батч обрабатывается из горячего кэша.

Строчное хранение (UnsafeRow) ломает кэш-локальность. При сканировании поля amount во всех строках нужно «перескакивать»: взяли amount из строки 0 (bytes 8-11), перескочили на bytes 0-11 строки 1 (там null bitmap + другие поля), взяли amount из строки 1 (bytes 8-11)... Если строка широкая (много полей), то между двумя последовательными значениями amount находится много байт других полей. Кэш загружает эти «лишние» байты, тратит кэш-строки на данные, которые не нужны.

Vectorized Reader в Spark: частичное решение

Spark знает об этой проблеме. Начиная с Spark 2.3 появился Vectorized Parquet Reader (spark.sql.parquet.enableVectorizedReader=true, включён по умолчанию).

Вместо построчного чтения Parquet он читает батчи колонок (ColumnarBatch) - структуру данных, где каждая колонка хранится как непрерывный массив:

// ColumnarBatch: внутри Spark
class ColumnarBatch {
    WritableColumnVector[] columns;  // Массив колонок
    int numRows;
    // columns[0] = вектор user_id: [1,2,3,4,5,...]
    // columns[1] = вектор amount:  [100,200,150,300,50,...]
}

Это позволяет читать Parquet без немедленной конвертации в UnsafeRow. Некоторые операторы Spark (в частности, некоторые агрегации и фильтры) могут работать напрямую с ColumnarBatch.

Но это половинчатое решение:

  • Поддержано только частично - многие операторы Spark всё равно конвертируют ColumnarBatch в UnsafeRow
  • Даже работая с ColumnarBatch, Spark не использует SIMD (по причинам JVM, описанным выше)
  • Когда данные попадают в Join, shuffle, sort - они конвертируются в UnsafeRow

Apache Arrow: единый стандарт колоночной памяти

Apache Arrow - это открытый стандарт колоночного формата для представления данных в оперативной памяти. Ключевые свойства:

  • Данные хранятся по колонкам как непрерывные буферы
  • Стандартизированный layout, одинаковый для C++, Java, Python, Rust
  • Zero-copy обмен данными между процессами (через shared memory или mmap)
  • Нативная поддержка null-значений через validity bitmap
  • Оптимизирован для SIMD: alignment гарантирован

Arrow играет критическую роль в двух местах Spark:

1. PySpark ↔ JVM обмен данных. До Arrow каждая передача данных из JVM в Python требовала сериализации через Pickle или неэффективный row-by-row обмен. С Arrow (pandas UDF) данные передаются как Arrow-буферы: JVM записывает колонки в Arrow-формате, Python читает те же буферы через PyArrow без копирования.

# Без Arrow (медленно):
@udf("double")
def slow_udf(x):          # Один вызов = одно значение, row-by-row
    return x * 2.0

# С Arrow (быстро: Pandas UDF):
@pandas_udf("double")
def fast_udf(s: pd.Series) -> pd.Series:   # Батч значений как numpy array
    return s * 2.0

Разница в производительности - от 10× до 100× для CPU-bound операций.

2. Spark SQL нативные движки. Arrow служит мостом между JVM Spark и нативными C++/Rust движками. Нативный движок (Velox, ClickHouse) получает данные в формате Arrow, обрабатывает их с полным SIMD, возвращает результат в Arrow обратно в JVM.


Выход из тупика: нативные движки

Нативный код вместо JVM

Понимание проблемы привело к появлению нового класса аналитических движков: нативных (написанных на C++ или Rust), которые компилируются в машинный код и имеют полный доступ к SIMD-инструкциям.

Нативный C++ код:

// C++: явное использование AVX-512
#include <immintrin.h>

void sum_column_avx512(const int32_t* data, int64_t n, int64_t& result) {
    __m512i sum_vec = _mm512_setzero_si512();     // Обнулить вектор
    int64_t i = 0;

    for (; i + 16 <= n; i += 16) {
        __m512i batch = _mm512_loadu_si512(data + i);  // Загрузить 16 × int32
        sum_vec = _mm512_add_epi32(sum_vec, batch);    // Сложить 16 за такт
    }
    // Горизонтальная сумма вектора
    result = _mm512_reduce_add_epi32(sum_vec);

    // Хвост: оставшиеся элементы (< 16) скалярно
    for (; i < n; i++) result += data[i];
}

Этот код компилятор (GCC, Clang) превращает в точные AVX-512 инструкции. Нет JVM, нет GC, нет скрытых checks - чистое железо.

Databricks Photon: пионер нативного Spark

Databricks Photon - проприетарный нативный векторизованный движок, разработанный Databricks для замены Tungsten-части Spark. Написан на C++, использует AVX2/AVX-512, работает напрямую с Apache Arrow в колоночном виде.

Photon прозрачен для пользователя: тот же Spark SQL API, те же DataFrame операции - но под капотом вычисления идут в нативном C++ коде с SIMD. Databricks рапортует об ускорении в 2-12× в зависимости от типа запроса по сравнению со стандартным Spark Tungsten.

Архитектура Photon: Catalyst генерирует физический план, как обычно. Но вместо Tungsten физический план передаётся в Photon, который выполняет операторы в нативном коде, работая с данными в Apache Arrow ColumnarBatch. Результат возвращается в JVM только для финального вывода.

Apache Gluten: Velox/ClickHouse backend для Spark

Apache Gluten - open-source проект, предоставляющий нативный execution backend для Apache Spark. Он заменяет физические операторы Tungsten на реализации из двух нативных движков:

  • Velox (Meta/Intel): универсальный C++ execution engine с полной поддержкой SIMD
  • ClickHouse backend: движок ClickHouse, адаптированный для Spark

Принцип работы Gluten: Spark Catalyst строит физический план. Gluten перехватывает план и трансформирует его в представление, понятное нативному движку (Velox Substrait Plan). Нативный движок выполняет план в C++ с SIMD. Результат возвращается в Apache Arrow формате в Spark.

[Spark SQL Query]
       │
       ▼
[Catalyst Optimizer]  ← Строит логический и физический план (JVM, как обычно)
       │
       ▼
[Gluten Plugin]       ← Перехватывает физический план
       │
       ▼
[Native Engine]       ← Velox или ClickHouse backend (C++, SIMD)
  ColumnarBatch → Filter (AVX-512) → HashAgg (SIMD) → Sort → Result
       │
       ▼
[Arrow Result]        ← Колоночный результат обратно в Spark

Gluten доступен как плагин, не требует изменения кода приложения. Тот же spark.sql("SELECT ...") - другой execution engine под капотом.

Экономический эффект

Переход с Tungsten на нативный движок (Photon, Gluten+Velox) даёт:

  • Снижение CPU-времени в 2-5× на типичных аналитических запросах (GROUP BY, JOIN, агрегации)
  • Снижение memory bandwidth usage за счёт колоночной обработки
  • Прямое снижение облачных затрат: меньше CPU-часов → меньше счёт AWS/GCP/Azure
  • Ускорение TTR (Time To Result) для интерактивных запросов

Для крупной платформы, тратящей $500 000 в месяц на Spark-кластеры - переход на Photon может сэкономить $200 000-300 000 в месяц. Это не теория: Databricks публикует кейсы с такими цифрами.


Диаграмма: Путь данных - Parquet до результата

Схема показывает два пути от Parquet до результата. Стандартный Spark (красный путь): колоночный Parquet → ColumnarBatch → конвертация в UnsafeRow → скалярный JVM-цикл через WSCG. Нативный путь (зелёный): ColumnarBatch → Apache Arrow → C++ движок с AVX-512. Ключевая развилка - момент после Vectorized Reader.


Практика: измерение CPU-эффективности в Spark UI

Анализ физического плана

Чтобы понять, как Spark реально выполняет запрос, используем EXPLAIN FORMATTED:

spark.sql("""
    SELECT user_id, SUM(amount) AS total, COUNT(*) AS cnt
    FROM orders
    WHERE status = 'completed'
    GROUP BY user_id
""").explain("formatted")

Ключевые операторы в плане:

== Physical Plan ==
*(2) HashAggregate(keys=[user_id#1], functions=[sum(amount#2), count(1)])    ← Финальная агрегация
+- Exchange hashpartitioning(user_id#1, 200), ENSURE_REQUIREMENTS             ← Shuffle
   +- *(1) HashAggregate(keys=[user_id#1], functions=[partial_sum(amount#2)]) ← Частичная агрегация
      +- *(1) Project [user_id#1, amount#2]                                   ← Проекция
         +- *(1) Filter (isnotnull(status#3) AND (status#3 = completed))      ← Фильтр
            +- *(1) ColumnarToRow                                              ← КОНВЕРТАЦИЯ!
               +- FileScan parquet [user_id,amount,status] ...
                  Format: Parquet, PartitionFilters: [], PushedFilters: [IsNotNull(status), EqualTo(status, completed)]
                  ReadSchema: struct<user_id:bigint,amount:int,status:string>

Обратите внимание на ColumnarToRow - это момент, когда колоночный ColumnarBatch от Vectorized Reader конвертируется в строчный UnsafeRow. После этого все операторы (Filter, Project, HashAggregate) работают в строчном режиме.

Символ *(1) означает, что оператор выполняется внутри Whole-Stage CodeGen (один скомпилированный метод). *(2) - второй стейдж.

Spark UI метрики для диагностики CPU-bound

В Spark UI (вкладка SQL → конкретный джоб) смотрим метрики:

Stage 0:
  Task Metrics:
    Duration:              2m 45s
    GC Time:               15s         ← GC занимает 9% времени - приемлемо
    CPU Time:              12m 30s     ← CPU time > Wall time → все ядра заняты
    Shuffle Write:         0 B         ← Нет shuffle
    Records Written:       500,000,000

  SQL Metrics:
    number of output rows:          500,000,000
    scan time total:                45s
    metadata time:                  0ms

Stage 1 (Shuffle + Final Agg):
  Duration:              8m 10s
  CPU Time:              45m          ← CPU времени в 5.5× больше wall time
  GC Time:               2m 30s       ← GC: 5% от CPU time - норм
  Shuffle Read:          850 GB       ← Объём shuffle
  Records Read:          500,000,000

CPU Time > Wall time × Core count означает, что все ядра загружены. Если при этом stage медленный - CPU-bound. Варианты диагностики:

  • Высокий GC Time (> 10-15% от CPU Time) → проблема с Java объектами/памятью
  • Низкий GC Time, высокий CPU Time, медленно → скалярное выполнение без SIMD, memory bandwidth bottleneck

Сравнение: Python UDF vs SQL built-in функция

Самый наглядный пример разницы в CPU-эффективности - сравнение Python UDF и эквивалентной SQL-функции:

from pyspark.sql import functions as F
from pyspark.sql.types import DoubleType
import time

# Генерируем датасет: 100 миллионов строк
df = spark.range(100_000_000).withColumn("amount", (F.rand() * 1000).cast("double"))

# Вариант 1: Python UDF (один вызов = одно значение, row-by-row)
@F.udf(returnType=DoubleType())
def tax_python_udf(amount):
    return amount * 0.2 if amount > 100.0 else 0.0

start = time.time()
df.withColumn("tax", tax_python_udf("amount")).agg(F.sum("tax")).show()
print(f"Python UDF: {time.time() - start:.1f}s")

# Вариант 2: Pandas UDF (батч через Arrow)
from pyspark.sql.functions import pandas_udf
import pandas as pd

@pandas_udf(DoubleType())
def tax_pandas_udf(amounts: pd.Series) -> pd.Series:
    return amounts.where(amounts <= 100.0, amounts * 0.2).where(amounts > 100.0, amounts * 0.2).fillna(0.0)
    # Или правильнее:
    # return (amounts * 0.2).where(amounts > 100.0, 0.0)

start = time.time()
df.withColumn("tax", tax_pandas_udf("amount")).agg(F.sum("tax")).show()
print(f"Pandas UDF: {time.time() - start:.1f}s")

# Вариант 3: SQL встроенная функция (JVM native, Catalyst оптимизирует)
start = time.time()
df.withColumn("tax",
    F.when(F.col("amount") > 100.0, F.col("amount") * 0.2).otherwise(0.0)
).agg(F.sum("tax")).show()
print(f"SQL built-in: {time.time() - start:.1f}s")

Типичные результаты на 100M строк:

  • Python UDF: ~90 секунд (row-by-row serialization, Python GIL)
  • Pandas UDF: ~15 секунд (Arrow batch transfer, vectorized numpy)
  • SQL built-in: ~8 секунд (JVM-native execution, Catalyst в Tungsten)

Разница Python UDF vs SQL - более 10×. Это не алгоритмическая разница - логика одинакова. Это чисто накладные расходы на JVM↔Python serialization.

Benchmark: column pruning в Parquet

Демонстрируем, насколько важно читать только нужные колонки:

import time

# Широкая таблица: 200 колонок, 50 миллионов строк
# (Предположим, она уже существует)
df_wide = spark.read.parquet("s3a://bucket/wide_table/")
print(f"Колонок в схеме: {len(df_wide.columns)}")  # 200

# Вариант 1: читаем все 200 колонок
start = time.time()
result1 = df_wide.agg(F.sum("revenue"), F.count("order_id")).collect()
elapsed1 = time.time() - start
print(f"Полный scan (все колонки): {elapsed1:.1f}s")

# Вариант 2: явно выбираем только нужные колонки
start = time.time()
result2 = df_wide.select("revenue", "order_id").agg(F.sum("revenue"), F.count("order_id")).collect()
elapsed2 = time.time() - start
print(f"Column pruning (2 колонки): {elapsed2:.1f}s")

# В EXPLAIN должны видеть разный ReadSchema:
df_wide.select("revenue", "order_id").explain("formatted")
# ReadSchema: struct<revenue:decimal(18,2),order_id:string>  ← только 2 колонки!

Column pruning - это одна из оптимизаций, где Parquet + Catalyst работают вместе эффективно. Spark читает только нужные Row Groups и только нужные Column Chunks. На таблице с 200 колонками и запросом по 2 - экономия I/O в 100 раз.

Включение и диагностика Vectorized Reader

# Проверить, включён ли Vectorized Reader
print(spark.conf.get("spark.sql.parquet.enableVectorizedReader"))  # true по умолчанию

# Если выключен - включить:
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "true")

# Размер батча (по умолчанию 4096 строк)
spark.conf.set("spark.sql.parquet.columnarReaderBatchSize", "4096")

# В EXPLAIN проверяем наличие ColumnarBatch и ColumnarToRow:
df = spark.read.parquet("s3a://bucket/orders/")
df.filter("amount > 0").groupBy("user_id").sum("amount").explain("formatted")

# При включённом Vectorized Reader увидим:
# +- *(1) ColumnarToRow            ← был ColumnarBatch, конвертировали в строки
#    +- FileScan parquet [...]
#       Batched: true              ← подтверждение: читаем батчами

Batched: true в EXPLAIN подтверждает, что Vectorized Reader активен. ColumnarToRow показывает момент, где колоночность теряется. В будущем (с Gluten/Photon) ColumnarToRow должен исчезнуть - данные остаются колоночными до самого результата.


Понимание узких мест: как мыслить об оптимизации

Три профиля аналитического джоба

I/O-bound: время преимущественно тратится на чтение/запись данных.

  • Симптомы: Scan time в Spark UI занимает 70%+ времени стейджа; I/O throughput близок к максимальной скорости дисков/сети
  • Решение: column pruning, partition pruning, форматы с лучшим сжатием (Zstandard), кэширование

Shuffle-bound: время преимущественно тратится на shuffle (Exchange).

  • Симптомы: Shuffle Write/Read в Spark UI огромный (сотни GB); стейдж shuffle занимает 70%+ времени
  • Решение: правильный spark.sql.shuffle.partitions, AQE, bucketing для Shuffle-Free Join, broadcast join для маленьких таблиц

CPU-bound: время преимущественно тратится на вычисления.

  • Симптомы: CPU Time >> Wall Time × Core count; низкий Shuffle/IO, но стейдж медленный; высокий GC time (> 15%)
  • Решение: убрать Python UDF в пользу SQL-функций, Pandas UDF вместо row-UDF, переход на нативный движок (Gluten/Photon)

Большинство реальных джобов - комбинация. Оптимизировать нужно в порядке узких мест: сначала устранить самый большой bottleneck.

Anti-patterns, ухудшающие CPU-эффективность

Python UDF там, где есть SQL-функция: заменяйте udf(lambda x: x.upper()) на F.upper("col").

Широкие таблицы без column pruning: явно выбирайте нужные колонки через .select() перед тяжёлыми операциями.

Row-by-row Pandas операции в Pandas UDF:

# ПЛОХО: row-by-row в pandas (не векторизованно)
@pandas_udf("double")
def bad_udf(s: pd.Series) -> pd.Series:
    result = []
    for v in s:          # Python цикл!
        result.append(v * 0.2 if v > 100 else 0.0)
    return pd.Series(result)

# ХОРОШО: векторизованные pandas операции (numpy под капотом)
@pandas_udf("double")
def good_udf(s: pd.Series) -> pd.Series:
    return (s * 0.2).where(s > 100, 0.0)  # Numpy SIMD-friendly

Отключение Arrow для pandas UDF:

# Убедиться, что Arrow включён:
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")

Взгляд в будущее: вектор развития индустрии

Тренд: columnar-native execution

Индустрия движется к полностью колоночному стеку:

  • Хранение: Parquet, Iceberg, Delta Lake (все колоночные)
  • Память: Apache Arrow (колоночный in-memory стандарт)
  • Выполнение: Velox, Photon, DuckDB, ClickHouse (колоночное выполнение с SIMD)

Spark эволюционирует в том же направлении: Spark 4.x активно развивает поддержку колоночного выполнения, интеграцию с Gluten.

Почему это знание важно для Senior/Staff инженера

Понимание физических лимитов - JVM overhead, SIMD невозможность, UnsafeRow vs ColumnarBatch - позволяет принимать правильные архитектурные решения:

  • Когда выбрать Databricks: нужен Photon? Это конкретное техническое обоснование.
  • Почему переходить на Gluten: не «потому что новое», а потому что конкретный запрос CPU-bound и нативный движок даст 3× ускорение.
  • Почему не писать Python UDF: не «это медленно», а потому что row-by-row serialization убивает SIMD-потенциал и добавляет JVM↔Python overhead.

Эти знания отличают инженера, который знает «что делать», от инженера, который знает «почему это работает так». Последнее - признак уровня Senior и выше.


Домашнее задание

Задача 1: Диагностика CPU-bound джоба

Напишите Spark-запрос, выполняющий тяжёлую агрегацию:

# Создать датасет с 200 миллионами строк
df = (
    spark.range(200_000_000)
    .withColumn("user_id", (F.rand() * 1_000_000).cast("int"))
    .withColumn("amount", (F.rand() * 1000).cast("double"))
    .withColumn("category", (F.rand() * 100).cast("int").cast("string"))
)

# Тяжёлый GROUP BY
result = df.groupBy("user_id", "category").agg(
    F.sum("amount").alias("total"),
    F.avg("amount").alias("avg"),
    F.count("*").alias("cnt")
)
result.write.format("noop").mode("overwrite").save()

В Spark UI для этого джоба найдите и запишите:

  • Общее время выполнения (Wall Time)
  • CPU Time для каждого стейджа
  • GC Time
  • Shuffle Write/Read объём
  • Операторы в физическом плане (explain("formatted"))

Ответьте на вопросы: какой стейдж является узким местом? Это I/O-bound, Shuffle-bound или CPU-bound? Как вы это определили?

Задача 2: Python UDF vs SQL benchmark

Реализуйте одну и ту же функцию тремя способами:

  1. Python udf (скалярный)
  2. Pandas pandas_udf (батч через Arrow)
  3. SQL built-in функции через F.when()

Функция: нормализация суммы - (amount - min) / (max - min), где min и max - константы (пусть будут 0 и 1000). Так что просто amount / 1000.0 с зажимом в [0, 1].

Измерьте время выполнения на 50 миллионах строк. Запишите результаты и объясните разницу через механизм исполнения.

Задача 3: Архитектурное эссе (200-400 слов)

Рассчитайте теоретическую разницу между скалярным и AVX-512 выполнением для SUM(amount) на 1 миллиарде строк (int32):

  • AVX-512 регистр: 512 бит = 16 × int32
  • Скалярное сложение: 1 такт на 1 элемент
  • SIMD сложение: ~1 такт на 16 элементов (в идеале)
  • Тактовая частота: 3 GHz

Вычислите теоретическое время скалярного и SIMD выполнения. Почему реальная разница меньше теоретической? Какие факторы ограничивают SIMD-ускорение (подсказка: memory bandwidth, latency, vector dependency chain)?

Объясните, почему JVM не может воспользоваться этим теоретическим ускорением для типичного Spark GROUP BY запроса.