Gluten + Velox: Substrait IR и offloading вычислений в C++ Velox
Gluten + Velox: Substrait IR и offloading вычислений в C++ Velox
Концепция проекта 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 аллоцирует два типа нативной памяти:
-
Velox Memory Pool: для операционных данных - хэш-таблицы join, промежуточные буферы агрегации, Arrow RecordBatch с результатами. Настраивается через
spark.gluten.memory.offHeap.size.in.bytes. -
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 минут (минимальное ускорение).
Что нужно сделать:
-
Диагностика: запустите скрипт с включённым логированием fallback-причин Gluten. Изучите вывод
result.explain()и лог Driver. Определите все узлы, которые падают в JVM, и для каждого укажите причину. -
Рефакторинг: замените все источники fallback на нативные аналоги:
- Python UDF →
F.when/F.expr -
Оконную функцию
rank()проверьте: поддерживает ли её Velox? Если нет - найдите совместимый аналог или способ убрать её из критического пути. -
Верификация корректности: убедитесь, что результаты исправленного скрипта совпадают с результатами оригинала (с допуском
rel_tol=1e-5для float-колонок). -
Бенчмарк: запустите три версии на TPC-H SF=10:
- Baseline JVM (без Gluten)
- Gluten + оригинальный скрипт (с UDF)
-
Gluten + исправленный скрипт
-
Отчёт: приложите:
- Вывод
explain()до и после рефакторинга - Таблицу с временем выполнения трёх версий
- Текстовое описание каждого найденного fallback и способа его устранения
- Итоговое ускорение исправленной версии относительно Baseline
Ожидаемый результат: ускорение не менее 4× по сравнению с Baseline JVM.