Gluten + Velox: Substrait IR и offloading вычислений в C++ Velox

Gluten + Velox: Substrait IR и offloading вычислений в C++ Velox

optimization

Концепция проекта Gluten: вынос вычислений за рамки JVM

Чтобы понять, зачем вообще существует Gluten, нужно признать фундаментальный факт: JVM - отличная платформа для написания распределённых систем, но плохая платформа для исполнения тяжёлых аналитических вычислений на современном железе. Spark был написан на JVM-языках (Scala, Java) - и это дало ему богатую экосистему, простоту развёртывания и десятки тысяч контрибьюторов. Но одновременно это оковы: JVM не умеет напрямую управлять SIMD-регистрами процессора, её сборщик мусора вносит непредсказуемые паузы при работе с большими датасетами, а объектная модель Java требует многократного разыменования указателей там, где C++ работает с линейными массивами в памяти.

Идея «разделения движков»

Одна крайность - переписать Spark целиком на C++ или Rust. Именно так сделан Sail (Rust + DataFusion + Spark Connect): он отказывается от всего JVM-слоя. Это радикально, но влечёт годы работы и неполную совместимость с существующими API.

Другая крайность - оставить всё как есть. Тогда всё, что Tungsten и WholeStageCodegen уже дают - это потолок. Он не низкий, но недостаточен для конкуренции с нативными движками ClickHouse, DuckDB или Presto/Trino.

Gluten выбирает средний путь: оставить Spark целиком - как оркестратор, планировщик, API-слой и средство общения с кластером - но заменить только ту часть, где JVM наиболее неэффективен: физическое исполнение операторов на уровне Executor. Spark по-прежнему:

  • принимает Python/SQL/DataFrame-запросы от пользователя
  • строит логический план через Catalyst optimizer
  • применяет AQE (Adaptive Query Execution)
  • управляет стейджами, задачами, Executor-ами
  • следит за отказоустойчивостью (retry, speculation)

Но в момент, когда нужно физически исполнить HashAggregate или SortMergeJoin над 100 миллиардами строк - Gluten перехватывает управление и передаёт вычисления в нативный C++ движок.

Роли участников: Gluten и Velox

В этом уравнении участвуют два совершенно разных проекта:

Apache Gluten - это мост: Java/C++ фреймворк, который:

  • интегрируется со Spark через Plugin API (SparkPlugin + SparkSessionExtensions)
  • перехватывает физический план Catalyst
  • транслирует его в стандартный промежуточный формат (Substrait IR)
  • передаёт этот IR через JNI в нативный backend
  • получает обратно результат в виде Arrow-совместимых колоночных батчей
  • умеет работать с разными нативными backends: Velox, ClickHouse, других

Velox - это вычислительный движок: монолитная высокопроизводительная C++ библиотека, созданная в Meta (Facebook) в 2019 году как единый execution runtime для Presto, Spark, Spark Connect и других аналитических систем в инфраструктуре Meta. Velox:

  • не является отдельным сервером или приложением - это библиотека, которую можно встроить в любой процесс
  • реализует полный стек векторизованных операторов: Filter, Project, HashAggregate, SortMergeJoin, BroadcastHashJoin, Window, Sort, OrderBy
  • управляет нативной памятью через собственный Memory Pool - без Java GC
  • использует SIMD-инструкции AVX2/AVX-512 через intrinsics в C++
  • поддерживает нативный сброс на диск (spilling) при нехватке RAM
  • открыт в Open Source под Apache 2.0

Почему Gluten + Velox - «хардкорное» решение

В сообществе сложился устойчивый взгляд на иерархию нативных ускорителей Spark:

Решение Backend Язык Сложность деплоя Потенциал ускорения
Apache Comet DataFusion Rust Низкая 2–4×
Gluten + Velox Velox C++ Высокая 3–8×
Photon (Databricks) Проприетарный C++ C++ Только облако 3–10×
Sail DataFusion Rust Средняя 3–6×, неполная совместимость

Gluten + Velox даёт наибольший потенциал ускорения среди open-source решений, но требует самого глубокого понимания системы: нужно правильно настроить нативную память, понять ограничения Substrait-трансляции и обеспечить сборку нативных бинарей под конкретную архитектуру CPU.


Сердце интеграции: что такое Substrait IR

Проблема Вавилонской башни

Представьте задачу: у вас есть сложный plan запроса, представленный в виде Java-объектов в JVM - SortMergeJoinExec, HashAggregateExec, FilterExec. Как передать логику этого плана в C++ процесс, который ничего не знает о Java?

Наивный подход - написать прямую конвертацию из каждого Java-оператора в соответствующий C++ вызов. Это работает, пока операторов десять. Но Spark содержит сотни вариантов физических планов, учитывает типы данных, null-семантику, порядок колонок, аннотации AQE - и прямая конвертация превращается в проект на несколько лет, который при каждом обновлении Spark нужно переписывать заново.

Именно эту проблему решает Substrait - открытый стандарт описания реляционных планов запросов, не привязанный ни к одному языку или движку. Это аналог LLVM IR в мире компиляторов: LLVM позволяет любому языку (C, C++, Rust, Swift) компилироваться в одно и то же промежуточное представление - и исполняться на любой архитектуре. Substrait делает то же самое для аналитических запросов.

