Spark Plugin API: SparkPlugin и QueryStage интерфейс для замены executor
Как Spark Plugin API позволяет заменять JVM-выполнение нативным кодом: DriverPlugin, ExecutorPlugin, ColumnarRules, AQE QueryStage, архитектура Gluten и RAPIDS.
Зачем Spark открыл внутренности: исторический контекст¶
До появления Spark Plugin API в версии 3.0 разработчики, которые хотели изменить поведение Spark на уровне выполнения, оказывались перед неприятным выбором: либо форкать Spark и поддерживать собственный патченый рантайм (что делает Facebook, Alibaba и другие гиганты - дорого и сложно), либо мириться с ограничениями JVM-движка.
Доступные точки расширения до 3.0 работали только выше уровня выполнения:
DataSource V1/V2- только для чтения и записи данных, не для изменения логики вычисленийCatalyst rulesчерезSparkSessionExtensions- оптимизация плана на уровне SQL, но не замена физических операторов--jarsи дополнительные классы - добавление кода в classpath, но без доступа к жизненному циклу executorsUDF- бизнес-логика поверх движка, но не замена самого движка
Ни один из этих механизмов не позволял сказать Spark: «когда ты собрался выполнять вот этот оператор HashAggregate на JVM - остановись, отдай мне данные, я выполню это на GPU/через Velox/через DataFusion и верну тебе результат».
Именно это стало возможным с Plugin API в Spark 3.0. Причиной появления стала конкретная проблема: NVIDIA хотела интегрировать GPU-ускорение (проект RAPIDS cuDF) в Spark без форка. Для этого им нужен был легальный способ перехватывать физический план и подменять операторы на GPU-реализации. Spark-коммьюнити признало, что это нужно не только NVIDIA - и создало Plugin API как общий механизм.
Почему UDF и --jars недостаточно¶
Чтобы понять ценность Plugin API, важно осознать фундаментальное ограничение: JVM-движок Spark - это монолит, который всегда выполняет операторы в JVM. UDF добавляет пользовательскую функцию поверх операторов, но сами операторы (HashAggregate, SortMergeJoin, BroadcastHashJoin) всегда выполняются в JVM байткоде через WholeStageCodegen.
Даже если вы хотите заменить один конкретный оператор на GPU-реализацию - нет механизма сказать Catalyst: «вот этот физический узел дерева - выполни его вот так, а не стандартно». До Plugin API такого интерфейса просто не существовало.
Plugin API создаёт этот интерфейс: плагин получает возможность обходить дерево физического плана и заменять узлы на кастомные реализации до того, как выполнение начнётся.
Архитектура Spark Plugin API: двусторонняя система¶
Spark Plugin API состоит из двух независимых интерфейсов, каждый из которых работает на своей стороне распределённой системы. Понимание этого разделения критически важно: Driver и Executor - это отдельные JVM-процессы (часто на разных машинах), и плагин должен инициализировать логику на каждой стороне независимо.
Регистрация плагина¶
Плагин объявляется через конфигурацию Spark. Это единственная точка входа: Spark читает список классов из spark.plugins и инстанцирует их при старте Driver и каждого Executor:
# ──────────────────────────────────────────────────────────
# Регистрация плагина через SparkSession конфигурацию
# ──────────────────────────────────────────────────────────
# spark.plugins - список полных имён Java/Scala-классов через запятую.
# Каждый класс должен реализовывать интерфейс SparkPlugin.
# Spark инстанцирует эти классы дважды:
# 1) В JVM-процессе Driver-а при создании SparkContext
# 2) В JVM-процессе каждого Executor-а при его старте
spark = SparkSession.builder \
.master("spark://host:7077") \
.config("spark.plugins", "com.mycompany.NativeEnginePlugin") \
.config("spark.plugins", "io.glutenproject.GlutenPlugin") \ # Gluten/Velox
# Или несколько плагинов через запятую:
# .config("spark.plugins", "com.a.PluginA,com.b.PluginB") \
.getOrCreate()
# Дополнительные расширения Catalyst (отдельный конфиг):
# SparkSessionExtensions - это другой механизм, работающий
# совместно с Plugin API для внедрения правил оптимизатора
spark = SparkSession.builder \
.config("spark.sql.extensions",
"io.glutenproject.GlutenSparkExtensions") \
.config("spark.plugins", "io.glutenproject.GlutenPlugin") \
.getOrCreate()
Класс, указанный в spark.plugins, должен реализовывать интерфейс SparkPlugin - Scala/Java интерфейс с двумя методами: driverPlugin() и executorPlugin(). Каждый из них возвращает объект соответствующего типа (или null, если плагин не нужен на этой стороне).
Интерфейс SparkPlugin: точка входа¶
В Scala/Java определение выглядит так (упрощённо):
// Интерфейс, который должен реализовывать ваш плагин-класс
trait SparkPlugin {
// Возвращает объект для Driver-стороны.
// Вызывается один раз при создании SparkContext на Driver.
// Может вернуть null если Driver-логика не нужна.
def driverPlugin(): DriverPlugin
// Возвращает объект для Executor-стороны.
// Вызывается при старте каждого Executor.
// Может вернуть null если Executor-логика не нужна.
def executorPlugin(): ExecutorPlugin
}
Важная деталь: SparkPlugin - это просто фабрика. Реальная логика живёт в DriverPlugin и ExecutorPlugin.
DriverPlugin: мозг плагина на стороне Driver¶
DriverPlugin отвечает за инициализацию плагина на стороне Driver. Driver - это главный JVM-процесс, который управляет всем Spark-приложением: строит планы, управляет задачами, координирует shuffle. Плагин, работающий на Driver, может:
- Регистрировать кастомные правила оптимизатора в Catalyst
- Получать метрики от Executor-ов
- Отправлять управляющие сообщения на Executor-ы
- Инициализировать Driver-сторону нативных библиотек
Интерфейс DriverPlugin¶
// Интерфейс DriverPlugin (Spark 3.0+)
trait DriverPlugin {
// Вызывается при создании SparkContext на Driver.
// sc - SparkContext (доступ ко всему Spark)
// pluginContext - вспомогательный контекст плагина
// Возвращает Map[String, String] - конфиг, который будет
// передан на Executor-ы через RPC при их старте.
def init(
sc: SparkContext,
pluginContext: PluginContext
): java.util.Map[String, String]
// Вызывается когда Driver получает RPC-сообщение от
// ExecutorPlugin через pluginContext.ask() или send().
// Позволяет организовать двунаправленный обмен данными.
def receive(message: AnyRef): AnyRef
// Вызывается при завершении приложения.
// Используйте для освобождения Driver-сторонних ресурсов.
def shutdown(): Unit
}
Пример: DriverPlugin для логирования планов¶
Рассмотрим учебный пример DriverPlugin, который логирует каждый сформированный физический план. Это полезно для отладки - вы видите все планы до и после оптимизации:
// LoggingDriverPlugin.java
import org.apache.spark.SparkContext;
import org.apache.spark.api.plugin.DriverPlugin;
import org.apache.spark.api.plugin.PluginContext;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.execution.QueryExecution;
import org.apache.spark.sql.util.QueryExecutionListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.HashMap;
import java.util.Map;
public class LoggingDriverPlugin implements DriverPlugin {
private static final Logger LOG =
LoggerFactory.getLogger(LoggingDriverPlugin.class);
@Override
public Map<String, String> init(SparkContext sc, PluginContext pluginContext) {
LOG.info("[Plugin] DriverPlugin.init() called on Driver");
// Получаем SparkSession из SparkContext
SparkSession spark = SparkSession.active();
// Регистрируем QueryExecutionListener - он вызывается
// ПОСЛЕ каждого SQL-запроса (для мониторинга и анализа).
// Это не то же самое что Catalyst rules - это post-execution hook.
spark.listenerManager().register(new QueryExecutionListener() {
@Override
public void onSuccess(
String funcName,
QueryExecution qe,
long durationNs) {
// qe.sparkPlan() - физический план ДО оптимизации AQE
// qe.executedPlan() - финальный план С оптимизацией AQE
LOG.info("[Plugin] Query '{}' succeeded in {}ms",
funcName, durationNs / 1_000_000);
LOG.debug("[Plugin] Physical plan:\n{}",
qe.executedPlan().toString());
}
@Override
public void onFailure(
String funcName,
QueryExecution qe,
Exception exception) {
LOG.error("[Plugin] Query '{}' FAILED: {}",
funcName, exception.getMessage());
}
});
// Возвращаем конфиг, который будет передан на все Executor-ы.
// Это Map<String, String> с произвольными параметрами.
// Executor-ы получат его в ExecutorPlugin.init().
Map<String, String> config = new HashMap<>();
config.put("plugin.driver.host", sc.master());
config.put("plugin.version", "1.0.0");
config.put("plugin.log.level", "DEBUG");
return config;
}
@Override
public Object receive(Object message) {
// Здесь обрабатываем сообщения от Executor-ов.
// Например, ExecutorPlugin может отправить статистику
// использования нативной памяти на Driver.
LOG.info("[Plugin] Received message from Executor: {}", message);
return null; // ответ Executor-у (опционально)
}
@Override
public void shutdown() {
LOG.info("[Plugin] DriverPlugin.shutdown() - cleaning up");
}
}
В методе init() обратите внимание на возвращаемый Map<String, String>. Это механизм передачи конфигурации от Driver к Executor-ам: когда Executor стартует и инициализирует свой ExecutorPlugin, он получает эту Map как параметр. Это единственный безопасный способ передать Driver-сторонние параметры (например, адрес нативного memory pool) на все воркеры.
ExecutorPlugin: мускулы плагина на воркерах¶
ExecutorPlugin - более мощная и одновременно более опасная часть Plugin API. Он работает внутри каждого JVM-процесса Executor, имеет доступ к выполнению Task-ов и может инициализировать нативные библиотеки.
Ключевой факт: ExecutorPlugin.init() вызывается до запуска первой задачи на Executor. Это означает, что плагин может инициализировать нативные ресурсы (аллоцировать Off-Heap память, загрузить C++/Rust библиотеку через JNI) до того, как Spark начнёт отправлять задачи на этот Executor.
Интерфейс ExecutorPlugin¶
// Интерфейс ExecutorPlugin (Spark 3.0+)
trait ExecutorPlugin {
// Вызывается при старте Executor до первой задачи.
// pluginContext - контекст для отправки сообщений на Driver
// extraConf - Map из DriverPlugin.init() (Driver → Executor transfer)
def init(
pluginContext: PluginContext,
extraConf: java.util.Map[String, String]
): Unit
// Вызывается ПЕРЕД запуском каждой Task на этом Executor.
// taskContext - информация о текущей задаче
def onTaskStart(): Unit
// Вызывается ПОСЛЕ завершения каждой Task (успешно или с ошибкой).
// succeeded - true если задача завершилась без исключений
def onTaskSucceeded(): Unit
def onTaskFailed(failureReason: TaskFailedReason): Unit
// Вызывается при завершении Executor (shutdown).
def shutdown(): Unit
}
Пример: ExecutorPlugin с инициализацией нативной библиотеки¶
Рассмотрим пример, который имитирует инициализацию нативного Arrow memory allocator - похожий код есть в реальных плагинах Gluten и RAPIDS:
// NativeExecutorPlugin.java
import org.apache.spark.api.plugin.ExecutorPlugin;
import org.apache.spark.api.plugin.PluginContext;
import org.apache.spark.TaskContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Map;
import java.util.concurrent.atomic.AtomicLong;
public class NativeExecutorPlugin implements ExecutorPlugin {
private static final Logger LOG =
LoggerFactory.getLogger(NativeExecutorPlugin.class);
// Счётчик аллоцированной нативной памяти (Off-Heap).
// AtomicLong потому что несколько задач могут работать параллельно.
private final AtomicLong nativeMemoryAllocated = new AtomicLong(0);
// Ссылка на контекст плагина для отправки сообщений на Driver
private PluginContext pluginContext;
// Указатель на нативный контекст (в реальности - long-адрес в памяти)
// В JNI это типичная практика: хранить void* как long в Java
private long nativeContextPtr = 0L;
@Override
public void init(
PluginContext pluginContext,
Map<String, String> extraConf) {
this.pluginContext = pluginContext;
LOG.info("[Plugin] ExecutorPlugin.init() - initializing native runtime");
// Читаем конфигурацию, переданную от DriverPlugin.init()
String logLevel = extraConf.getOrDefault("plugin.log.level", "INFO");
String driverHost = extraConf.getOrDefault("plugin.driver.host", "unknown");
LOG.info("[Plugin] Driver host: {}, log level: {}", driverHost, logLevel);
// В реальном плагине (Gluten, RAPIDS) здесь происходит:
// 1. Загрузка нативной библиотеки через System.loadLibrary()
// 2. Инициализация нативного memory allocator
// 3. Регистрация JNI-коллбэков
// Пример загрузки нативной библиотеки:
// System.loadLibrary("native_spark_plugin"); // libNative_spark_plugin.so
// this.nativeContextPtr = initNativeContext(offHeapSizeBytes);
// Мы имитируем инициализацию:
long offHeapSizeMb = Long.parseLong(
extraConf.getOrDefault("plugin.offheap.mb", "1024")
);
LOG.info("[Plugin] Reserving {}MB Off-Heap for native allocator", offHeapSizeMb);
this.nativeContextPtr = 1L; // имитация указателя
// Отправляем сообщение на Driver о готовности Executor-а
// pluginContext.send() - асинхронная отправка
pluginContext.send("EXECUTOR_READY:" + offHeapSizeMb + "MB");
}
@Override
public void onTaskStart() {
// Вызывается перед каждой Task.
// Полезно для сброса per-task счётчиков или
// инициализации task-level нативных ресурсов.
LOG.debug("[Plugin] Task starting on executor");
}
@Override
public void onTaskSucceeded() {
// Вызывается после успешной Task.
// Можно логировать статистику использования нативной памяти.
long allocatedMb = nativeMemoryAllocated.get() / (1024 * 1024);
if (allocatedMb > 0) {
LOG.info("[Plugin] Task completed, native memory peak: {}MB", allocatedMb);
}
// В реальном плагине: освобождение task-level нативных буферов
}
@Override
public void onTaskFailed(TaskFailedReason failureReason) {
// Критически важно: при падении Task нативные ресурсы
// могут не освободиться автоматически (нет GC для нативной памяти).
// Здесь ОБЯЗАТЕЛЬНО нужно освобождать нативные буферы.
LOG.error("[Plugin] Task FAILED: {} - releasing native resources",
failureReason);
// В реальном плагине: освободить нативную память принудительно
// freeNativeTaskResources(nativeContextPtr);
nativeMemoryAllocated.set(0);
}
@Override
public void shutdown() {
LOG.info("[Plugin] ExecutorPlugin.shutdown() - freeing native context");
// В реальном плагине: освободить нативный контекст полностью
// destroyNativeContext(nativeContextPtr);
this.nativeContextPtr = 0L;
}
}
Обратите особое внимание на onTaskFailed: в нативных плагинах это критически важный метод. GC Java не знает о памяти, аллоцированной нативным кодом через JNI. Если задача упала - нативные буферы должны быть освобождены вручную. Без этого Executor будет медленно утекать по Off-Heap памяти, пока не упадёт по OOM.
Перехват планов: SparkSessionExtensions и ColumnarRules¶
DriverPlugin и ExecutorPlugin - это инфраструктурный уровень плагина. Но чтобы заменить физические операторы Spark на нативные, нужен другой механизм: расширения Catalyst.
SparkSessionExtensions: внедрение правил в pipeline оптимизатора¶
Catalyst optimizer в Spark - это pipeline правил (rules), которые последовательно трансформируют дерево плана. Plugin API позволяет внедрять кастомные правила в этот pipeline через SparkSessionExtensions.
Существуют несколько точек внедрения - каждая работает на разной стадии:
// Регистрация расширений Catalyst в SparkSessionExtensions
// (обычно делается внутри DriverPlugin.init() или отдельным классом)
spark.withExtensions { extensions =>
// injectOptimizerRule - добавление правила в Catalyst LOGICAL optimizer.
// Работает с Logical Plan (до выбора физических операторов).
// Используется для: predicate pushdown, constant folding,
// join reordering и других логических оптимизаций.
extensions.injectOptimizerRule { session =>
new MyCustomLogicalOptimizationRule(session)
}
// injectPlannerStrategy - добавление стратегии Physical Planner.
// Определяет, как конкретный Logical оператор превращается
// в Physical оператор. Используется для замены стандартных
// физических операторов на нативные.
extensions.injectPlannerStrategy { session =>
new NativePlannerStrategy(session)
}
// injectQueryStagePrepRule - правила, применяемые к КАЖДОМУ
// QueryStage ПЕРЕД его материализацией (в режиме AQE).
// Это самая поздняя точка, где можно изменить план до выполнения.
extensions.injectQueryStagePrepRule { session =>
new PreStageNativeOptimizationRule(session)
}
// injectColumnar - специальные правила для замены ROW-based
// операторов на COLUMNAR операторы (работающие с Arrow ColumnarBatch).
// Это главный механизм, через который Gluten/RAPIDS заменяют операторы.
extensions.injectColumnar { session =>
NativeColumnarOverrideRules(session)
}
// injectPostHocResolutionRule - правила после анализа,
// но до оптимизации. Используется редко.
extensions.injectPostHocResolutionRule { session =>
new MyPostHocRule(session)
}
}
Ключевой механизм: injectColumnar и ColumnarRule¶
Наиболее важный для нативных движков механизм - injectColumnar. Он позволяет внедрить ColumnarRule, который получает весь физический план и может заменять операторы.
ColumnarRule - это два метода:
preColumnarTransitions(plan)- вызывается ДО применения стандартных ColumnarToRow/RowToColumnar переходовpostColumnarTransitions(plan)- вызывается ПОСЛЕ этих переходов
В Gluten и RAPIDS именно preColumnarTransitions используется для замены JVM-операторов на нативные:
// Упрощённый пример ColumnarRule для замены операторов
// (реальная реализация в Gluten значительно сложнее)
class NativeColumnarOverrideRules(session: SparkSession)
extends ColumnarRule {
// Вызывается с физическим планом.
// Возвращает трансформированный план.
override def preColumnarTransitions(plan: SparkPlan): SparkPlan = {
// transformDown обходит дерево план сверху вниз
// и применяет partial function к каждому узлу
plan.transformDown {
// Если встречаем стандартный HashAggregateExec (JVM) —
// заменяем на нативный NativeHashAggregateExec
case agg: HashAggregateExec
if canNativelyExecute(agg) =>
NativeHashAggregateExec(
agg.requiredChildDistributionExpressions,
agg.groupingExpressions,
agg.aggregateExpressions,
agg.aggregateAttributes,
agg.initialInputBufferOffset,
agg.resultExpressions,
agg.child
)
// SortMergeJoinExec → NativeSortMergeJoinExec
case join: SortMergeJoinExec
if canNativelyExecute(join) =>
NativeSortMergeJoinExec(join.leftKeys, join.rightKeys,
join.joinType, join.condition, join.left, join.right)
// FilterExec → NativeFilterExec
case filter: FilterExec
if canNativelyExecute(filter) =>
NativeFilterExec(filter.condition, filter.child)
// Все остальные операторы - оставляем JVM-реализацию.
// Это принципиальная особенность plug-in подхода:
// незамена → автоматический fallback на JVM, не ошибка.
}
}
// Проверяет, может ли нативный движок выполнить оператор.
// Если нет - оставляем JVM-версию (graceful degradation).
private def canNativelyExecute(plan: SparkPlan): Boolean = {
// Проверяем типы данных всех колонок
val unsupportedTypes = plan.output.exists { attr =>
attr.dataType match {
case _: MapType => true // MapType не поддерживается нативно
case _: UserDefinedType[_] => true // кастомные типы - нет
case _ => false
}
}
!unsupportedTypes
}
}
Обратите внимание на canNativelyExecute: это ключевая особенность правильно написанного плагина. Нативные движки не поддерживают все операторы Spark - у каждого есть пробелы. Graceful degradation (возврат к JVM для неподдерживаемых операторов) обеспечивает 100% корректность запросов даже при неполной поддержке.
Adaptive Query Execution и QueryStage¶
Чтобы полностью понять, где именно плагин перехватывает выполнение, нужно понять AQE (Adaptive Query Execution) и концепцию QueryStage.
Что такое AQE и почему он изменил всё¶
До AQE (Spark 3.0) весь физический план строился один раз до начала выполнения. Это создавало проблему: количество shuffle partitions, стратегия JOIN (broadcast vs sort-merge) и другие решения принимались на основе статистики, собранной до выполнения. Статистика могла быть неточной, и план часто был субоптимальным.
AQE решает это: он разбивает физический план на QueryStage - изолированные сегменты, разделённые по границам shuffle. Каждый QueryStage выполняется независимо, после его завершения Spark смотрит на реальную статистику результата (размер shuffle output) и принимает решения для следующих стадий. Это называется adaptive (адаптивная) оптимизация.
QueryStage как единица выполнения¶
QueryStage - это фрагмент физического плана между двумя Exchange-узлами (shuffle-операциями). Все операторы внутри одного QueryStage выполняются как одна единая задача (набор Task-ов) без промежуточных shuffle.
Пример физического плана с AQE:
SortMergeJoinExec
├── Exchange (ShuffleExchangeExec) ← граница QueryStage 2
│ └── FilterExec + ParquetScan ← QueryStage 1 (orders)
└── Exchange (BroadcastExchangeExec) ← граница QueryStage 3
└── FilterExec + ParquetScan ← QueryStage 0 (customers)
AQE выполняет:
1. QueryStage 0 и 1 параллельно
2. После завершения: видит что customers < 10MB → меняет SortMergeJoin на BroadcastHashJoin
3. Выполняет финальный join с broadcast
QueryStagePrepRule: последний шанс изменить план¶
injectQueryStagePrepRule - это правило, применяемое к каждому QueryStage непосредственно перед его выполнением. Это самая поздняя точка в lifecycle плана, где плагин может вмешаться.
Именно здесь нативные плагины делают финальную замену операторов - потому что к этому моменту AQE уже применил свои динамические оптимизации (например, сменил JOIN стратегию), и план максимально близок к финальному виду:
// QueryStagePrepRule - применяется к каждому QueryStage перед выполнением
class NativeQueryStagePrepRule(session: SparkSession)
extends Rule[SparkPlan] {
override def apply(plan: SparkPlan): SparkPlan = {
// plan - физический план ОДНОГО QueryStage
// (уже с AQE-оптимизациями)
// Логируем план для диагностики
LOG.info("QueryStage plan before native transformation:\n{}",
plan.treeString)
// Применяем нативные трансформации к этому stage
val nativePlan = applyNativeTransformations(plan)
LOG.info("QueryStage plan after native transformation:\n{}",
nativePlan.treeString)
nativePlan
}
private def applyNativeTransformations(plan: SparkPlan): SparkPlan = {
plan.transformDown {
// Замена операторов - аналогично ColumnarRule
case agg: HashAggregateExec if isNativeSupported(agg) =>
NativeHashAggregateExec.from(agg)
// ... другие замены
}
}
}
Физический уровень: JNI и передача Arrow буферов¶
Когда Spark Executor получает Task для выполнения - он вызывает метод compute() физического оператора. Для стандартных JVM-операторов compute() возвращает Iterator[InternalRow] (строковый итератор) или Iterator[ColumnarBatch] (батчевый итератор для columnar операторов).
Для нативных операторов (например, NativeHashAggregateExec из плагина) compute() реализует переход через JNI (Java Native Interface) к нативному коду:
На схеме показан полный цикл выполнения нативного оператора. Ключевые моменты:
Передача данных через указатель (zero-copy). Входные данные уже находятся в Arrow ColumnarBatch - это структура, которая хранит данные в Off-Heap памяти в Arrow формате. Нативному коду не нужно копировать данные - достаточно получить указатель (long в Java = void* в C++) на Arrow буфер.
JNI вызов. Java Native Interface - это стандартный механизм вызова C/C++/Rust кода из JVM. Оверхед одного JNI-вызова - несколько микросекунд. Для батча из 8192 строк это пренебрежимо мало по сравнению с выигрышем от SIMD-выполнения.
Возврат через указатель. Нативный код аллоцирует Arrow буфер для результата в Off-Heap памяти и возвращает указатель. JVM оборачивает указатель в ColumnarBatch без копирования - данные по-прежнему в нативной памяти.
Пример нативного оператора (упрощённый псевдокод)¶
// NativeHashAggregateExec - упрощённый учебный пример
// Реальный код в Gluten/RAPIDS значительно сложнее
public class NativeHashAggregateExec extends SparkPlan
implements ColumnarSupport {
private final HashAggregateExec originalPlan;
// Флаг: используем ли columnar (Arrow) или row-based путь
@Override
public boolean supportsColumnar() {
return true; // говорим Spark: мы работаем с ColumnarBatch
}
@Override
public RDD<ColumnarBatch> executeColumnar() {
// Получаем входные данные в columnar формате от child оператора
RDD<ColumnarBatch> childRDD = child().executeColumnar();
return childRDD.mapPartitions(iter -> {
// Собираем все батчи партиции в нативной памяти
List<Long> arrowPtrs = new ArrayList<>();
while (iter.hasNext()) {
ColumnarBatch batch = iter.next();
// extractArrowBufPtr() возвращает long-указатель
// на Arrow буфер в Off-Heap памяти
arrowPtrs.add(ArrowUtils.extractArrowBufPtr(batch));
}
// Один JNI-вызов для всей партиции
// nativeAggregate() - метод, объявленный как native
long resultPtr = nativeAggregate(
arrowPtrs.stream().mapToLong(Long::longValue).toArray(),
buildAggregateConfig()
);
// Оборачиваем результат в ColumnarBatch без копирования
ColumnarBatch resultBatch = ArrowUtils.wrapArrowPtr(resultPtr);
return Collections.singletonList(resultBatch).iterator();
});
}
// JNI-объявление: этот метод реализован в libNativePlugin.so
// При вызове JVM передаёт управление в нативный код
private native long nativeAggregate(
long[] arrowInputPtrs,
byte[] aggregateConfig // protobuf-сериализованный конфиг
);
}
Реальные плагины: Gluten, RAPIDS, Comet¶
Рассмотрим, как Plugin API используется тремя ключевыми production-плагинами. Каждый реализует схожую архитектуру (DriverPlugin + ExecutorPlugin + ColumnarRules + JNI), но с разными нативными backends и разными компромиссами.
Apache Gluten + Velox: C++ движок от Meta¶
Gluten (Intel, 2022) - это слой адаптации, который позволяет Spark использовать Velox как backend выполнения. Velox - это vectorized execution engine на C++, созданный Meta для Presto/Trino и обрабатывающий экзабайты данных в Meta Production.
Архитектура Gluten:
PySpark / SparkSQL
↓
Spark Catalyst (Logical Plan → Physical Plan)
↓
GlutenPlugin (ColumnarRule):
SortMergeJoinExec → VeloxSortMergeJoinExec (C++ JNI)
HashAggregateExec → VeloxHashAggregateExec (C++ JNI)
BroadcastHashJoinExec → VeloxBroadcastHashJoin (C++ JNI)
ParquetScan → VeloxParquetScan (C++ JNI)
[UDFs, MLlib, ...] → JVM fallback (без изменений)
↓
Velox Execution Engine (C++)
- Velox expression evaluator (vectorized)
- Velox hash table (cache-friendly)
- Velox Parquet reader (columnar, no JVM deserialization)
- AVX-512/AVX2 SIMD через Velox operators
Конфигурация Gluten:
spark = SparkSession.builder \
.config("spark.plugins", "io.glutenproject.GlutenPlugin") \
.config("spark.sql.extensions",
"io.glutenproject.GlutenSparkExtensions") \
# Выбор нативного backend: velox (CPU) или другие
.config("spark.gluten.sql.columnar.backend.lib", "velox") \
# Off-Heap размер для Velox - Velox не использует JVM heap
.config("spark.memory.offHeap.enabled", "true") \
.config("spark.memory.offHeap.size", "8g") \
# Включаем конкретные операторы (можно контролировать фолбек)
.config("spark.gluten.sql.columnar.hashagg", "true") \
.config("spark.gluten.sql.columnar.sort", "true") \
.config("spark.gluten.sql.columnar.shuffleManager",
"org.apache.spark.shuffle.sort.ColumnarShuffleManager") \
.getOrCreate()
Ключевая особенность Gluten: ColumnarShuffleManager. Стандартный Spark Shuffle преобразует Arrow ColumnarBatch в строки (ColumnarToRow), сериализует строки, передаёт по сети, десериализует, конвертирует обратно в строки (RowToColumnar). Gluten заменяет это нативным shuffle, который передаёт Arrow IPC напрямую - без двойной конвертации.
NVIDIA RAPIDS cuDF: GPU-ускорение¶
RAPIDS cuDF (NVIDIA) - самый известный пример Plugin API. Он заменяет CPU-выполнение Spark на GPU через CUDA:
spark = SparkSession.builder \
# RAPIDS Plugin для Spark
.config("spark.plugins", "com.nvidia.spark.SQLPlugin") \
# GPU resource discovery
.config("spark.worker.resource.gpu.discoveryScript",
"/opt/sparkRapidsPlugin/getGpusResources.sh") \
.config("spark.executor.resource.gpu.amount", "1") \
.config("spark.task.resource.gpu.amount", "0.5") \
# Размер GPU-батча
.config("spark.rapids.sql.batchSizeBytes", "512m") \
# Fallback для неподдерживаемых операций
.config("spark.rapids.sql.incompatibleOps.enabled", "false") \
.getOrCreate()
RAPIDS работает по тому же принципу: ColumnarRule заменяет CPU-операторы GPU-операторами. Вместо JNI → C++ (как в Gluten) используется JNI → CUDA. Данные перемещаются CPU RAM → GPU VRAM (через PCIe или NVLink), операции выполняются на GPU tensor cores, результат возвращается как Arrow ColumnarBatch.
Ограничение RAPIDS: GPU - это дорогой ресурс с жёсткими требованиями. Минимум 1 GPU на Executor, данные должны помещаться в GPU VRAM (обычно 16–80 GB). RAPIDS выгоден на очень широких таблицах (тысячи колонок) и ML-операциях (XGBoost, RAPIDS ML).
Apache Comet: Rust/DataFusion плагин¶
Comet (часть Apache Arrow проекта) - это плагин для Spark, который заменяет физические операторы на реализации из Apache DataFusion (Rust):
spark = SparkSession.builder \
.config("spark.plugins", "org.apache.spark.CometPlugin") \
.config("spark.sql.extensions",
"org.apache.comet.CometSparkSessionExtensions") \
# Включаем нативное выполнение
.config("spark.comet.enabled", "true") \
# Нативный shuffle (Arrow IPC между воркерами)
.config("spark.comet.exec.shuffle.enabled", "true") \
# Размер батча для нативных операторов
.config("spark.comet.batchSize", "8192") \
.getOrCreate()
Comet использует DataFusion как execution engine - тот же движок, что используется в Sail (который мы изучали в предыдущем уроке). Разница между Comet и Sail: Comet - это плагин к Spark (JVM-драйвер остаётся), Sail - полная замена (JVM отсутствует).
ColumnarBatch: интерфейс между JVM и нативным кодом¶
ColumnarBatch - это центральная абстракция для взаимодействия JVM Spark и нативных плагинов. Понимание этого класса критически важно для понимания всей архитектуры.
ColumnarBatch представляет батч данных в columnar формате. Каждая колонка - это ColumnVector (интерфейс Spark). Реализации ColumnVector могут хранить данные в JVM heap (OnHeapColumnVector), Off-Heap памяти (OffHeapColumnVector) или нативной памяти через Arrow (ArrowColumnVector).
# Python-сторона: как ColumnarBatch виден в PySpark
# (обычно прозрачно, но можно увидеть в explain())
# Запрос, который вызывает переходы Row ↔ Columnar:
df.filter("amount > 1000") \
.groupBy("region") \
.agg({"amount": "sum"}) \
.explain(extended=True)
# В плане при включённом нативном плагине (например Gluten):
# == Physical Plan ==
# VeloxColumnarToRowExec ← конвертация из Columnar в Row
# +- VeloxHashAggregateExec ← нативная агрегация (Velox, columnar)
# +- VeloxColumnarShuffleExchange ← нативный shuffle (Arrow IPC)
# +- VeloxHashAggregateExec ← partial aggregate (нативная)
# +- VeloxFilterExec ← нативный filter (SIMD)
# +- VeloxParquetScan ← нативное чтение Parquet
#
# vs без плагина:
# == Physical Plan ==
# *(2) HashAggregate(...)
# +- Exchange hashpartitioning(region, 200)
# +- *(1) HashAggregate(...)
# +- *(1) Filter (amount > 1000)
# +- *(1) ColumnarToRow ← конвертация Arrow → Row
# +- FileScan parquet
Переходы Row ↔ Columnar: неизбежные издержки¶
Одна из главных проблем plug-in архитектуры - границы между JVM-операторами и нативными операторами. Когда данные переходят с JVM-стороны на нативную (или обратно), их нужно конвертировать.
ColumnarToRow - конвертация Arrow ColumnarBatch в JVM InternalRow. Это дорогая операция: нужно распаковать каждый Arrow буфер в Java-объекты, проверить null-биты, построить InternalRow для каждой строки.
RowToColumnar - обратная конвертация. Строки из JVM нужно упаковать в Arrow columnar формат.
Хорошо написанный плагин (Gluten, RAPIDS) минимизирует эти переходы, объединяя максимум операторов в нативном пути. Но некоторые операторы (UDF, MLlib) неизбежно остаются на JVM, создавая переходы в середине плана.
Практика: запуск с плагином и анализ планов¶
Конфигурирование тестовой среды¶
# ──────────────────────────────────────────────────────────
# Запуск с Comet (наиболее простой для локального тестирования)
# ──────────────────────────────────────────────────────────
# Установка: pip install pyspark apache-comet
# (или скачать JAR: https://github.com/apache/datafusion-comet/releases)
import os
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
# Путь к JAR-файлу Comet
comet_jar = "/path/to/comet-spark-spark3.5_2.12-0.4.0.jar"
spark = SparkSession.builder \
.master("local[4]") \
.config("spark.driver.extraClassPath", comet_jar) \
.config("spark.executor.extraClassPath", comet_jar) \
.config("spark.plugins", "org.apache.spark.CometPlugin") \
.config("spark.sql.extensions",
"org.apache.comet.CometSparkSessionExtensions") \
.config("spark.comet.enabled", "true") \
.config("spark.comet.exec.shuffle.enabled", "true") \
# Off-Heap для нативных буферов
.config("spark.memory.offHeap.enabled", "true") \
.config("spark.memory.offHeap.size", "2g") \
# Уровень логирования плагина
.config("spark.comet.logLevel", "INFO") \
.getOrCreate()
spark.sparkContext.setLogLevel("WARN")
print(f"Spark version: {spark.version}")
print(f"Comet enabled: {spark.conf.get('spark.comet.enabled')}")
Анализ планов: с плагином и без¶
# ──────────────────────────────────────────────────────────
# Тестовый запрос: тяжёлая аналитика
# ──────────────────────────────────────────────────────────
# Создаём DataFrame для тестирования
import random
from pyspark.sql.types import *
# Генерируем данные: транзакции с нескольких регионов
n_rows = 1_000_000
data = [(
i,
random.choice(["North", "South", "East", "West", "Central"]),
random.choice(["Electronics", "Clothing", "Food", "Sports"]),
random.uniform(10.0, 10000.0),
random.randint(1, 100)
) for i in range(n_rows)]
schema = StructType([
StructField("id", LongType(), False),
StructField("region", StringType(), True),
StructField("category", StringType(), True),
StructField("amount", DoubleType(), True),
StructField("quantity", IntegerType(), True),
])
df = spark.createDataFrame(data, schema).cache()
df.count() # форсируем кэширование
# ──────────────────────────────────────────────────────────
# Запрос с JOIN + GROUP BY + оконными функциями
# ──────────────────────────────────────────────────────────
from pyspark.sql.window import Window
# Таблица регионов (для JOIN)
regions = spark.createDataFrame([
("North", "Moscow"),
("South", "Krasnodar"),
("East", "Vladivostok"),
("West", "Saint Petersburg"),
("Central", "Voronezh")
], ["region", "city"])
# Сложный аналитический запрос
query = df \
.join(regions, on="region", how="left") \
.filter(F.col("amount") > 100) \
.groupBy("region", "city", "category") \
.agg(
F.sum("amount").alias("total_revenue"),
F.count("*").alias("order_count"),
F.avg("amount").alias("avg_order"),
F.sum(F.col("amount") * F.col("quantity")).alias("gross_revenue")
) \
.withColumn(
"revenue_rank",
F.rank().over(
Window.partitionBy("region").orderBy(F.desc("total_revenue"))
)
)
# ──────────────────────────────────────────────────────────
# Анализ плана: ищем нативные операторы
# ──────────────────────────────────────────────────────────
print("=" * 60)
print("PHYSICAL PLAN WITH COMET PLUGIN:")
print("=" * 60)
query.explain(extended=True)
# Что искать в плане с Comet:
# - CometScanExec или CometParquetScan - нативное чтение
# - CometFilterExec - нативная фильтрация (SIMD)
# - CometHashAggregateExec - нативная агрегация (DataFusion)
# - CometSortMergeJoinExec или CometBroadcastHashJoinExec
# - CometColumnarToRowExec - переход из нативного в JVM row-based
# (присутствует только на финальном выходе или у неподдерживаемых операторов)
# Что НЕ должно быть много в плане с плагином:
# - ColumnarToRow / RowToColumnar в СЕРЕДИНЕ плана
# (это означает неподдерживаемый оператор создал разрыв)
Сравнение производительности: с плагином и без¶
# ──────────────────────────────────────────────────────────
# Бенчмарк: Spark vs Spark+Comet
# ──────────────────────────────────────────────────────────
import time
def run_and_time(label, query_fn, n_runs=3):
"""Запускает запрос n_runs раз, возвращает медиану."""
times = []
for i in range(n_runs):
start = time.time()
result = query_fn()
result.count() # форсируем полное выполнение
elapsed = time.time() - start
times.append(elapsed)
print(f" {label} run {i+1}: {elapsed:.2f}s")
return sorted(times)[n_runs // 2]
# Запускаем с плагином (Comet включён в конфиге)
print("\n--- С Comet Plugin ---")
comet_time = run_and_time("Comet", lambda: query)
# Отключаем Comet на лету (не перезапуская SparkSession)
# и запускаем без нативного движка
spark.conf.set("spark.comet.enabled", "false")
print("\n--- Без плагина (JVM Spark) ---")
jvm_time = run_and_time("JVM Spark", lambda: query)
# Включаем обратно
spark.conf.set("spark.comet.enabled", "true")
print(f"\nSpeedUp: {jvm_time / comet_time:.1f}×")
print(f"JVM Spark: {jvm_time:.2f}s")
print(f"Comet: {comet_time:.2f}s")
Диагностика: что происходит при fallback¶
# ──────────────────────────────────────────────────────────
# Диагностика fallback: почему оператор не пошёл в нативный движок
# ──────────────────────────────────────────────────────────
# Включаем детальное логирование Comet
spark.conf.set("spark.comet.logLevel", "DEBUG")
# Запрос с UDF - UDF не может быть нативным (Python UDF остаётся в JVM)
from pyspark.sql.types import DoubleType
@F.udf(returnType=DoubleType())
def custom_score(amount, quantity):
"""Кастомная Python-логика - остаётся в JVM."""
return amount * quantity * 1.1
query_with_udf = df \
.withColumn("score", custom_score("amount", "quantity")) \
.groupBy("region") \
.agg(F.sum("score").alias("total_score"))
print("Plan с UDF (часть в JVM, часть в нативном):")
query_with_udf.explain(extended=True)
# В логах Comet при DEBUG уровне:
# [DEBUG] CometSparkSessionExtensions: Operator BatchEvalPython not supported
# by Comet - falling back to JVM execution for this subtree
# [DEBUG] CometSparkSessionExtensions: Inserted RowToColumnarExec before
# CometHashAggregateExec to handle JVM→Native transition
# В плане это выглядит как:
# CometColumnarToRowExec
# CometHashAggregateExec ← нативная агрегация после UDF
# RowToColumnarExec ← конвертация JVM → Arrow перед нативным
# BatchEvalPython [custom_score] ← UDF остаётся в JVM
# CometFilterExec ← нативный filter перед UDF
# CometScanExec ← нативное чтение
# Это показывает "разрыв" плана из-за UDF:
# [нативный] → [JVM UDF] → [конвертация] → [нативный]
# Каждый разрыв = ColumnarToRow + RowToColumnar = дополнительный overhead
spark.conf.set("spark.comet.logLevel", "WARN")
Мониторинг нативной памяти¶
# ──────────────────────────────────────────────────────────
# Мониторинг Off-Heap памяти (нативной памяти плагина)
# ──────────────────────────────────────────────────────────
# Spark UI показывает Off-Heap использование в разделе Executor:
# http://localhost:4040/executors/
# Через Spark Metrics API:
def get_executor_memory_stats(spark):
"""Возвращает статистику памяти всех executors."""
sc = spark.sparkContext
status_tracker = sc.statusTracker()
# executorInfos() доступен через JVM через Py4J
executor_info = sc._jvm.org.apache.spark.SparkContext \
.getOrCreate() \
.statusStore() \
.executorList(True)
# В реальном коде: парсим ExecutorSummary из REST API
# GET http://localhost:4040/api/v1/applications/{appId}/executors
import urllib.request
import json
app_id = sc.applicationId
url = f"http://localhost:4040/api/v1/applications/{app_id}/executors"
try:
with urllib.request.urlopen(url) as response:
executors = json.loads(response.read())
for ex in executors:
print(f"Executor {ex['id']}:")
print(f" JVM Heap Used: {ex['memoryUsed'] / 1024 / 1024:.0f} MB")
print(f" Off-Heap Used: "
f"{ex.get('peakMemoryMetrics', {}).get('OffHeapMemoryUsed', 0) / 1024 / 1024:.0f} MB")
print(f" JVM Heap Max: "
f"{ex['maxMemory'] / 1024 / 1024:.0f} MB")
except Exception as e:
print(f"Spark UI not accessible: {e}")
get_executor_memory_stats(spark)
Риски и ограничения Plugin API¶
Самый опасный инструмент в Spark ecosystem¶
Plugin API - это самый мощный и одновременно самый опасный механизм расширения Spark. Понимание рисков критически важно.
Segmentation Fault роняет весь Executor. Ошибка в нативном C++/Rust коде (неправильный указатель, выход за границы массива, двойное освобождение памяти) приводит к SIGSEGV - Segmentation Fault. В JVM это немедленное завершение процесса Executor. Spark не может поймать SIGSEGV через try/catch - это не Java-исключение. Весь Executor падает, все его текущие Tasks теряются, Spark перезапускает их на других Executor-ах (или не может, если нет).
Утечки нативной памяти. GC Java не знает о памяти, аллоцированной через JNI. Если плагин не освобождает нативные буферы при ошибках (в onTaskFailed) или при ранней остановке (в shutdown), Executor медленно утекает по Off-Heap памяти. Симптом: Executor работает стабильно часами, потом внезапно падает по OOM на Off-Heap.
JNI-дедлоки. JNI запрещает вызов определённых JVM-методов из нативных потоков без явного прикрепления потока к JVM (AttachCurrentThread). Неправильное управление JNI-потоками приводит к дедлокам или JVM crashes.
Несовместимость версий. Plugin API - это JVM-уровень. Плагин, скомпилированный для Spark 3.4, может не работать с Spark 3.5 из-за изменений в internal API (физические операторы, AQE internals). Это заставляет plugin-авторов поддерживать отдельные сборки для каждой минорной версии Spark.
Матрица рисков по типам плагинов¶
| Тип плагина | Риск краша | Риск утечек | Сложность | Кто пишет |
|---|---|---|---|---|
| DriverPlugin (только Java/Scala) | Низкий | Низкий | Средняя | Команды платформы |
| ExecutorPlugin (только Java/Scala) | Средний | Средний | Высокая | Spark infrastructure инженеры |
| ExecutorPlugin + JNI (C++) | Высокий | Высокий | Очень высокая | Специализированные команды (RAPIDS, Meta) |
| ExecutorPlugin + JNI (Rust) | Средний | Средний | Очень высокая | Специализированные команды (Comet, Gluten Rust) |
Rust снижает риск по сравнению с C++ благодаря системе Ownership, которая устраняет большинство ошибок памяти на этапе компиляции. Но JNI-взаимодействие по-прежнему требует unsafe блоков в Rust.
Когда Plugin API - правильный выбор¶
Plugin API подходит для инфраструктурных задач уровня платформы, а не для бизнес-логики:
✅ Правильные применения:
- Execution acceleration (Gluten, RAPIDS, Comet) - замена JVM операторов на нативные
- GPU resource management - распределение GPU между Tasks
- Custom shuffle service - замена стандартного shuffle на disaggregated (например, Remote Shuffle Service)
- Custom metrics и observability - сбор детальных метрик выполнения
- Security instrumentation - аудит доступа к данным на уровне оператора
❌ Неправильные применения:
- Бизнес-логика обработки данных - для этого есть
@pandas_udfи SQL - Feature engineering - нет смысла писать нативный JNI-код для трансформаций, которые решает Spark SQL
- ETL-логика - Plugin API слишком тяжёлый инструмент
- Замена
DataFrame.filter()на кастомную фильтрацию - встроенные функции лучше
Пишем минимальный учебный плагин¶
Для закрепления напишем полный работающий плагин, который логирует жизненный цикл. Это именно тот скелет, на основе которого строятся реальные production-плагины.
Полный код плагина на Scala¶
// MySparkPlugin.scala
package com.example.plugin
import org.apache.spark.SparkContext
import org.apache.spark.api.plugin._
import org.slf4j.LoggerFactory
import java.util
// ── Главный класс плагина - точка регистрации ─────────────
// Должен быть в spark.plugins конфиге
class MySparkPlugin extends SparkPlugin {
override def driverPlugin(): DriverPlugin = new MyDriverPlugin()
override def executorPlugin(): ExecutorPlugin = new MyExecutorPlugin()
}
// ── Driver-сторона ─────────────────────────────────────────
class MyDriverPlugin extends DriverPlugin {
private val log = LoggerFactory.getLogger(classOf[MyDriverPlugin])
override def init(
sc: SparkContext,
pluginContext: PluginContext
): util.Map[String, String] = {
log.info("=== [Driver Plugin] INIT ===")
log.info("Spark master: {}", sc.master)
log.info("App name: {}", sc.appName)
log.info("Executor memory: {} bytes", sc.executorMemory)
// Инициализируем Driver-сторонние ресурсы
// (в реальном плагине: нативные библиотеки, connection pools)
// Передаём конфигурацию на все Executor-ы
val config = new util.HashMap[String, String]()
config.put("plugin.driver.initialized.at",
System.currentTimeMillis().toString)
config.put("plugin.driver.app.name", sc.appName)
config.put("plugin.native.enabled", "true")
config
}
override def receive(message: AnyRef): AnyRef = {
log.info("=== [Driver Plugin] Received from Executor: {} ===", message)
// Возвращаем ответ Executor-у (опционально)
"ACK:" + message
}
override def shutdown(): Unit = {
log.info("=== [Driver Plugin] SHUTDOWN ===")
// Освобождаем Driver-сторонние ресурсы
}
}
// ── Executor-сторона ───────────────────────────────────────
class MyExecutorPlugin extends ExecutorPlugin {
private val log = LoggerFactory.getLogger(classOf[MyExecutorPlugin])
private var pluginContext: PluginContext = _
private var executorId: String = "unknown"
private var taskCount: Long = 0
override def init(
pluginContext: PluginContext,
extraConf: util.Map[String, String]
): Unit = {
this.pluginContext = pluginContext
// Получаем ID этого Executor (если доступен через контекст)
log.info("=== [Executor Plugin] INIT ===")
// Читаем конфиг, переданный DriverPlugin
val initTime = extraConf.get("plugin.driver.initialized.at")
val appName = extraConf.get("plugin.driver.app.name")
val nativeEnabled = extraConf.get("plugin.native.enabled")
log.info("Config from Driver: appName={}, nativeEnabled={}, driverInitTime={}",
appName, nativeEnabled, initTime)
// Инициализируем Executor-сторонние ресурсы
// В реальном плагине: System.loadLibrary("native_plugin_lib")
log.info("Native runtime initialized (simulated)")
// Отправляем сообщение на Driver
pluginContext.send("EXECUTOR_READY:executor-" + Thread.currentThread().getId)
}
override def onTaskStart(): Unit = {
taskCount += 1
log.debug("[Executor Plugin] Task #{} starting", taskCount)
}
override def onTaskSucceeded(): Unit = {
log.debug("[Executor Plugin] Task #{} succeeded", taskCount)
}
override def onTaskFailed(failureReason: TaskFailedReason): Unit = {
log.warn("[Executor Plugin] Task #{} FAILED: {}", taskCount, failureReason)
// КРИТИЧЕСКИ ВАЖНО: освободить нативные ресурсы при ошибке
// В реальном плагине: freeNativeTaskResources()
}
override def shutdown(): Unit = {
log.info("[Executor Plugin] SHUTDOWN after {} tasks", taskCount)
// Освобождаем нативный контекст
// В реальном плагине: destroyNativeContext()
}
}
Для использования этого плагина нужно:
- Скомпилировать в JAR:
sbt package(или Maven) - Добавить JAR в classpath Spark
- Указать класс в конфиге
# Использование учебного плагина из Python
spark = SparkSession.builder \
.master("local[2]") \
.config("spark.driver.extraClassPath", "/path/to/my-plugin.jar") \
.config("spark.executor.extraClassPath", "/path/to/my-plugin.jar") \
.config("spark.plugins", "com.example.plugin.MySparkPlugin") \
.config("spark.driver.extraJavaOptions",
"-Dlog4j.logger.com.example.plugin=DEBUG") \
.getOrCreate()
# Запускаем запрос - плагин логирует каждый шаг
df = spark.range(100000).groupBy((F.col("id") % 10).alias("group")).count()
df.show()
# В логах увидим:
# [Driver Plugin] INIT
# [Executor Plugin] INIT (для каждого Executor)
# [Executor Plugin] Received: EXECUTOR_READY:executor-123
# [Executor Plugin] Task #1 starting
# [Executor Plugin] Task #1 succeeded
# ...
# [Driver Plugin] SHUTDOWN
# [Executor Plugin] SHUTDOWN after N tasks
Домашнее задание¶
Задание 1: Анализ плана с нативным плагином (аналитическое)¶
Установите Apache Comet (JAR доступен на GitHub releases) или используйте готовую Spark-дистрибуцию с RAPIDS/Gluten в облаке. Выполните следующий запрос с включённым и выключенным плагином:
SELECT
region,
category,
SUM(amount) as total,
COUNT(*) as orders,
AVG(amount) as avg_amount,
RANK() OVER (PARTITION BY region ORDER BY SUM(amount) DESC) as rank
FROM transactions
WHERE amount > 0
GROUP BY region, category
Получите explain(extended=True) для обоих случаев. Напишите отчёт:
- Какие операторы заменены нативными (с именами нативных операторов)
- Где остался JVM-код и почему
- Где появились ColumnarToRow/RowToColumnar переходы
- Как изменилось время выполнения (в секундах)
Задание 2: Диагностика OOM по логу (аналитическое)¶
Вам дан лог Spark-приложения с установленным Gluten-плагином, в котором произошёл OOM на Executor:
[INFO] GlutenExecutorPlugin: Initialized Velox runtime, offHeap=4g
[INFO] GlutenExecutorPlugin: Task#142 starting (Stage 3, queryStage=5)
[INFO] VeloxHashAggregateExec: Processing partition 0/8, inputRows=14_500_000
[INFO] VeloxHashJoinExec: Building hash table, rightSide=2_100_000 rows
[WARN] VeloxMemoryPool: Memory pressure: 92% of offHeap used
[INFO] VeloxHashJoinExec: Building hash table, rightSide=2_100_000 rows (RETRY)
[ERROR] VeloxMemoryPool: OOM - failed to allocate 512MB, available=180MB
[ERROR] GlutenExecutorPlugin: Native OOM in Task#142, Stage 3 —
executor will be terminated
[ERROR] KubernetesExternalClusterManager: Executor pod killed (OOMKilled)
Ответьте на вопросы:
- На каком физическом операторе произошёл OOM?
- Сколько памяти было сконфигурировано и сколько требовалось?
- Какие параметры конфигурации нужно изменить, чтобы исправить проблему? Укажите конкретные
spark.*илиspark.gluten.*параметры с объяснением. - Как определить оптимальный размер
spark.memory.offHeap.sizeдля этой нагрузки на будущее?
Задание 3: Матрица выбора плагина (проектировочное)¶
Для каждого из следующих сценариев выберите подходящий нативный плагин (Gluten+Velox, RAPIDS cuDF, Apache Comet, или никакого) и обоснуйте выбор:
Сценарий A: SQL-аналитика на CPU-кластере (AWS r6g Graviton ARM), 10 TB/день, тяжёлые JOIN'ы и GROUP BY, бюджет ограничен (CPU-инстансы дешевле GPU).
Сценарий B: ML feature engineering + XGBoost training, 500 GB/день, очень широкие таблицы (500 колонок), инфраструктура с NVIDIA A100 GPU.
Сценарий C: Streaming pipeline (Structured Streaming), Kafka → S3, 100k events/sec, низкая latency критична, никаких изменений в существующем коде.
Сценарий D: Data Quality checks - 200 тяжёлых SQL-запросов с GROUP BY и COUNT DISTINCT, 1 TB данных, команда работает только с Python, нет Scala/Java разработчиков.
Резюме¶
Spark Plugin API - это легальный механизм для фундаментальной замены JVM-движка Spark без форка исходного кода. Он создан для одной цели: позволить специализированным командам интегрировать нативные acceleration engines (GPU через RAPIDS, CPU vectorized через Velox/DataFusion) в Spark-экосистему.
Архитектура двусторонняя: DriverPlugin управляет инициализацией на стороне Driver и координирует Executor-ы через RPC, ExecutorPlugin перехватывает жизненный цикл задач на каждом воркере и инициализирует нативные рантаймы до запуска первой задачи. Замена физических операторов происходит через ColumnarRule (внедряется через SparkSessionExtensions), который обходит дерево физического плана и заменяет JVM-операторы на нативные аналоги.
Ключевое свойство правильно написанного плагина - graceful degradation: неподдерживаемые операторы оставляются JVM-реализации вместо ошибки. Это обеспечивает 100% корректность запросов даже при неполном покрытии.
Plugin API не предназначен для бизнес-логики. Если вы думаете о написании плагина для ETL-трансформаций - рассмотрите @pandas_udf или Spark SQL. Плагины - это инструмент для инфраструктурных команд, которые знают Spark internals на уровне физических операторов, управляют JNI-взаимодействием и готовы поддерживать совместимость с каждым минорным релизом Spark.