Spark Plugin API: SparkPlugin и QueryStage интерфейс для замены executor

Как Spark Plugin API позволяет заменять JVM-выполнение нативным кодом: DriverPlugin, ExecutorPlugin, ColumnarRules, AQE QueryStage, архитектура Gluten и RAPIDS.

optimization

Зачем Spark открыл внутренности: исторический контекст

До появления Spark Plugin API в версии 3.0 разработчики, которые хотели изменить поведение Spark на уровне выполнения, оказывались перед неприятным выбором: либо форкать Spark и поддерживать собственный патченый рантайм (что делает Facebook, Alibaba и другие гиганты - дорого и сложно), либо мириться с ограничениями JVM-движка.

Доступные точки расширения до 3.0 работали только выше уровня выполнения:

  • DataSource V1/V2 - только для чтения и записи данных, не для изменения логики вычислений
  • Catalyst rules через SparkSessionExtensions - оптимизация плана на уровне SQL, но не замена физических операторов
  • --jars и дополнительные классы - добавление кода в classpath, но без доступа к жизненному циклу executors
  • UDF - бизнес-логика поверх движка, но не замена самого движка

Ни один из этих механизмов не позволял сказать 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()
  }
}

Для использования этого плагина нужно:

  1. Скомпилировать в JAR: sbt package (или Maven)
  2. Добавить JAR в classpath Spark
  3. Указать класс в конфиге
# Использование учебного плагина из 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.