Что такое Substrait

Substrait (substrait.io) - это открытая спецификация (Apache 2.0), описывающая реляционную алгебру в виде Protocol Buffers-схемы. Она определяет:

  • Relations (отношения): Scan, Filter, Project, Aggregate, Join, Sort, Limit - все стандартные операции над таблицами
  • Expressions: арифметика, сравнения, функции, CASE WHEN, подзапросы
  • Types: полная система типов (int8–int64, float32/64, string, decimal, date, timestamp, list, map, struct)
  • Functions: стандартные функции (sum, count, substring, etc.) с точным описанием семантики

Ключевое свойство Substrait: план запроса сериализуется в бинарный формат (Protocol Buffers) - компактный, быстрый для сериализации/десериализации, языконезависимый. Этот бинарный blob можно передать через JNI между JVM и нативным C++ кодом - и на другой стороне его десериализуют за микросекунды.

Цепочка трансляции Gluten

Подробно разберём каждый шаг этой цепочки.

Шаг 1 - Spark Catalyst строит физический план. Пользователь пишет df.groupBy("region").agg(F.sum("amount")). Catalyst преобразует это в логический план, оптимизирует его и строит физический план из стандартных JVM-операторов: HashAggregateExec(ScanExec(...)).

Шаг 2 - Gluten перехватывает план. ColumnarRule, зарегистрированный через CometSparkSessionExtensions, обходит каждый узел физического дерева и проверяет, может ли он быть трансформирован в Substrait. Если да - помечает его для offloading. Если нет - оставляет как JVM-оператор с fallback.

Шаг 3 - Трансляция в Substrait. Gluten обходит дерево поддерживаемых операторов и сериализует их в Substrait Protocol Buffers. Каждый оператор Spark имеет соответствующее Relation в Substrait:

HashAggregateExec(grouping=[region], funcs=[sum(amount)])
    ↓
AggregatRelation {
    groupings: [FieldReference{field=0}],  // колонка region
    measures: [AggregateFunction{name="sum", args=[FieldReference{field=1}]}]
    input: FilterRelation { ... }
}

Шаг 4 - JNI-вызов. Бинарный protobuf-blob (byte[]) передаётся в нативный C++ код через JNI. JNI работает в пределах одного OS-процесса - это не IPC, не сетевой вызов, это вызов функции в другой части памяти того же процесса.

Шаг 5 - Velox десериализует и строит нативный план. На C++ стороне Substrait-blob десериализуется в объекты Velox PlanNode. Это дерево нативных C++ объектов, полностью аналогичных Spark-плану по логике, но реализованных через Velox-операторы.

Шаг 6 - Векторизованное выполнение. Velox исполняет план пакетами (RowVector - аналог Arrow RecordBatch), применяя SIMD-инструкции к каждому батчу колоночных данных.

Шаг 7 - Возврат результата. Результирующие RowVector конвертируются в Arrow-совместимые ColumnarBatch и возвращаются в JVM через JNI. Spark получает результат как обычный набор ColumnarBatch и продолжает оркестрацию.

Почему Substrait - правильный выбор

До Substrait каждый нативный ускоритель писал свою собственную трансляцию Spark → C++. Comet писал свою (protobuf для DataFusion). Gluten с ClickHouse-бэкендом - свою. Результат: несовместимые решения, которые невозможно переиспользовать.

Substrait создаёт общий язык между движками. Если вы описали запрос в Substrait, вы можете исполнить его в Velox, DataFusion, ClickHouse, DuckDB - в любом движке, который реализует Substrait-reader. Это не просто удобство - это архитектурное решение, которое позволяет Gluten подключать разные нативные backends (spark.gluten.sql.columnar.backend.lib=velox или =ch для ClickHouse) без изменения логики трансляции.


Нативная магия Velox: оптимизация на уровне железа

Почему JVM execution - не оптимум для аналитики

Прежде чем разобраться, почему Velox быстрый, нужно понять, где именно JVM теряет производительность на аналитических задачах.

Проблема 1: Объектная модель Java. Запись в Spark UnsafeRow - это область памяти, но добраться до неё нужно через несколько уровней разыменования указателей. При агрегации миллиарда строк CPU тратит значительную часть времени на cache miss - данные находятся не в L1/L2 кэше, а в основной RAM. Velox хранит колонки как плоские C++ массивы (int64_t* data) - одно разыменование, максимальная локальность кэша.

Проблема 2: Java GC. При обработке больших join-операций Spark создаёт миллионы промежуточных Java-объектов для хэш-таблиц, буферов сортировки. GC периодически останавливает все потоки на сборку (STW - stop-the-world паузы). Velox управляет памятью через свой MemoryPool - детерминированные аллокации/освобождения без GC.

Проблема 3: JIT vs. SIMD. JVM JIT-компилятор умеет генерировать SIMD-инструкции, но это происходит недетерминированно и зависит от профиля выполнения, версии JVM, паттернов доступа к памяти. C++-компилятор (Clang/GCC) с флагами -O3 -mavx512f генерирует предсказуемые SIMD-инструкции для любого кода, написанного с использованием intrinsics.

SIMD: обработка данных пачками на уровне процессора

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

