SIMD и векторные регистры: AVX-512, как процессор обрабатывает Parquet-колонку пачкой за один такт

SIMD и векторные регистры: AVX-512, как процессор обрабатывает Parquet-колонку пачкой за один такт

optimization

Физика процессора: от частоты к параллелизму данных

Термальный барьер и конец эпохи частот

В 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:

  1. JVM Spark executor собирает batch из N строк в Arrow RecordBatch (columnar format)
  2. Arrow RecordBatch передаётся в Python-процесс через shared memory (или сокет) - без копирования данных
  3. Python видит pd.Series backed by NumPy array, который backed by Arrow buffer
  4. x * 0.2 выполняется NumPy, который компилируется с AVX2 (или AVX-512 при наличии в системе) - это настоящая векторизация!
  5. Результат передаётся обратно как 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 уникальных значений).

Рассчитайте:

  1. Сколько элементов amount обработает AVX-512 за один такт? А AVX2?
  2. Сколько теоретически тактов нужно на фильтрацию amount > 0 для 500M элементов через SIMD vs скалярно?
  3. Если процессор работает на 3 ГГц - сколько секунд займёт фильтрация в обоих случаях теоретически?
  4. Почему реальное время будет больше теоретического? Назовите 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).

Объясните:

  1. Почему CPI нативного движка < 1 возможен (что такое суперскалярность + SIMD)?
  2. Откуда берётся CPI = 3-4 в Tungsten (назовите 3 компонента overhead)?
  3. Если переход с Tungsten на Photon даёт CPI с 3.5 до 0.6 - какой процент ускорения это даёт по времени?

Что сдавать

  1. Расчёты из задачи 1 с объяснением
  2. Скриншоты Spark UI из задачи 2 с выделенными ключевыми метриками
  3. Ответы на вопросы задачи 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.