Проблема JVM: почему Tungsten не может использовать SIMD и при чём тут колоночное хранение
Проблема JVM: почему Tungsten не может использовать SIMD и при чём тут колоночное хранение
Исторический контекст: эволюция узких мест в 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:
- Spark читает Parquet: получает только колонку
amountкак непрерывный массив[100, 200, 150, 300, 50, ...](отлично для SIMD!) - Spark конвертирует в UnsafeRow: каждое значение
amountупаковывается в строку[null_bitmap | amount_value](теряем непрерывность колонки!) - 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¶
Реализуйте одну и ту же функцию тремя способами:
- Python
udf(скалярный) - Pandas
pandas_udf(батч через Arrow) - 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 запроса.