Что это даёт на практике: регистр ZMM (512 бит) вмещает 8 значений double (64-бит каждое) или 16 значений float (32-бит). Одна инструкция _mm512_add_pd складывает сразу 8 пар double. Для SUM(amount) по батчу из 8192 строк это означает: не 8192 последовательных сложений, а 8192 / 8 = 1024 векторных инструкции - примерно в 8 раз меньше тактов CPU.

Реальное ускорение меньше теоретического (3–6× вместо 8×) из-за накладных расходов на загрузку/выгрузку данных из регистров, проверки условий и ветвления для null-значений. Но даже 3–6× - это огромный выигрыш на CPU-bound задачах.

Velox использует SIMD через C++ intrinsics - типобезопасные обёртки над ассемблерными инструкциями:

// Суммирование колонки amount (double) в батче с 8 элементами за такт
// Эквивалент: SELECT SUM(amount) FROM batch

void sum_double_avx512(const double* amounts, int64_t n_rows,
                       const uint64_t* validity_bitmap,
                       double* result) {
    __m512d acc = _mm512_setzero_pd();   // обнуляем 512-бит аккумулятор

    for (int64_t i = 0; i < n_rows; i += 8) {
        // Загружаем 8 значений double из колоночного массива
        __m512d vals = _mm512_loadu_pd(amounts + i);

        // Читаем 8 битов маски null из bitmap (1 = not null)
        __mmask8 valid = _mm_maskz_set1_epi8(0xFF, 1);  // упрощённо

        // Прибавляем только ненулевые значения (null-safe)
        acc = _mm512_mask_add_pd(acc, valid, acc, vals);
    }

    // Горизонтальная сумма 8 lane-ов в один double
    *result = _mm512_reduce_add_pd(acc);
}

Эквивалентный JVM-код в WholeStageCodegen выглядит как цикл с проверкой if (!isNull[i]) sum += values[i] - компилятор может автовекторизовать это в AVX2, но не гарантированно. В Velox vectorization гарантирована архитектурой.

Velox Memory Management: жизнь без GC

Velox организует всю нативную память через иерархию MemoryPool:

VeloxMemoryManager (глобальный)
├── QueryPool (на запрос)
│   ├── AggregationOperatorPool (для хэш-таблицы)
│   ├── JoinOperatorPool (для probe-side буфера)
│   └── ScanOperatorPool (для Parquet декодирования)
└── SharedPool (для временных буферов shuffle)

Каждый оператор запрашивает и освобождает память через свой pool. Когда суммарное потребление достигает лимита, Velox не падает с OOM - он начинает нативный spill: сериализует промежуточные данные (хэш-таблицу join, буфер сортировки) в бинарный Arrow IPC-формат на локальный диск Executor, освобождает RAM и продолжает работу. После того как нагрузка спадёт - перечитывает с диска.

Это качественно отличается от JVM-spilling: нет Java-сериализации, нет GC-давления во время spill, нет дополнительных STW-пауз в самый неподходящий момент.

Принципиально важно: нативная память Velox не видна JVM. JVM думает, что Executor использует, например, 8 GB heap. На самом деле к этим 8 GB добавляется ещё, скажем, 12 GB нативной Velox-памяти. Контейнерный менеджер (YARN/Kubernetes) видит суммарное потребление - и убьёт контейнер, если оно превысит выделенный лимит. Именно поэтому правильная настройка памяти критична.

Векторы Velox: нативный аналог Arrow

Velox работает с данными в виде RowVector - это нативный C++ аналог Arrow RecordBatch. RowVector содержит:

  • Массив VectorPtr для каждой колонки (FlatVector, DictionaryVector, ArrayVector...)
  • FlatVector<int64_t> - плоский массив int64, FlatVector<double> - плоский double
  • Отдельный null-bitmap (один бит на строку)
  • Размер батча (по умолчанию 10000 строк в Velox, настраивается)

Конвертация между Arrow RecordBatch и Velox RowVector выполняется с минимальным копированием: оба формата используют колоночные массивы в native memory, поэтому Gluten может передавать указатели на буферы без полного копирования данных.


Конфигурация и тюнинг Gluten в Enterprise-среде

Получение артефактов Gluten

В отличие от Comet (который публикуется в Maven Central как готовый JAR), Gluten требует сборки под конкретную архитектуру CPU. Нативные библиотеки Velox компилируются с флагами -mavx2 или -mavx512f - и бинарь, собранный под AVX-512, не запустится на CPU без поддержки AVX-512.

# Вариант 1: скачать готовые бинари из GitHub Releases (для x86-64 + AVX2)
# https://github.com/apache/incubator-gluten/releases
# Доступны JAR под: Ubuntu 22.04, Spark 3.3/3.4/3.5, Scala 2.12/2.13

# Вариант 2: сборка из исходников (рекомендуется для production)
git clone https://github.com/apache/incubator-gluten
cd incubator-gluten

# Установка зависимостей (Ubuntu)
sudo apt-get install -y cmake ninja-build libssl-dev libcurl4-openssl-dev \
    libboost-all-dev libevent-dev libgflags-dev

# Сборка с Velox backend под Spark 3.5
./dev/builddeps-veloxbe.sh  # загружает и компилирует Velox (~30-60 минут)
mvn clean package -Pspark-3.5 -Pbackends-velox -DskipTests

