SIMD и векторные регистры: AVX-512, как процессор обрабатывает Parquet-колонку пачкой за один такт
SIMD и векторные регистры: AVX-512, как процессор обрабатывает Parquet-колонку пачкой за один такт
Физика процессора: от частоты к параллелизму данных¶
Термальный барьер и конец эпохи частот¶
В 1990-е и 2000-е годы производительность процессоров росла предсказуемо: каждые 18 месяцев тактовая частота удваивалась. Intel Pentium работал на 60 МГц в 1993 году, Pentium 4 достиг 3.8 ГГц в 2004 году. Это была эпоха, когда одна программа просто становилась быстрее на новом железе без каких-либо изменений.
В 2004 году всё остановилось. Проблема оказалась физической: при увеличении частоты выше 4 ГГц кремниевый чип начинает выделять такое количество тепла, что никакая система охлаждения не справляется. Это явление называют термальным барьером (thermal wall). Инженеры Intel и AMD столкнулись с фундаментальным физическим ограничением: P = C × V² × f, где P - потребляемая мощность, C - ёмкость цепей, V - напряжение, f - частота. Рост частоты ведёт к квадратичному росту тепловыделения.
Индустрия нашла выход не через рост частоты, а через параллелизм:
- Многоядерность: Intel Core 2 Duo в 2006 году - два ядра на одном чипе. Сегодня серверные процессоры имеют 64-128 ядер.
- Out-of-order execution: процессор сам переставляет независимые инструкции и выполняет несколько за один такт (суперскалярность).
- SIMD (Single Instruction, Multiple Data): одна инструкция обрабатывает несколько значений данных одновременно - параллелизм внутри одного ядра.
Именно SIMD стал ключевым механизмом роста производительности для аналитических (OLAP) рабочих нагрузок. Понимание SIMD - это понимание, почему нативные движки (ClickHouse, DuckDB, Velox, Photon) в 3-6 раз быстрее Spark Tungsten на одинаковом железе.
Архитектура вычислительного конвейера¶
Чтобы понять, как работает SIMD, нужно понять базовый принцип работы процессора. Каждая инструкция проходит через конвейер:
Fetch → Decode → Issue → Execute → Writeback
↑ ↑ ↑ ↑ ↑
Из I-кэша Разбор Очередь ALU/FPU В регистр
байт задач Unit
ALU (Arithmetic Logic Unit) - это «рабочий орган» процессора, который непосредственно выполняет сложение, умножение, сравнение. В современных процессорах есть несколько ALU, работающих параллельно (суперскалярность).
Скалярное выполнение: каждая инструкция обрабатывает одно значение (одно число). ADD EAX, EBX сложит два 32-битных числа. За один такт - одна операция над одной парой чисел.
SIMD: один специализированный блок (vector execution unit) выполняет одну инструкцию, которая применяется одновременно к нескольким парам чисел. VPADDD ZMM0, ZMM1, ZMM2 сложит 16 пар 32-битных чисел за один такт.
Скалярные и векторные регистры: анатомия¶
Скалярные регистры: классика архитектуры x86-64¶
Скалярные регистры - это хранилища для одного значения. В архитектуре x86-64 основные регистры общего назначения:
RAX, RBX, RCX, RDX, RSI, RDI, RSP, RBP - 64-битные (8 байт каждый)
R8, R9, R10, R11, R12, R13, R14, R15 - ещё 8 регистров по 64 бит
EAX - младшие 32 бита RAX (для int32)
AX - младшие 16 бит RAX (для int16)
AL - младшие 8 бит RAX (для byte)
64-битный регистр хранит ровно одно значение: один long, один double, один int32 (в младших 32 битах). Это то, с чем работает обычный Java/Python/C код.
При скалярной обработке массива из N элементов нужно N итераций, N загрузок в регистр, N операций, N сохранений:
; Скалярный цикл: сложение двух массивов a[] + b[] → c[]
; EAX = счётчик (i), RBX = base(a), RCX = base(b), RDX = base(c)
loop_start:
mov eax_temp, [rbx + rax*4] ; загрузить a[i]
add eax_temp, [rcx + rax*4] ; прибавить b[i]
mov [rdx + rax*4], eax_temp ; сохранить c[i]
inc rax ; i++
cmp rax, N ; i < N?
jl loop_start ; повторить
Для массива из 1 000 элементов - ровно 1 000 итераций.
Эволюция SIMD-расширений: от MMX до AVX-512¶
SIMD-расширения появлялись постепенно, каждое поколение расширяло ширину вектора:
| Поколение | Год | Ширина регистра | Регистры | Примечание |
|---|---|---|---|---|
| MMX | 1996 | 64 бит | MM0-MM7 | Только целые, делил FPU |
| SSE | 1999 | 128 бит | XMM0-XMM7 | Добавлены float |
| SSE2 | 2001 | 128 бит | XMM0-XMM15 | Добавлены double, int64 |
| SSE4 | 2007 | 128 бит | XMM0-XMM15 | Новые инструкции |
| AVX | 2011 | 256 бит | YMM0-YMM15 | Расширение XMM до YMM |
| AVX2 | 2013 | 256 бит | YMM0-YMM15 | Целые числа 256-bit |
| AVX-512 | 2016+ | 512 бит | ZMM0-ZMM31 | Маски, gather/scatter |
Каждое следующее поколение расширяло существующие регистры: XMM0 (128 бит) - это младшие 128 бит YMM0 (256 бит), который является младшими 256 битами ZMM0 (512 бит). Это важно: старый SSE-код продолжает работать на новом процессоре.
Векторные регистры: физическая ширина¶
Векторный регистр - это физически более широкое хранилище. Если скалярный RAX = 64 бита, то ZMM0 (AVX-512) = 512 бит. Это не магия - это буквально больше транзисторов и более широкие шины данных в процессоре.
Что можно поместить в один ZMM0 регистр (512 бит):
int8 (1 байт): 64 значения
int16 (2 байта): 32 значения
int32 (4 байта): 16 значений ← типичный INT в Spark
int64 (8 байт): 8 значений ← типичный BIGINT/LONG в Spark
float (4 байта): 16 значений
double (8 байт): 8 значений ← типичный DOUBLE в Spark
В YMM0 (AVX2, 256 бит) - вдвое меньше каждого типа.
Вот что это означает для аналитики: если Spark выполняет SUM(amount) по таблице, где amount - это INT (int32), то при AVX-512:
- Скалярно (Tungsten): 1 значение amount за итерацию
- AVX-512: 16 значений amount за одну инструкцию
Теоретически - 16× ускорение только от этого. Реальный выигрыш меньше (об этом позже), но порядок величин верный.
Как работают SIMD-инструкции: пошаговый разбор¶
Векторная загрузка: поместить массив в регистр¶
Первый шаг - загрузить данные из памяти в векторный регистр. Это требует, чтобы данные лежали последовательно в памяти:
; Загрузить 16 × int32 из памяти по адресу [rsi] в ZMM0
VMOVDQU32 ZMM0, [rsi] ; Unaligned load (данные могут быть по любому адресу)
VMOVDQA32 ZMM0, [rsi] ; Aligned load (адрес должен быть выровнен по 64 байта)
VMOVDQU32 загружает 512 бит = 64 байта = 16 × int32 из непрерывного участка памяти за одну инструкцию. Это и есть ключевое требование: данные должны лежать непрерывно.
Колоночный массив Parquet/Arrow - идеально подходит: все значения одной колонки лежат непрерывно. Строчный формат (UnsafeRow) - плохо подходит: нужное поле разделено другими полями строки.
Векторное сложение: суммирование 16 пар за такт¶
; a[] + b[] → c[], 16 элементов за раз
VMOVDQU32 ZMM0, [rsi] ; загрузить a[0..15] в ZMM0
VMOVDQU32 ZMM1, [rdi] ; загрузить b[0..15] в ZMM1
VPADDD ZMM2, ZMM0, ZMM1 ; ZMM2 = ZMM0 + ZMM1 (поэлементно, 16 × int32)
VMOVDQU32 [rdx], ZMM2 ; сохранить c[0..15] из ZMM2 в память
Что происходит внутри VPADDD:
ZMM0: [a0 | a1 | a2 | ... | a15] ← 16 × 32-bit значений
ZMM1: [b0 | b1 | b2 | ... | b15] ← 16 × 32-bit значений
↓ ↓ ↓ ↓ ← Параллельное сложение за 1 такт
ZMM2: [a0+b0 | a1+b1 | a2+b2 | ... | a15+b15]
Все 16 сложений происходят одновременно, в параллельных ALU-lane внутри vector execution unit.
Полный векторный цикл для сложения двух массивов:
; Предположим N = 1024 элементов
xor rax, rax ; i = 0
loop_avx512:
VMOVDQU32 ZMM0, [rsi + rax*4] ; загрузить a[i..i+15]
VMOVDQU32 ZMM1, [rdi + rax*4] ; загрузить b[i..i+15]
VPADDD ZMM2, ZMM0, ZMM1 ; 16 сложений за 1 такт
VMOVDQU32 [rdx + rax*4], ZMM2 ; сохранить c[i..i+15]
add rax, 16 ; i += 16
cmp rax, N ; i < 1024?
jl loop_avx512 ; повторить
Для 1024 элементов: 1024 / 16 = 64 итерации вместо 1024 скалярных. Ускорение в ~16× только от уменьшения числа итераций.
Векторное умножение и FMA: вычисление revenue¶
Рассмотрим вычисление revenue = price * quantity * (1 - discount) для 8 строк (double):
; Используем AVX-512 с doubles: ZMM вмещает 8 × double
VMOVUPD ZMM0, [price_array] ; price[0..7]
VMOVUPD ZMM1, [quantity_array] ; quantity[0..7]
VMOVUPD ZMM2, [discount_array] ; discount[0..7]
; FMA: ZMM3 = price * quantity (vmulpd)
VMULPD ZMM3, ZMM0, ZMM1 ; ZMM3 = price * quantity, 8 за такт
; Вычислить (1 - discount)
VBROADCASTSD ZMM4, [ONE] ; ZMM4 = [1.0, 1.0, 1.0, 1.0, 1.0, 1.0, 1.0, 1.0]
VSUBPD ZMM5, ZMM4, ZMM2 ; ZMM5 = 1.0 - discount[0..7]
; FMA: revenue = price*quantity * (1-discount)
VMULPD ZMM6, ZMM3, ZMM5 ; ZMM6 = revenue[0..7]
VMOVUPD [revenue_array], ZMM6 ; сохранить
8 строк, 3 умножения на строку = 24 операции. Выполнено за ~4 такта (4 VMOVUPD + 3 VMULPD + 1 VSUBPD). Скалярно это заняло бы 24+ тактов.
Особый класс инструкций - FMA (Fused Multiply-Add): VFMADD231PD ZMM_acc, ZMM_a, ZMM_b выполняет acc += a * b за один такт вместо двух. Это критично для агрегаций: SUM(price * quantity) вычисляется через FMA без промежуточного сохранения.
Идеальный брак: колоночное хранение и SIMD¶
Принцип непрерывности памяти¶
SIMD-загрузка (VMOVDQU32) берёт непрерывный блок из 64 байт (для ZMM). Это абсолютное требование. Рассмотрим два варианта хранения:
Строчный формат (UnsafeRow, как в Tungsten):
Память:
[user_id=1][amount=100][date=18000][status=1] ← строка 0: 4 поля × 4 байт = 16 байт
[user_id=2][amount=200][date=18000][status=1] ← строка 1
[user_id=3][amount=150][date=18001][status=0] ← строка 2
...
Чтобы загрузить 16 значений amount в ZMM0, нужно:
- Взять байты 4-7 из строки 0 (это
amount=100) - Перепрыгнуть 12 байт (user_id + date + status строки 0, user_id строки 1)
- Взять байты 4-7 из строки 1 (это
amount=200) - ...
Данные разбросаны! Одна инструкция VMOVDQU32 загрузит 64 байта подряд, но там окажутся вперемежку user_id, amount, date, status - бесполезно. Нужна специальная инструкция GATHER (о ней позже), которая значительно медленнее обычной загрузки.
Колоночный формат (Arrow, Parquet):
Память (колонка amount):
[100][200][150][300][ 50][400][250][120][600][350][80][450][700][90][180][500]
↑ ↑
Начало колонки 16-й элемент
Все 16 значений amount лежат подряд. VMOVDQU32 загружает все 16 за одну инструкцию. Идеально.
Это не теоретическое рассуждение - это буквально то, почему колоночные форматы были изобретены: они позволяют CPU работать максимально эффективно с аналитическими запросами.
Кэш-иерархия и локальность данных¶
Между процессором и оперативной памятью стоит иерархия кэшей:
Регистры: ~0.1 нс, ~256 байт (сами регистры - буквально внутри ALU)
L1 кэш: ~1 нс, 32-64 KB (на каждое ядро)
L2 кэш: ~4 нс, 256-512 KB (на каждое ядро)
L3 кэш: ~15 нс, 8-64 MB (общий для нескольких ядер)
DRAM RAM: ~60 нс, 64-512 GB
Разница между L1 и DRAM - 60 раз по задержке. Для SIMD, который обрабатывает 16 элементов за такт, эта разница критична.
Cache line - минимальная единица передачи данных между кэшем и памятью: 64 байта. Когда процессор обращается к amount[0], кэш автоматически загружает amount[0..15] (64 байта = 16 × int32). Следующие 15 обращений к amount[1..15] - попадут в кэш бесплатно.
Prefetcher - аппаратный предсказатель: если процессор последовательно читает amount[0], amount[16], amount[32] - prefetcher замечает паттерн и заблаговременно загружает следующие 64-байтные блоки в L1 до того, как они понадобятся. Процессор практически никогда не ждёт данных при последовательном доступе.
При колоночном хранении сканирование amount[0], amount[1], ..., amount[N] - это идеальный последовательный доступ. Prefetcher всегда опережает CPU. Кэш-промахи (cache misses) - минимальны.
При строчном хранении (UnsafeRow) для колонки amount нужно «перескакивать» через другие поля каждой строки. Шаг доступа = ширина строки. Если строка широкая (20+ полей по 8 байт = 160 байт), prefetcher может не успевать, L1 кэш заполнен ненужными полями, эффективная полезная нагрузка кэша снижается.
Пайплайн «без декодирования в строки»¶
В нативных движках (Velox, Photon) пайплайн обработки Parquet-файла выглядит так:
1. Чтение с диска:
Считать Row Group (например, 128 MB Parquet)
2. Векторизованная декомпрессия:
Decompress Column Chunk (Snappy/Zstd/LZ4) → contiguous byte buffer
(Алгоритмы специально оптимизированы под SIMD)
3. Dictionary decode (если словарное кодирование):
Перевести индексы в значения через AVX GATHER
(или оставить как индексы для filtered scan)
4. Прямая загрузка в ZMM-регистр:
VMOVDQU32 ZMM0, [column_buffer + offset]
Данные идут в регистр без конвертации в объекты
5. SIMD-операции:
Filter, Project, Aggregate - прямо над ZMM-регистрами
6. Результат:
Агрегат остаётся в регистрах или записывается в Arrow-буфер
Нет Java-объектов. Нет строчной конвертации. Нет GC. Данные идут с диска в регистры процессора максимально напрямую.
Глубокий разбор: SIMD-фильтрация и агрегация¶
Векторизованная фильтрация: WHERE amount > 100¶
Рассмотрим, как нативный движок выполняет WHERE amount > 100 над 16 значениями int32:
; Входные данные (16 × int32) в ZMM0:
; ZMM0 = [50, 200, 100, 300, 75, 150, 450, 80, 120, 90, 600, 200, 30, 180, 90, 500]
; Загрузить константу 100 во все 16 lane:
VPBROADCASTD ZMM1, [threshold] ; ZMM1 = [100, 100, 100, ..., 100] (16 раз)
; Сравнение 16 элементов одновременно → 16-битная маска:
VPCMPD K1, ZMM0, ZMM1, 6 ; K1 = (ZMM0 > ZMM1) как битовая маска
; Результат K1 = 0b0111010110101110 (1 = прошёл фильтр, 0 = отфильтрован)
; ↑ ↑
; [15] [0]
; VPCOMPRESSD: упаковать прошедшие элементы в начало регистра
VPCOMPRESSD ZMM2 {K1}, ZMM0 ; ZMM2 = [200, 300, 150, 450, 120, 600, 200, 180, 500, ...]
; Непрошедшие элементы удалены
Ключевые детали:
K1 - маска-регистр (mask register). AVX-512 добавил специальные 16-битные регистры K0-K7 для хранения результатов сравнений. Каждый бит маски соответствует одному элементу вектора.
VPCMPD сравнивает 16 пар int32 без ветвлений. Нет if/else, нет branch misprediction. Результат - просто биты в K-регистре. Это branchless execution - принципиальное преимущество SIMD над скалярным if.
VPCOMPRESSD - специальная операция «сжатия»: берёт только те элементы ZMM0, где K1 = 1, и упаковывает их в начало ZMM2. Это аппаратная реализация WHERE условия.
Скалярный эквивалент тех же операций:
# Что делает Tungsten (упрощённо):
result = []
for value in row_batch: # 16 итераций
if value > 100: # branch prediction miss для каждого!
result.append(value)
Vs SIMD: 3 инструкции, ~3 такта для 16 элементов. Скалярно: 16 итераций × (загрузка + сравнение + branch + append) ≈ 64 такта.
Проблема ветвлений: Branch Misprediction¶
Современные процессоры используют предсказание ветвлений (branch prediction): специальный блок предсказывает, будет ли if истинным, и загружает предполагаемую ветвь в конвейер заранее. Если предсказание верное - нет потерь. Если ошибочное - конвейер нужно сбросить (pipeline flush), что стоит 15-20 тактов штрафа.
В аналитических данных ветвления часто непредсказуемы: WHERE amount > 100 - данные могут быть случайными, предиктор угадывает правильно лишь в 50% случаев. Это называется branch misprediction, и на горячих циклах это серьёзный источник потерь производительности.
SIMD полностью устраняет эту проблему: сравнение и маска не требуют ветвлений - это арифметические операции.
Векторизованная агрегация: SUM с горизонтальной редукцией¶
SUM(amount) над колонкой - это операция редукции: взять N элементов, получить 1 сумму. В SIMD это делается в два этапа:
Этап 1: векторизованное накопление (вертикальная сумма):
; Обнуляем аккумулятор
VPXORD ZMM_acc, ZMM_acc, ZMM_acc ; ZMM_acc = [0, 0, 0, ..., 0]
; Цикл: суммируем 16 элементов за раз в 16 параллельных аккумуляторах
loop:
VMOVDQU32 ZMM0, [data + i*4]
VPADDD ZMM_acc, ZMM_acc, ZMM0 ; acc[0..15] += data[i..i+15]
add i, 16
cmp i, N
jl loop
; ZMM_acc = [sum_0, sum_16, sum_32, ...] - частичные суммы
После цикла в ZMM_acc хранятся 16 частичных сумм: ZMM_acc[0] = data[0] + data[16] + data[32] + ..., ZMM_acc[1] = data[1] + data[17] + data[33] + ... и т.д.
Этап 2: горизонтальная редукция (вертикальная сумма → одно число):
; Сложить 16 значений внутри ZMM_acc в одно число
; Используем последовательное сжатие: сложить попарно, повторить
VEXTRACTI64X4 YMM1, ZMM_acc, 1 ; верхняя половина
VPADDD YMM0, YMM_acc, YMM1 ; сложить верхние 8 с нижними 8
; ... повторить для YMM → XMM → горизонтальное сложение
VPHADDW XMM0, XMM0, XMM0
; Итог: XMM0[0] = полная сумма
Горизонтальная редукция требует log2(16) = 4 шагов. Это небольшой overhead, который не влияет на производительность при обработке миллионов элементов.
FMA: Fused Multiply-Add для сложных агрегаций¶
Для агрегаций вида SUM(price * quantity) существует специализированная инструкция FMA:
; VFMADD231PD: ZMM_acc += ZMM_price * ZMM_quantity
; Это НЕ умножение + сложение (2 инструкции)
; Это ОДНА инструкция с одним тактом задержки
VFMADD231PD ZMM_acc, ZMM_price, ZMM_quantity
; Эквивалентно: ZMM_acc = ZMM_acc + (ZMM_price * ZMM_quantity)
; Для 8 doubles за один такт (AVX-512)
FMA важна не только потому, что заменяет две инструкции одной. FMA также более точна (нет промежуточного округления после умножения) и более эффективна по регистрам (не нужен промежуточный регистр для произведения).
Для типичного SQL запроса SELECT SUM(price * quantity) FROM orders:
- Скалярно: N × (загрузка price + загрузка quantity + умножение + сложение с аккумулятором) = 4N операций
- SIMD FMA: N/8 × (загрузка price вектора + загрузка quantity вектора + VFMADD) = 3N/8 операций
Ускорение: 4N / (3N/8) ≈ 10.7×
AVX-512 Throttling: миф и реальность¶
Что такое частотное снижение при AVX-512¶
AVX-512 инструкции задействуют более широкие исполнительные блоки и потребляют больше энергии, чем скалярные. На ранних процессорах (Skylake-X, 2017) это приводило к AVX-512 throttling: при активации AVX-512 инструкций процессор снижал тактовую частоту на всём чипе (иногда до 1.8 ГГц против штатных 3.5-4 ГГц).
Это создало парадоксальную ситуацию: теоретически AVX-512 в 2× быстрее AVX2 (вдвое шире регистр), но практически снижение частоты съедало часть выигрыша. Некоторые инженеры даже утверждали, что AVX-512 медленнее AVX2 - что было ошибочным обобщением.
Современная картина¶
Начиная с Ice Lake (2019) и особенно Sapphire Rapids (2023), ситуация кардинально изменилась:
- Разделение frequency domains: AVX-512 больше не снижает частоту на всём чипе, только на ядрах, активно использующих 512-битные операции
- EVEX encoding оптимизирован: новые процессоры обрабатывают AVX-512 команды эффективнее
- Два AVX-512 порта: Ice Lake и новее имеют два векторных исполнительных порта по 512 бит - теоретически до 2 AVX-512 инструкций за такт
На практике для аналитических рабочих нагрузок (длинные циклы, регулярный доступ к памяти, без крайних тепловых всплесков):
- Skylake-X (старые серверы): AVX-512 даёт +50-80% vs AVX2 с учётом throttling
- Ice Lake / Sapphire Rapids: AVX-512 даёт +80-120% vs AVX2, throttling минимален
- AMD EPYC (Zen 4 / Genoa): AVX-512 с минимальным throttling, конкурентный с Intel
Для Data Engineer практическое следствие: на современных серверных платформах AVX-512 - это реальный выигрыш, а не маркетинговая цифра.
Проверка поддержки SIMD на хосте¶
# Проверить поддержку AVX-512 на Linux
grep -m1 'flags' /proc/cpuinfo | grep -o 'avx512[a-z_]*'
# Типичный вывод для Intel Ice Lake:
# avx512f avx512dq avx512ifma avx512pf avx512er avx512cd avx512bw avx512vl avx512vbmi
# Показать все поддерживаемые расширения
lscpu | grep -i avx
# В Python через numpy:
import numpy as np
# NumPy собирается с поддержкой AVX2/AVX-512 при наличии в системе
np.__config__.show()
# Проверка из PySpark (через sys.platform info):
spark.conf.get("spark.executor.cpu") # Количество CPU на executor
# Получить CPU info через shell:
import subprocess
result = subprocess.run(['lscpu'], capture_output=True, text=True)
print(result.stdout)
Apache Arrow: стандарт SIMD-дружественной памяти¶
Что такое Arrow и почему он важен¶
Apache Arrow - это открытый стандарт колоночного формата для представления данных в оперативной памяти. Ключевое слово - стандарт: Arrow определяет не просто «как хранить», а точный физический layout байт в памяти, одинаковый для C++, Java, Python, Rust, Go.
До Arrow каждая система имела свой внутренний формат: Spark - UnsafeRow, Pandas - NumPy arrays с обёртками, R - SEXP объекты. Передача данных между системами требовала сериализации/десериализации. Arrow устранил это: если две системы понимают Arrow, они могут обмениваться данными через общую память без копирования.
Физический layout Arrow-буферов¶
Arrow буфер для колонки состоит из трёх частей:
Validity Bitmap Buffer (null mask):
[bit0 | bit1 | bit2 | ... | bitN] ← 1 = значение, 0 = NULL
Каждый бит = 1 байт, упакованы 8 бит в байт
Aligned по 64 байта
Values Buffer (для fixed-size типов):
[val0 | val1 | val2 | ... | valN] ← непрерывный массив примитивов
int32: каждый val = 4 байта
int64: каждый val = 8 байт
Aligned по 64 байта (SIMD alignment!)
Offsets Buffer (для variable-size типов, например string):
[offset0 | offset1 | ... | offsetN+1] ← индексы в Data buffer
Data Buffer (для строк и байтовых данных):
[bytes of val0 | bytes of val1 | ...]
Критически важная деталь: выравнивание по 64 байта (aligned to 64 bytes). Это именно размер кэш-строки процессора и AVX-512 регистра. VMOVDQA32 (aligned load) требует выравнивания по 64 байта - Arrow гарантирует это. Aligned load ~15% быстрее unaligned load на типичных рабочих нагрузках.
Arrow Validity Bitmap как SIMD-маска¶
Arrow использует битмап для кодирования NULL-значений. Это не случайно: этот битмап имеет ровно ту структуру, которую нужно для SIMD-фильтрации:
Validity bitmap: [1, 1, 0, 1, 0, 1, 1, 1, 1, 0, 1, 1, 1, 0, 1, 1]
↕
AVX-512 K-mask: bit15..bit0 = 0b1101110110111011
Arrow validity bitmap можно напрямую загрузить как AVX-512 маску K. SIMD-операции с маской автоматически пропустят NULL-значения. Нет лишних ветвлений, нет условных проверок - NULL обрабатываются через маску.
PySpark и Arrow: pandas_udf¶
В Spark Arrow используется для оптимизации PySpark-интеграции:
# Без Arrow: row-by-row pickle serialization (медленно)
@F.udf("double")
def without_arrow(x: float) -> float:
return x * 0.2
# С Arrow: batch transfer через Arrow IPC
@pandas_udf("double")
def with_arrow(x: pd.Series) -> pd.Series:
return x * 0.2 # NumPy vectorized operation на стороне Python
Что происходит при pandas_udf:
- JVM Spark executor собирает batch из N строк в Arrow
RecordBatch(columnar format) - Arrow RecordBatch передаётся в Python-процесс через shared memory (или сокет) - без копирования данных
- Python видит
pd.Seriesbacked by NumPy array, который backed by Arrow buffer x * 0.2выполняется NumPy, который компилируется с AVX2 (или AVX-512 при наличии в системе) - это настоящая векторизация!- Результат передаётся обратно как Arrow buffer
Разница в производительности: 10-100× по сравнению с обычным Python UDF на типичных вычислительных задачах.
# Убедиться, что Arrow включён (включён по умолчанию в Spark 3.x):
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
# Размер batch для Arrow transfer:
spark.conf.set("spark.sql.execution.arrow.maxRecordsPerBatch", "10000")
Архитектура нативных движков: Photon и Velox¶
Общая архитектура нативного execution engine¶
Нативные движки (Photon, Velox, DuckDB) реализуют vectorized execution model - вместо Volcano iterator model (строка за строкой) они передают между операторами векторные батчи (ColumnarBatch / Arrow RecordBatch):
Традиционный Volcano Model (Tungsten):
SortAgg → запрашивает строку у → HashJoin → запрашивает у → Scan
← возвращает 1 строку ← ← возвращает 1 строку ←
Vectorized Model (Photon / Velox):
SortAgg ← получает batch 4096 строк ← HashJoin ← получает batch ← Scan
одна передача = 4096 строк одна передача = 4096 строк
Разница: при Volcano Model операторы вызываются N раз (по одному на строку). При vectorized model - N/4096 раз (по одному на батч). Overhead на вызов функции и пересылку данных уменьшается в 4096 раз.
Databricks Photon¶
Photon - это закрытый нативный движок от Databricks, написанный на C++. Он является partial replacement для Spark Tungsten: выбранные операторы физического плана заменяются нативными реализациями.
Что умеет Photon (поддерживаемые операторы):
- Scan: чтение Delta/Parquet без конвертации в строки, прямо в columnar format
- Filter: SIMD-фильтрация через маски
- Project: колоночные проекции и вычисления
- HashAggregate: векторизованная хэш-агрегация
- HashJoin: векторизованный hash join
- Sort/TopK: векторизованная сортировка
- Exchange: columnar shuffle через Arrow
Что ещё НЕ в Photon (выполняется Tungsten):
- Python UDF / Scala UDF (нет обхода: UDF - это пользовательский код)
- Некоторые оконные функции
- Сложные nested типы (MapType, complex ArrayType)
Когда Photon «берёт» оператор - в физическом плане появляется Photon*:
== Physical Plan ==
PhotonResultStage
+- PhotonProject [user_id, sum(amount)]
+- PhotonHashAggregate [user_id, sum(amount)]
+- PhotonExchange hashpartitioning(user_id, 200)
+- PhotonHashAggregate [user_id, partial_sum(amount)]
+- PhotonFilter (amount > 100)
+- PhotonScan parquet [user_id, amount] ...
Весь план выполняется в C++ с SIMD. Данные покидают C++ только в самом конце для вывода результата.
Apache Velox: Meta's open-source engine¶
Velox (Meta, 2022, open-source под Apache License) - это C++ библиотека векторизованного выполнения запросов. Не самостоятельная СУБД - именно библиотека, которую встраивают в другие системы:
- Presto (Meta's query engine): Presto использует Velox как native execution backend
- Spark via Apache Gluten: Gluten подменяет Tungsten на Velox
- PyTorch ExecuTorch: Velox для ML workloads
Ключевые свойства Velox:
- Apache Arrow как основной формат данных in-memory
- SIMD-первичность: каждый оператор оптимизирован под AVX2/AVX-512
- Adaptive operator fusion: несколько операторов объединяются в pipeline без промежуточных материализаций
- Dictionary encoding awareness: если данные Dictionary-encoded, Velox может фильтровать по словарю без декодирования
Apache Gluten: мост между Spark и Velox¶
Apache Gluten (incubating, 2022) - плагин для Apache Spark, который заменяет Tungsten физические операторы на Velox (или ClickHouse) реализации.
Обычный Spark:
SQL → Catalyst (JVM) → Physical Plan → Tungsten (JVM) → Result
Spark + Gluten:
SQL → Catalyst (JVM) → Physical Plan → Gluten Plugin → Velox (C++) → Result
↕
Substrait plan format
(common intermediate representation)
Substrait - это portable intermediate representation для query plans, аналог LLVM IR для компиляторов. Gluten конвертирует Spark physical plan в Substrait, Velox выполняет Substrait plan.
Для пользователя - ничего не меняется:
# Тот же код
df = spark.read.parquet("s3a://bucket/data/")
result = df.filter("amount > 100").groupBy("user_id").agg(F.sum("amount"))
result.show()
# Но в логах:
# Gluten: offloading plan to Velox backend
# Velox: executing with AVX-512 vectorization
Диаграмма: путь от Parquet до регистра процессора¶
Схема показывает полный путь данных в нативном движке: от сжатого Parquet на диске через декомпрессию, dictionary decode, Arrow-буфер до загрузки в ZMM-регистр. Каждый переход - с минимальными копированиями. SIMD-операции (Filter, Aggregation) работают непосредственно над данными в векторных регистрах. Никаких Java-объектов, никакого GC - только bytes, registers и CPU execution units.
Бенчмарки: сколько стоит векторизация¶
TPC-H: стандартный аналитический benchmark¶
TPC-H - индустриальный benchmark для OLAP, состоящий из 22 сложных SQL-запросов над синтетическими данными (shipping, orders, lineitem). Используется для сравнения движков.
Типичные результаты на TPC-H SF=100 (100 GB данных):
| Движок | Q1 (SUM аггр.) | Q6 (фильтр+SUM) | Q12 (JOIN+GROUP) | Геом. среднее |
|---|---|---|---|---|
| Spark 3.x Tungsten | 45 сек | 38 сек | 95 сек | ~55 сек |
| Spark + Photon (Databricks) | 8 сек | 6 сек | 18 сек | ~10 сек |
| Spark + Gluten/Velox | 12 сек | 9 сек | 25 сек | ~14 сек |
| DuckDB (нативный) | 5 сек | 4 сек | 12 сек | ~7 сек |
| ClickHouse (нативный) | 3 сек | 2 сек | 8 сек | ~4 сек |
Ускорение Photon vs Tungsten: ~5.5× в среднем. Это при одинаковом железе, одинаковых данных, одинаковых SQL-запросах - разница только в execution engine.
Понимание источников ускорения¶
Ускорение 4-6× складывается из нескольких факторов:
1. SIMD параллелизм (2-4×): обработка 8-16 элементов за такт вместо 1. На чистых агрегациях без сложных условий - близко к теоретическому максимуму.
2. Columnar execution (1.5-2×): данные не конвертируются в строчный UnsafeRow. Лучшая кэш-локальность, эффективное prefetching.
3. Отсутствие JVM overhead (1.2-1.5×): нет GC пауз, нет object headers, нет bounds checks.
4. Branchless execution (1.2-1.5×): предикаты через SIMD-маски вместо условных переходов.
5. SIMD-дружественные алгоритмы (1.2-2×): hash aggregation, sort, join заново написаны под векторизацию - не просто применён SIMD к старым алгоритмам.
Произведение факторов: 4 × 1.7 × 1.3 × 1.3 × 1.5 ≈ 17×. Почему реально 5-6×? Потому что не весь план векторизован (некоторые операторы в JVM), и потому что узкое место смещается на память (memory bandwidth bottleneck) при достижении CPU-предела.
Практика: анализ SIMD-потенциала в вашем кластере¶
Проверка поддержки векторизации в Spark¶
# 1. Проверить Vectorized Reader
print("Vectorized Parquet Reader:",
spark.conf.get("spark.sql.parquet.enableVectorizedReader")) # true
print("Vectorized ORC Reader:",
spark.conf.get("spark.sql.orc.enableVectorizedReader")) # true
print("Arrow enabled:",
spark.conf.get("spark.sql.execution.arrow.pyspark.enabled")) # true
# 2. Размер batch для Vectorized Reader
print("Batch size:",
spark.conf.get("spark.sql.parquet.columnarReaderBatchSize")) # 4096
# 3. Проверить через EXPLAIN наличие columnar операторов
df = spark.read.parquet("s3a://bucket/large_table/")
df.filter("amount > 100").groupBy("user_id").sum("amount").explain("formatted")
В EXPLAIN ищите:
== Physical Plan ==
*(2) HashAggregate ...
+- Exchange hashpartitioning ...
+- *(1) HashAggregate ...
+- *(1) Filter amount > 100
+- *(1) ColumnarToRow ← здесь заканчивается columnar
+- FileScan parquet ... Batched: true ← Vectorized Reader активен
Batched: true - подтверждение, что Vectorized Reader работает. ColumnarToRow - граница, после которой начинается построчный UnsafeRow. Именно этот шаг нативные движки устраняют.
Benchmark: columnar vs row operations¶
import time
# Создать датасет: 100M строк, типичная transaction таблица
n = 100_000_000
df = (
spark.range(n)
.withColumn("user_id", (F.rand() * 1_000_000).cast("int"))
.withColumn("amount", (F.rand() * 10_000).cast("int"))
.withColumn("quantity", (F.rand() * 100).cast("int"))
.withColumn("category", (F.rand() * 50).cast("int").cast("string"))
.withColumn("is_valid", (F.rand() > 0.1).cast("boolean"))
)
# Записать в Parquet для реалистичного теста
df.write.mode("overwrite").parquet("/tmp/benchmark_data/")
df_parquet = spark.read.parquet("/tmp/benchmark_data/")
# Тест 1: простая агрегация (SIMD-friendly: чистый SUM)
start = time.time()
df_parquet.filter("is_valid = true") \
.groupBy("category") \
.agg(F.sum("amount").alias("total"), F.count("*").alias("cnt")) \
.write.format("noop").mode("overwrite").save()
t1 = time.time() - start
print(f"SQL aggregation (SIMD-friendly): {t1:.1f}s")
# Тест 2: та же логика через Python UDF (разрушает SIMD)
@F.udf("long")
def py_amount_valid(amount, is_valid):
return amount if is_valid else 0
start = time.time()
df_parquet.withColumn("amount_clean", py_amount_valid("amount", "is_valid")) \
.groupBy("category") \
.agg(F.sum("amount_clean").alias("total"),
F.count("*").alias("cnt")) \
.write.format("noop").mode("overwrite").save()
t2 = time.time() - start
print(f"Python UDF aggregation (SIMD destroyed): {t2:.1f}s")
print(f"Speedup SQL vs UDF: {t2/t1:.1f}x")
# Тест 3: pandas_udf (частично восстанавливает batching)
@pandas_udf("long")
def pandas_amount_valid(amount: pd.Series, is_valid: pd.Series) -> pd.Series:
return amount.where(is_valid, 0)
start = time.time()
df_parquet.withColumn("amount_clean", pandas_amount_valid("amount", "is_valid")) \
.groupBy("category") \
.agg(F.sum("amount_clean").alias("total"),
F.count("*").alias("cnt")) \
.write.format("noop").mode("overwrite").save()
t3 = time.time() - start
print(f"Pandas UDF aggregation (batch Arrow): {t3:.1f}s")
Ожидаемые результаты:
- SQL: ~15-25 секунд
- Python UDF: ~120-180 секунд (~8-10× медленнее)
- Pandas UDF: ~25-40 секунд (~1.5-2× медленнее SQL)
Анализ Spark UI для выявления CPU bottleneck¶
После запуска benchmark открываем Spark UI:
SQL Tab → click on query → click on Stage → Task Metrics
Для SQL aggregation:
Task Metrics:
Duration: 22s
GC Time: 0.5s ← < 5% от total, хорошо
CPU Time: 176s ← 8 cores × 22s = 176s (все ядра загружены)
Records Read: 100,000,000
Bytes Read: 3.2 GB
SQL Metrics:
Batched: true ← Vectorized reader работает
scan time total: 8s
filter time: 5s
aggregate time: 9s
Для Python UDF:
Duration: 148s
GC Time: 12s ← ~8%, начинает давить
CPU Time: 1184s ← 8 × 148s = все ядра заняты, но медленно
Shuffle bytes: 2.1 GB
Task deserialization: 45s ← ОГРОМНО: это Python overhead!
Task deserialization time - время на JVM↔Python сериализацию. 45 секунд из 148 - это 30% времени только на переброску данных между JVM и Python. Именно это убивает производительность Python UDF.
SIMD limitations: когда векторизация не помогает¶
Нерегулярный доступ к памяти: gather/scatter¶
SIMD load (VMOVDQU32) требует непрерывных данных. Но hash join, shuffle, random access по ключу требуют «подбора» данных из произвольных адресов. Для этого существует инструкция GATHER:
; VGATHERDPS: загрузить floats по индексам в ZMM_idx
; ZMM_idx = [i0, i1, i2, ..., i15] - индексы
; base_addr = начало массива
VGATHERDPS ZMM_result, [base_addr + ZMM_idx * 4]
; Результат: ZMM_result[k] = *(base_addr + ZMM_idx[k] * 4)
GATHER - медленнее обычного load в 3-5×. Он всё равно лучше, чем 16 последовательных скалярных загрузок (особенно если данные в кэше), но далеко от идеала.
Для операций типа JOIN (random access в probe table) даже нативные движки не достигают полного SIMD-ускорения - данные поступают нерегулярно.
Сложные ветвления и вложенные структуры¶
Вложенные JSON, ArrayType, MapType - всё это требует обработки переменной длины и условной логики. SIMD плохо справляется с переменной длиной (нельзя загрузить 16 строк разной длины в ZMM регистр).
# SIMD-friendly: фиксированные типы
df.select("user_id", "amount", "quantity") # int32, int32, int32 - отлично
# SIMD-unfriendly: переменная длина
df.select(F.get_json_object("payload", "$.user_id")) # Строки → плохо
df.select(F.explode("items_array")) # Array → переменная длина → плохо
Когда Python UDF неизбежен¶
Если бизнес-логика требует Python-кода (ML-модели, сложные алгоритмы, внешние библиотеки), Python UDF неизбежен. В этом случае:
- Всегда используйте
pandas_udfвместо обычногоudf- батч через Arrow в 5-20× быстрее - Делайте максимум операций в SQL до вызова UDF (predicate pushdown, projection)
- Вызывайте UDF как можно реже, на как можно меньшем количестве строк
# ПЛОХО: UDF применяется ко всем строкам до фильтрации
df.withColumn("score", complex_ml_udf("features")) \
.filter("score > 0.8") # фильтруем ПОСЛЕ UDF
# ХОРОШО: фильтруем ДО UDF, UDF применяется к меньшему числу строк
df.filter("amount > 100 AND status = 'active'") \ # SQL filter сначала (SIMD!)
.withColumn("score", complex_ml_udf("features")) # UDF только для отобранных строк
Anti-patterns и best practices¶
Anti-pattern 1: строчное мышление в аналитике¶
# ПЛОХО: обработка строк по одной через Python
for row in df.collect(): # collect() переносит данные в driver
process_row(row) # Python цикл без параллелизма
# ХОРОШО: операции над колонками
df.withColumn("result", F.col("amount") * F.col("quantity")) # Columnar, SIMD-friendly
Правило: никогда не использовать df.collect() + Python цикл для трансформаций. Это не только медленно из-за отсутствия SIMD - это переносит терабайты данных на driver и разрушает распределённость.
Anti-pattern 2: JSON везде¶
# ПЛОХО: хранить структурированные данные как JSON строки
# Колонка payload: '{"user_id": 123, "amount": 456.78, "category": "electronics"}'
df.select(
F.get_json_object("payload", "$.user_id").cast("int"),
F.get_json_object("payload", "$.amount").cast("double")
)
# → Каждая операция парсит JSON строку: нет SIMD, нет column pruning
# ХОРОШО: распаковать JSON при ingestion, хранить как struct
df.select(
F.from_json("payload", schema).alias("data")
).select("data.user_id", "data.amount", "data.category")
# → После распаковки: типизированные колонки → SIMD-friendly
JSON-парсинг - принципиально не SIMD-friendly: строки переменной длины, сложный синтаксис, условная обработка. Распаковывайте JSON на Bronze→Silver этапе, храните в нативных типах.
Anti-pattern 3: ненужные wide rows¶
# ПЛОХО: читать все 200 колонок для агрегации по 3
df_200_cols = spark.read.parquet("s3a://bucket/wide_table/")
df_200_cols.groupBy("user_id").sum("amount").show()
# Spark прочитает все 200 колонок! Хотя нужны только 2.
# ХОРОШО: column pruning через явный select
df_200_cols.select("user_id", "amount") \
.groupBy("user_id").sum("amount").show()
# Spark прочитает только 2 колонки из Parquet
Column pruning в Parquet работает автоматически через Catalyst, если вы явно указываете нужные колонки. Но если между read и groupBy есть withColumn или UDF, которые требуют всех колонок - pruning не применится.
Best practices: максимизация SIMD в Spark¶
# 1. Используйте встроенные SQL-функции, не UDF
spark.sql("SELECT SUM(price * qty * (1 - discount)) FROM orders") # JVM-native, SIMD-ready
# 2. Приводите типы явно
df.withColumn("amount", F.col("amount").cast("int")) # int32: 16 за такт
# vs F.col("amount").cast("long") # int64: 8 за такт
# 3. Включите AQE для автоматической оптимизации
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
# 4. Убедитесь в Vectorized Reader
spark.conf.set("spark.sql.parquet.enableVectorizedReader", "true")
spark.conf.set("spark.sql.parquet.columnarReaderBatchSize", "4096")
# 5. Включите Arrow для pandas UDF
spark.conf.set("spark.sql.execution.arrow.pyspark.enabled", "true")
# 6. Предпочитайте Parquet/ORC/Delta/Iceberg CSV/JSON для аналитики
# Parquet = columnar, compressed, SIMD-friendly
# JSON = row-oriented, uncompressed, SIMD-unfriendly
Домашнее задание¶
Задача 1: Расчёт теоретического ускорения¶
Дана таблица orders с 500 миллионами строк. Нужно выполнить:
SELECT category, SUM(amount), COUNT(*), AVG(quantity)
FROM orders
WHERE status = 'completed' AND amount > 0
GROUP BY category
Колонки: amount (int32), quantity (int32), status (string, dictionary encoded), category (string, dictionary encoded 50 уникальных значений).
Рассчитайте:
- Сколько элементов
amountобработает AVX-512 за один такт? А AVX2? - Сколько теоретически тактов нужно на фильтрацию
amount > 0для 500M элементов через SIMD vs скалярно? - Если процессор работает на 3 ГГц - сколько секунд займёт фильтрация в обоих случаях теоретически?
- Почему реальное время будет больше теоретического? Назовите 3 фактора.
Задача 2: Профилирование Spark-запроса¶
Запустите два варианта одного запроса:
# Вариант A: SQL built-in functions
df.filter("amount > 100 AND status = 'completed'") \
.groupBy("user_id") \
.agg(F.sum(F.col("amount") * F.col("quantity")).alias("revenue")) \
.write.format("noop").mode("overwrite").save()
# Вариант B: эквивалентный Python UDF
@F.udf("long")
def calc_revenue(amount, quantity, status):
if status == 'completed' and amount > 100:
return amount * quantity
return 0
df.withColumn("revenue", calc_revenue("amount", "quantity", "status")) \
.groupBy("user_id") \
.agg(F.sum("revenue")) \
.write.format("noop").mode("overwrite").save()
Для каждого варианта запишите из Spark UI:
- Wall time (общее время)
- CPU Time
- GC Time
- Task deserialization time (если есть)
- Операторы в Physical Plan
Ответьте: в чём конкретно разница? Какой метрикой подтверждается overhead Python UDF?
Задача 3: Анализ CPI (Cycles Per Instruction)¶
Концептуальный вопрос. Идеальный аналитический движок с AVX-512 выполняет SUM(amount) с CPI ≈ 0.5 (то есть выдаёт 2 результата на такт через суперскалярность). Spark Tungsten в скалярном режиме показывает CPI ≈ 3-4 на той же операции (такты на 1 результат, с учётом cache misses и overhead).
Объясните:
- Почему CPI нативного движка < 1 возможен (что такое суперскалярность + SIMD)?
- Откуда берётся CPI = 3-4 в Tungsten (назовите 3 компонента overhead)?
- Если переход с Tungsten на Photon даёт CPI с 3.5 до 0.6 - какой процент ускорения это даёт по времени?
Что сдавать¶
- Расчёты из задачи 1 с объяснением
- Скриншоты Spark UI из задачи 2 с выделенными ключевыми метриками
- Ответы на вопросы задачи 3 с математикой
Задача со звёздочкой: запустите perf stat -e instructions,cycles,cache-misses,cache-references java -jar your-spark-app.jar (Linux perf tool) для обоих вариантов из задачи 2 и сравните instructions per cycle и cache-miss rate. Это профилирование уровня hardware performance counters.