# Результат: gluten-velox-bundle-spark3.5_2.12-1.x.x-SNAPSHOT.jar

Тяжесть сборки - одна из причин, почему Gluten + Velox считается более "хардкорным" решением: он требует DevOps-компетенций и настроенного CI для сборки под каждую версию Spark и каждую модель CPU в кластере.

Подключение плагина и базовые параметры

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("gluten-velox-demo")

    # --- Регистрация Gluten ---
    # Spark Plugin API: инициализирует Velox runtime на каждом Executor
    .config("spark.plugins", "org.apache.gluten.GlutenPlugin")
    # Расширения Catalyst: перехватывает и трансформирует физические планы
    .config("spark.sql.extensions",
            "io.glutenproject.extension.GlutenSessionExtensions")
    # Выбор нативного backend (velox или ch для ClickHouse)
    .config("spark.gluten.sql.columnar.backend.lib", "velox")

    # --- Активация ---
    .config("spark.gluten.enabled", "true")
    # Включить колоночное выполнение для всех поддерживаемых операторов
    .config("spark.gluten.sql.columnar.force", "true")

    # --- Нативный Shuffle ---
    # Заменяет Java-сериализацию на Arrow IPC + C++ сжатие
    .config("spark.shuffle.manager",
            "org.apache.spark.shuffle.sort.ColumnarShuffleManager")

    # --- AQE совместимость ---
    .config("spark.sql.adaptive.enabled", "true")

    .getOrCreate()
)

Математика памяти: критическая зона

Настройка памяти для Gluten + Velox - самая важная и самая часто ошибочная часть конфигурации. Нужно понять, что Velox аллоцирует два типа нативной памяти:

  1. Velox Memory Pool: для операционных данных - хэш-таблицы join, промежуточные буферы агрегации, Arrow RecordBatch с результатами. Настраивается через spark.gluten.memory.offHeap.size.in.bytes.

  2. JVM Off-Heap: для Arrow-буферов на границе JVM ↔ native (передача батчей через JNI), буферов нативного shuffle. Настраивается через spark.memory.offHeap.size.

Правильная схема расчёта для Executor с 24 GB RAM:

# Executor с 24 GB RAM

# JVM heap: минимально необходимый для Spark-коммуникации
# Catalyst, AQE, метаданные задач - Java-код
.config("spark.executor.memory", "4g")

# Off-Heap для JVM Arrow-буферов (граница JVM↔native)
.config("spark.memory.offHeap.enabled", "true")
.config("spark.memory.offHeap.size", "4g")

# Velox Memory Pool: основная операционная память нативного движка
# Это самый "жадный" параметр: join hash tables могут быть огромными
.config("spark.gluten.memory.offHeap.size.in.bytes",
        str(12 * 1024 * 1024 * 1024))  # 12 GB

# memoryOverhead: буфер для JVM накладных расходов + загрузки нативной либы
.config("spark.executor.memoryOverhead", "2g")

# Итого: 4 (JVM heap) + 4 (offHeap) + 12 (Velox Pool) + 2 (overhead) = 22 GB
# Оставляем 2 GB буфер под OS и системные процессы

# Для YARN: явно указываем суммарный контейнер
.config("spark.yarn.executor.memoryOverhead", "2g")

Если Velox Pool закончится во время тяжёлого join - Velox начнёт нативный spill. Это медленнее, чем работа из памяти, но безопаснее OOM. Если же превысить суммарный лимит контейнера - YARN/Kubernetes убьёт процесс жёстко.

# Настройка нативного spill
.config("spark.gluten.sql.columnar.backend.velox.spillEnabled", "true")
.config("spark.gluten.sql.columnar.backend.velox.aggregationSpillEnabled", "true")
.config("spark.gluten.sql.columnar.backend.velox.joinSpillEnabled", "true")
# Директория для spill-файлов на локальном диске Executor
.config("spark.local.dir", "/nvme/spark-spill")

Нативный Shuffle: ColumnarShuffleManager

Стандартный Spark Shuffle работает так: данные сериализуются в Java-формат, записываются на диск, читаются с диска на другом Executor, десериализуются в Java-объекты. При работе с Velox это означает: конвертация Arrow → Java при записи, Java → Arrow при чтении - двойное копирование.

ColumnarShuffleManager заменяет этот процесс:

  • Запись: Arrow ColumnarBatch сжимается нативно (LZ4 или Zstd) и пишется в shuffle-файл как Arrow IPC
  • Чтение: Arrow IPC читается нативно и передаётся прямо в Velox Memory Pool
  • Нулевое копирование: данные никогда не конвертируются в Java-объекты и обратно

Результат: shuffle-интенсивные запросы (SortMergeJoin, repartition, groupBy с перемешиванием) ускоряются дополнительно - помимо ускорения от самих операторов.


Диагностика планов и траблшутинг

Анализ плана выполнения с Gluten

Давайте разберём конкретный пример: запрос с join и агрегацией.

df_orders = spark.read.parquet("/data/orders/")
df_customers = spark.read.parquet("/data/customers/")
df_nations = spark.read.parquet("/data/nations/")

result = (
    df_orders
    .join(df_customers, df_orders.o_custkey == df_customers.c_custkey)
    .join(df_nations, df_customers.c_nationkey == df_nations.n_nationkey)
    .filter(df_orders.o_orderdate >= "2023-01-01")
    .groupBy("n_name")
    .agg(
        F.sum("o_totalprice").alias("total_revenue"),
        F.count("*").alias("order_count"),
        F.avg("o_totalprice").alias("avg_order")
    )
    .orderBy(F.col("total_revenue").desc())
)

result.explain()

Вывод с включённым Gluten:

== Physical Plan ==
VeloxColumnarToRowExec
+- VeloxSortExec[total_revenue DESC]
   +- VeloxHashAggregateExec(keys=[n_name],
   |     funcs=[sum(o_totalprice), count(1), avg(o_totalprice)])
   +- ColumnarExchange hashpartitioning(n_name, 200)
      +- VeloxHashAggregateExec(keys=[n_name],
      |     funcs=[partial_sum(o_totalprice), partial_count(1), partial_avg(...)])
         +- VeloxProjectExec[n_name, o_totalprice]
            +- VeloxBroadcastHashJoinExec[c_nationkey = n_nationkey]
               :- VeloxSortMergeJoinExec[o_custkey = c_custkey]
               :  :- VeloxFilterExec(o_orderdate >= 2023-01-01)
               :  :  +- VeloxScan parquet[o_custkey,o_orderdate,o_totalprice]
               :  +- VeloxScan parquet[c_custkey,c_nationkey]
               +- BroadcastExchange HashedRelationBroadcastMode
                  +- VeloxScan parquet[n_nationkey,n_name]

Разберём структуру:

  • VeloxScan - нативный Parquet-ридер. Читает только нужные колонки (column pruning), применяет фильтры на уровне метаданных Parquet (predicate pushdown).
  • VeloxFilterExec - SIMD-ускоренная фильтрация по дате.
  • VeloxSortMergeJoinExec - нативный Sort Merge Join между orders и customers.
  • VeloxBroadcastHashJoinExec - nations (маленькая таблица) транслируется на все Executor-ы как нативная хэш-таблица.
  • VeloxHashAggregateExec - два прохода агрегации: partial (на Executor) и final (после shuffle).
  • VeloxSortExec - финальная сортировка по revenue.
  • VeloxColumnarToRowExec - конвертация результата обратно в JVM строки.

Весь план выполняется нативно в C++. Единственный переход JVM ↔ native - в самом конце, когда результат небольшого числа строк (ORDER BY + LIMIT) возвращается в JVM.

Идентификация fallback-узлов

Теперь добавим Python UDF и посмотрим, что изменится:

from pyspark.sql.types import StringType

@udf(returnType=StringType())
def categorize_order(price):
    if price > 100000: return "premium"
    elif price > 10000: return "standard"
    else: return "small"

result_with_udf = (
    result
    .withColumn("order_category", categorize_order(F.col("avg_order")))
)

result_with_udf.explain()
== Physical Plan ==
VeloxColumnarToRowExec
+- RowToVeloxColumnarExec                      <-- возврат в native
   +- BatchEvalPython [categorize_order]       <-- Python UDF (JVM-мир)
      +- VeloxColumnarToRowExec                <-- выход из native
         +- VeloxSortExec[...]
            +- VeloxHashAggregateExec(...)
               ...

Паттерн тот же, что и в Comet: два перехода на UDF - выход из native, Python-обработка, возврат в native. Решение идентично: заменить UDF на F.when:

result_fixed = result.withColumn(
    "order_category",
    F.when(F.col("avg_order") > 100000, "premium")
     .when(F.col("avg_order") > 10000, "standard")
     .otherwise("small")
)
# Весь план снова полностью нативный

Настройка логирования fallback-причин

# Включаем детальный лог о том, почему операторы деградируют
spark.conf.set("spark.gluten.sql.columnar.backend.velox.logLevel", "INFO")

# Gluten-специфичное логирование деградации
spark.conf.set("spark.gluten.sql.columnar.fallback.ignoreRowToColumnar", "false")
spark.conf.set("spark.gluten.sql.columnar.wholeStage.numaBinding", "true")

В логах Driver при fallback появятся строки:

[GLUTEN] Operator PythonEvalUDF cannot be offloaded to Velox:
  Reason: Python UDF execution requires JVM Python worker communication
  Impact: 3 downstream operators forced to JVM execution
  Operators affected: [BatchEvalPython, HashAggregateExec, ProjectExec]

Моделирование OOM: симуляция нехватки нативной памяти

Понимание того, как Gluten реагирует на нехватку памяти, важно для production-эксплуатации. Давайте смоделируем ситуацию:

# Намеренно занижаем нативный memory pool (в реальности так не делайте!)
spark_oom = (
    SparkSession.builder
    .appName("gluten-oom-simulation")
    .config("spark.plugins", "org.apache.gluten.GlutenPlugin")
    .config("spark.sql.extensions",
            "io.glutenproject.extension.GlutenSessionExtensions")
    .config("spark.gluten.sql.columnar.backend.lib", "velox")
    .config("spark.executor.memory", "4g")
    .config("spark.memory.offHeap.enabled", "true")
    .config("spark.memory.offHeap.size", "2g")
    # СПЕЦИАЛЬНО занижаем Velox Pool для демонстрации
    .config("spark.gluten.memory.offHeap.size.in.bytes",
            str(256 * 1024 * 1024))  # только 256 MB!
    .config("spark.gluten.sql.columnar.backend.velox.spillEnabled", "false")
    .getOrCreate()
)

# Этот join потребует гораздо больше 256 MB для hash table
# При нехватке и отключённом spill: Velox бросает исключение
try:
    big_join = (
        spark_oom.read.parquet("/data/large_table_a/")
        .join(spark_oom.read.parquet("/data/large_table_b/"), "key")
        .count()
    )
except Exception as e:
    print(f"Ожидаемая ошибка: {e}")
    # [VELOX] Memory pool exhausted: requested 512.0 MB, available 128.0 MB
    # [VELOX] VeloxMemoryCapExceededException
    # После перехода в Spark logs: Executor lost, Stage failed

В логах Executor при подобной ошибке вы увидите:

ERROR VeloxMemoryPool: Memory pool exhausted during HashBuild
  Pool name: HashJoinOperatorPool
  Requested: 524288000 bytes (512 MB)
  Available:  134217728 bytes (128 MB)
  Total pool capacity: 268435456 bytes (256 MB)

FATAL: VeloxMemoryCapExceededException
  native library threw uncaught C++ exception
  Process: Executor-7 terminated with SIGABRT

Исправление: включите spill и увеличьте pool:

# Правильная конфигурация
.config("spark.gluten.memory.offHeap.size.in.bytes",
        str(12 * 1024 * 1024 * 1024))  # 12 GB
.config("spark.gluten.sql.columnar.backend.velox.spillEnabled", "true")
.config("spark.gluten.sql.columnar.backend.velox.aggregationSpillEnabled", "true")
.config("spark.gluten.sql.columnar.backend.velox.joinSpillEnabled", "true")

Полная конфигурация для production

from pyspark.sql import SparkSession

def create_gluten_spark(
    app_name: str,
    executor_cores: int = 8,
    executor_memory_gb: int = 24,
) -> SparkSession:
    """
    Production-ready SparkSession с Gluten + Velox.

    Расчёт памяти:
    - JVM heap:      4 GB  (Catalyst, AQE, метаданные)
    - JVM offHeap:   4 GB  (Arrow-буферы на JNI-границе)
    - Velox Pool:   14 GB  (операционная память C++-операторов)
    - memOverhead:   2 GB  (OS + JVM накладные расходы)
    Итого:          24 GB
    """
    velox_pool_bytes = (executor_memory_gb - 4 - 4 - 2) * 1024 ** 3

    return (
        SparkSession.builder
        .appName(app_name)

        # Gluten plugin
        .config("spark.plugins", "org.apache.gluten.GlutenPlugin")
        .config("spark.sql.extensions",
                "io.glutenproject.extension.GlutenSessionExtensions")
        .config("spark.gluten.sql.columnar.backend.lib", "velox")
        .config("spark.gluten.enabled", "true")

        # Память
        .config("spark.executor.memory", "4g")
        .config("spark.memory.offHeap.enabled", "true")
        .config("spark.memory.offHeap.size", "4g")
        .config("spark.executor.memoryOverhead", "2g")
        .config("spark.gluten.memory.offHeap.size.in.bytes", str(velox_pool_bytes))

        # Nативный spill (безопасность при пиковой нагрузке)
        .config("spark.gluten.sql.columnar.backend.velox.spillEnabled", "true")
        .config("spark.gluten.sql.columnar.backend.velox.aggregationSpillEnabled", "true")
        .config("spark.gluten.sql.columnar.backend.velox.joinSpillEnabled", "true")
        .config("spark.gluten.sql.columnar.backend.velox.orderBySpillEnabled", "true")

        # Нативный Shuffle
        .config("spark.shuffle.manager",
                "org.apache.spark.shuffle.sort.ColumnarShuffleManager")
        .config("spark.gluten.sql.columnar.shuffle.codec", "lz4")

        # AQE
        .config("spark.sql.adaptive.enabled", "true")
        .config("spark.sql.adaptive.coalescePartitions.enabled", "true")

        # Диагностика (отключить в hot path)
        .config("spark.gluten.sql.columnar.backend.velox.logLevel", "WARNING")

        .getOrCreate()
    )

Практика: TPC-H бенчмарк с Gluten + Velox

Запустим те же запросы TPC-H, что использовали для Comet, теперь с Gluten + Velox:

import time
from pyspark.sql import functions as F

def run_tpch_gluten(spark, data_path: str) -> dict:
    """TPC-H запросы оптимизированные под Gluten + Velox (без Python UDF)."""

    for table in ["lineitem", "orders", "customer", "part",
                  "supplier", "partsupp", "nation", "region"]:
        spark.read.parquet(f"{data_path}/{table}").createOrReplaceTempView(table)

    queries = {
        "Q1": """
            SELECT l_returnflag, l_linestatus,
                   SUM(l_quantity) AS sum_qty,
                   SUM(l_extendedprice) AS sum_base_price,
                   SUM(l_extendedprice * (1 - l_discount)) AS sum_disc_price,
                   SUM(l_extendedprice * (1 - l_discount) * (1 + l_tax)) AS sum_charge,
                   AVG(l_quantity) AS avg_qty,
                   AVG(l_extendedprice) AS avg_price,
                   AVG(l_discount) AS avg_disc,
                   COUNT(*) AS count_order
            FROM lineitem
            WHERE l_shipdate <= date_sub(date('1998-12-01'), 90)
            GROUP BY l_returnflag, l_linestatus
            ORDER BY l_returnflag, l_linestatus
        """,
        "Q6": """
            SELECT SUM(l_extendedprice * l_discount) AS revenue
            FROM lineitem
            WHERE l_shipdate >= date('1994-01-01')
              AND l_shipdate < date('1995-01-01')
              AND l_discount BETWEEN 0.05 AND 0.07
              AND l_quantity < 24
        """,
        "Q12": """
            SELECT l_shipmode,
                   SUM(CASE WHEN o_orderpriority IN ('1-URGENT','2-HIGH') THEN 1 ELSE 0 END) AS high_line_count,
                   SUM(CASE WHEN o_orderpriority NOT IN ('1-URGENT','2-HIGH') THEN 1 ELSE 0 END) AS low_line_count
            FROM orders JOIN lineitem
              ON o_orderkey = l_orderkey
            WHERE l_shipmode IN ('MAIL','SHIP')
              AND l_commitdate < l_receiptdate
              AND l_shipdate < l_commitdate
              AND l_receiptdate >= date('1994-01-01')
              AND l_receiptdate < date('1995-01-01')
            GROUP BY l_shipmode
            ORDER BY l_shipmode
        """,
    }

    results = {}
    for qid, sql in queries.items():
        times = []
        for _ in range(3):
            t0 = time.time()
            spark.sql(sql).count()
            times.append(time.time() - t0)
        results[qid] = sorted(times)[1]
        print(f"{qid}: {results[qid]:.2f}s (median of 3)")

    return results

Сравнительная таблица ускорения: Baseline vs Comet vs Gluten+Velox

Типичные результаты на TPC-H SF=100 на кластере из 4 Executor × 32 GB / 16 cores:

Запрос Baseline JVM Apache Comet Gluten + Velox Velox vs Baseline
Q1 (pure agg) 180s 49s (3.7×) 27s 6.7×
Q6 (filter+sum) 72s 20s (3.6×) 12s 6.0×
Q12 (join+case) 240s 86s (2.8×) 48s 5.0×
Q18 (subquery) 420s 168s (2.5×) 105s 4.0×
Среднее - 3.1× 5.4×

Velox устойчиво обгоняет DataFusion/Comet на 30–60% за счёт более агрессивной SIMD-оптимизации (C++ даёт полный контроль над intrinsics) и более зрелого нативного Memory Pool.

Почему Velox опережает DataFusion на join-heavy запросах

Join - это операция, где управление памятью наиболее критично. Velox Hash Join строит хэш-таблицу в специально аллоцированном MemoryPool с заранее известным размером. DataFusion (Rust, jemalloc) также работает с нативной памятью, но Velox дополнительно применяет hardware prefetching (__builtin_prefetch) - загрузку следующих блоков хэш-таблицы в кэш CPU до того, как они понадобятся:

// Velox: SIMD probe в хэш-таблице с prefetching
for (int32_t i = 0; i < batch_size; ++i) {
    // Заблаговременно загружаем в кэш следующий элемент
    __builtin_prefetch(hash_table_ptr + next_hash[i], 0, 1);
    // Проверяем текущий элемент
    bool found = probe_hash_table(hash_table_ptr, keys[i], hash[i]);
    matches[i] = found;
}

На практике prefetching на join-тяжёлых запросах даёт дополнительные 15–25% ускорения.


Best Practices и анти-паттерны

Чек-лист готовности пайплайна к Gluten + Velox

Перед переносом production-пайплайна проверьте:

Данные и форматы:

  • Данные хранятся в Parquet или ORC (не CSV, не JSON row-by-row)
  • Parquet-файлы имеют разумный размер (128–512 MB), не слишком маленькие
  • Схема без экзотических типов: нет UserDefinedType, минимум глубоко вложенных Map<String, Array<Struct<...>>>
  • Колонки с числами имеют явные типы (LongType, DoubleType), не StringType с числами

Код запросов:

  • Нет Python UDF (замените на F.when, F.expr, встроенные функции)
  • Нет Pandas UDF в критическом пути (если нужны - изолируйте в отдельный Stage)
  • Строковые функции используют ISO-форматы дат
  • Нет collect_list с последующей Python-обработкой

Конфигурация:

  • spark.gluten.memory.offHeap.size.in.bytes выставлен явно и корректно
  • spark.gluten.sql.columnar.backend.velox.spillEnabled = true включён как страховка
  • spark.local.dir указывает на NVMe-диск (не сетевую FS) для spill

Верификация:

  • Запущена проверка корректности: результаты Gluten == результаты Baseline (с допуском для float)
  • df.explain() для ключевых запросов проверен: нет неожиданных fallback-операторов
  • Первый запуск в production - на подмножестве данных (canary)

Анти-паттерны

Анти-паттерн 1: Не проверять explain() после включения Gluten

Частая ошибка: включить spark.gluten.enabled=true, запустить задачу, измерить время - "почти не ускорилось". Открыть explain и обнаружить, что 90% плана выполняется в JVM из-за одного MapType в схеме или Pandas UDF на входе.

Анти-паттерн 2: Blindly доверять Velox на float-арифметике

Velox и JVM используют разные реализации floating-point arithmetic. Результаты SUM(double_column) могут отличаться в последних знаках после запятой из-за разного порядка операций (SIMD меняет ассоциативность сложения). Всегда проверяйте корректность с math.isclose(rel_tol=1e-6), а не ==.

Анти-паттерн 3: Не учитывать Velox Pool при расчёте памяти

Самая опасная ошибка - забыть про spark.gluten.memory.offHeap.size.in.bytes. Тогда Velox будет аллоцировать неограниченное количество нативной памяти - пока YARN/Kubernetes не убьёт контейнер жёстким сигналом.

Анти-паттерн 4: Использовать Gluten на IO-bound задачах

Если узкое место - сеть или диск, Velox не поможет: он ускоряет CPU. Задачи типа "читаем CSV, конвертируем в Parquet, пишем" почти не выиграют от Gluten. Задачи типа "читаем Parquet, делаем 5 JOIN, 3 GROUP BY, пишем результат" - выиграют в разы.

Gluten + Velox vs. другие нативные ускорители: итоговое сравнение

Критерий Gluten + Velox Apache Comet Photon Sail
Backend C++ Velox Rust DataFusion C++ (проприетарный) Rust DataFusion
Open Source Да (Apache) Да (Apache) Нет Да
Язык runtime C++ Rust C++ Rust
Ускорение (TPC-H) 4–8× 3–5× 3–10× 3–6×
Сложность деплоя Высокая Низкая Только cloud Средняя
Полнота покрытия Spark API ~85% ~80% ~95% ~75%
RDD API Нет Нет Нет Нет
Управление памятью Velox MemPool + spill Arrow offHeap Проприетарный Arrow offHeap
Native Shuffle Да (ColumnarShuffle) Да (CometShuffleManager) Да Нет
Substrait IR Да Нет (protobuf custom) Нет Нет

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

Условие:

Вам предоставлен Docker-контейнер с Spark 3.5 + Gluten + Velox (готовый образ) и PySpark-скрипт с тяжёлым аналитическим запросом:

# Данный скрипт - стартовая точка для исследования
from pyspark.sql import functions as F
from pyspark.sql.types import StringType, DoubleType
from pyspark.sql import SparkSession

spark = SparkSession.builder \
    .config("spark.plugins", "org.apache.gluten.GlutenPlugin") \
    .config("spark.gluten.enabled", "true") \
    .getOrCreate()

# Псевдо-UDF для классификации заказов
@udf(returnType=StringType())
def order_segment(priority, status):
    if priority in ['1-URGENT', '2-HIGH'] and status == 'O':
        return 'active_priority'
    elif status == 'F':
        return 'fulfilled'
    return 'standard'

# Оконные функции
from pyspark.sql.window import Window
w = Window.partitionBy("c_mktsegment").orderBy("o_totalprice")

result = (
    spark.read.parquet("/data/tpch/sf10/orders")
    .join(spark.read.parquet("/data/tpch/sf10/customer"),
          "o_custkey == c_custkey")
    .join(spark.read.parquet("/data/tpch/sf10/nation"),
          "c_nationkey == n_nationkey")
    # Оконная функция + UDF
    .withColumn("price_rank", F.rank().over(w))
    .withColumn("segment", order_segment(F.col("o_orderpriority"),
                                         F.col("o_orderstatus")))
    # Тяжёлая агрегация
    .groupBy("n_name", "c_mktsegment", "segment")
    .agg(
        F.sum("o_totalprice").alias("total"),
        F.avg("price_rank").alias("avg_rank"),
        F.count("*").alias("cnt")
    )
    .orderBy("total", ascending=False)
)

result.show()

Задача работает 18 минут на Baseline-Spark. После включения spark.gluten.enabled=true - 16 минут (минимальное ускорение).

Что нужно сделать:

  1. Диагностика: запустите скрипт с включённым логированием fallback-причин Gluten. Изучите вывод result.explain() и лог Driver. Определите все узлы, которые падают в JVM, и для каждого укажите причину.

  2. Рефакторинг: замените все источники fallback на нативные аналоги:

  3. Python UDF → F.when / F.expr
  4. Оконную функцию rank() проверьте: поддерживает ли её Velox? Если нет - найдите совместимый аналог или способ убрать её из критического пути.

  5. Верификация корректности: убедитесь, что результаты исправленного скрипта совпадают с результатами оригинала (с допуском rel_tol=1e-5 для float-колонок).

  6. Бенчмарк: запустите три версии на TPC-H SF=10:

  7. Baseline JVM (без Gluten)
  8. Gluten + оригинальный скрипт (с UDF)
  9. Gluten + исправленный скрипт

  10. Отчёт: приложите:

  11. Вывод explain() до и после рефакторинга
  12. Таблицу с временем выполнения трёх версий
  13. Текстовое описание каждого найденного fallback и способа его устранения
  14. Итоговое ускорение исправленной версии относительно Baseline

Ожидаемый результат: ускорение не менее 4× по сравнению с Baseline JVM.