Apache DataFusion Comet: архитектура Rust + Arrow, как подключается к Spark

Apache Comet как open-source Rust-ускоритель Spark: DataFusion внутри, Arrow zero-copy, Plugin API интеграция, JVM fallback и TPC-H бенчмарки.

optimization

Концепция и предпосылки: зачем ещё один нативный движок

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

  • Photon (Databricks) - C++ engine, заменяет execution layer целиком, только в Databricks Runtime
  • Gluten + Velox (Intel + Meta) - C++ плагин, заменяет физические операторы через Plugin API
  • Sail - Rust engine, заменяет JVM полностью через Spark Connect

У каждого из них есть фундаментальный недостаток для open-source мира: Photon проприетарный и недоступен вне Databricks, Gluten+Velox требует C++ toolchain и остаётся сложным в развёртывании, Sail не совместим с RDD API и находится в alpha-статусе.

Apache DataFusion Comet появился как ответ на конкретный запрос: дайте open-source-пользователям ванильного Apache Spark нативное ускорение без vendor lock-in, без сложных C++ зависимостей, без необходимости переходить на другой API.

Коmet - инкубируемый проект Apache Software Foundation (передан из Databricks в 2023 году как open-source). Он реализует ровно ту идею, которую мы изучали в уроке про Spark Plugin API: использует Plugin API и SparkSessionExtensions для замены JVM физических операторов на Rust-нативные реализации на основе Apache DataFusion.

Проблема «нативной интеграции» в Enterprise

Enterprise-компании оказываются в сложном положении: с одной стороны, они видят 3–10× разницу в производительности между JVM Spark и нативными движками. С другой стороны, у них есть реальные ограничения:

  • Тысячи строк PySpark-кода, который нельзя переписать быстро - любая миграция на новый API (DuckDB, Sail) требует месяцев работы
  • Зависимость от Spark-экосистемы - Airflow DAGs, Spark History Server, Spark UI, интеграции с Hive Metastore, YARN/K8s операторы - всё это построено вокруг Spark
  • Требования стабильности - production-кластеры обрабатывают критические данные, экспериментальные движки неприемлемы
  • C++ рисковы - Segmentation Fault в нативном коде роняет весь Executor без шанса на retry (как мы изучали в уроке про Spark Plugin API)

Comet решает все эти ограничения через инкрементальную стратегию: ваш PySpark/Spark SQL код не меняется ни на строчку. Comet подключается как JAR-плагин, берёт поддерживаемые операторы из физического плана Catalyst и выполняет их через DataFusion на Rust. Операторы, которые Comet не поддерживает, продолжают работать через JVM Tungsten - без ошибок.

Почему Rust, а не C++

Выбор Rust вместо C++ (как в Velox) - не случаен. Рассмотрим конкретные последствия:

Memory Safety без GC. Rust гарантирует отсутствие dangling pointers, double-free, buffer overflow на этапе компиляции. В C++ эти ошибки приводят к Segmentation Fault в production - и роняют весь JVM-процесс Executor'а. Rust-ошибки превращаются либо в compile error (до production), либо в controlled panic (который Spark может обработать как Task failure, а не Executor crash).

DataFusion как готовая база. Apache DataFusion - это зрелый Rust query engine, уже используемый в InfluxDB IOx, Delta Lake Rust reader, GlareDB, Apache Comet, Ballista. Команде Comet не нужно писать query engine с нуля - DataFusion предоставляет готовые векторизованные операторы (HashAggregate, SortMergeJoin, BroadcastHashJoin, Filter, Project) на Apache Arrow.

Apache Arrow как универсальный стандарт. Rust DataFusion нативно работает с Arrow RecordBatch. Spark 3.x тоже поддерживает Arrow через ColumnarBatch API. Это создаёт естественный мост между JVM и Rust без дополнительного протокола.


Архитектура: как Rust уживается с JVM

Самый важный вопрос практического использования Comet: как именно Rust-код исполняется внутри JVM-процесса Spark Executor? Ответ - через многоуровневую архитектуру с чётким разделением ответственности.

Двухслойный рантайм: JVM координирует, Rust вычисляет

На схеме видна ключевая особенность Comet: Rust DataFusion runtime работает внутри того же OS-процесса, что и JVM Executor. Это не отдельный микросервис, не sidecar - это нативная библиотека (.so файл на Linux, .dylib на macOS), загруженная в JVM-процесс через System.loadLibrary(). JNI-вызовы происходят внутри одного процесса - без network overhead, без копирования данных через socket.

Apache DataFusion: Rust query engine под капотом

Apache DataFusion - это Rust-нативный SQL query engine, являющийся частью Apache Arrow экосистемы. Он реализует полный стек аналитического выполнения:

Парсер и планировщик. DataFusion может принимать SQL строки и строить Logical Plan. Но Comet использует его иначе: Catalyst строит физический план, Comet транслирует его в DataFusion Physical Plan через protobuf сериализацию, передаёт через JNI в Rust, DataFusion десериализует и выполняет.

Векторизованные операторы. DataFusion реализует Volcano iterator model, но в батчевом варианте: каждый оператор получает RecordBatch (аналог Spark's ColumnarBatch) и возвращает RecordBatch. Внутри каждого оператора - loop по Arrow arrays с SIMD-инструкциями.

Встроенные функции. DataFusion имеет сотни встроенных функций (арифметика, строки, даты, агрегаты), реализованных на Rust с Arrow vectorization. Comet маппирует Spark SQL функции на соответствующие DataFusion функции.

Parquet reader. DataFusion содержит нативный Parquet reader, написанный на Rust (использует crate parquet). Он читает columnar данные прямо в Arrow buffers, применяет predicate pushdown через Parquet statistics и Dictionary encoding optimization.

Apache Arrow: единый язык памяти

Центральная роль Arrow в архитектуре Comet - быть форматом нулевого копирования между JVM и Rust. Рассмотрим, что именно происходит:

  1. Spark Executor получает Task с физическим планом
  2. CometExec оператор (JVM Scala-класс) вызывается Spark'ом
  3. Внутри CometExec.executeColumnar() происходит JNI-вызов в Rust
  4. Rust-код создаёт DataFusion physical plan из protobuf
  5. DataFusion читает Parquet напрямую в Arrow buffers - минуя JVM heap полностью
  6. DataFusion выполняет операторы (Filter, Aggregate, Join) над Arrow buffers
  7. Результат - Arrow RecordBatch в нативной памяти
  8. JNI возвращает в JVM long-указатель на этот RecordBatch
  9. JVM оборачивает указатель в ColumnarBatch без копирования данных

Шаг 9 критически важен: JVM не копирует данные из нативной памяти в JVM heap. Вместо этого ColumnarBatch в JVM содержит wrapper, который смотрит на ту же нативную память через ArrowColumnVector (который поддерживает Off-Heap доступ). Это и есть zero-copy: данные остаются в одном месте в памяти, а Java и Rust просто имеют разные "взгляды" на них.

Protobuf как язык трансляции планов

Механизм передачи физического плана от Catalyst к DataFusion заслуживает отдельного разбора. Catalyst строит дерево SparkPlan объектов на Scala. Rust не умеет читать Scala-объекты напрямую. Comet решает это через Protocol Buffers:

  1. CometSparkSessionExtensions обходит физический план Catalyst
  2. Поддерживаемые узлы сериализуются в protobuf схему (.proto файлы описывают операторы: CometFilter, CometAggregate, CometJoin, ...)
  3. Protobuf байты передаются через JNI в Rust как byte[] массив
  4. Rust десериализует protobuf в DataFusion physical plan объекты (используя crate prost)
  5. DataFusion выполняет план

Protobuf - универсальный бинарный формат с поддержкой в обоих языках. Это решение надёжнее, чем, например, передача плана как JSON строки: protobuf компактнее, быстрее десериализуется, строго типизирован.


Механизм подключения: Spark Plugin API и Extensions

Comet использует тот же Plugin API, который мы изучали в уроке 5, в сочетании со SparkSessionExtensions для замены физических операторов.

Конфигурация: минимальный набор для активации

# ──────────────────────────────────────────────────────────
# Минимальная конфигурация для включения Comet
# ──────────────────────────────────────────────────────────
# Предполагается, что comet-spark-spark3.5_2.12-X.Y.Z.jar
# добавлен в classpath через --jars или extraClassPath

spark = SparkSession.builder \
    .master("local[4]") \
    # Spark Plugin API: регистрируем Comet runtime plugin.
    # Этот класс реализует DriverPlugin и ExecutorPlugin.
    # На каждом Executor при старте загружается native lib.
    .config("spark.plugins", "org.apache.spark.CometPlugin") \
    # SparkSessionExtensions: регистрируем правила оптимизатора.
    # Этот класс внедряет ColumnarRule для замены JVM операторов.
    .config("spark.sql.extensions",
            "org.apache.spark.CometSparkSessionExtensions") \
    # Главный флаг включения Comet
    .config("spark.comet.enabled", "true") \
    # Разрешаем Comet перехватывать все поддерживаемые операторы:
    # JOIN, Aggregate, Sort, Window и не только фильтры и проекции.
    # По умолчанию false - только базовые операторы.
    .config("spark.comet.exec.all.enabled", "true") \
    # Нативный Parquet reader на Rust (Apache Arrow Rust crate).
    # Читает данные напрямую в Arrow buffers без JVM десериализации.
    .config("spark.comet.parquet.enable.vectorized.reader", "true") \
    # Arrow shuffle: передача результатов между стадиями
    # через Arrow IPC вместо Spark serialization.
    # Даёт дополнительный 10–30% прирост на shuffle-heavy запросах.
    .config("spark.comet.exec.shuffle.enabled", "true") \
    .getOrCreate()

# Проверка активации:
comet_enabled = spark.conf.get("spark.comet.enabled", "false")
print(f"Comet enabled: {comet_enabled}")

# В логах Executor при успешной инициализации:
# [INFO] CometPlugin: Loading native library from resources...
# [INFO] CometPlugin: Native library loaded successfully
# [INFO] CometPlugin: Comet native runtime initialized

Обратите внимание на два раздельных механизма:

spark.plugins=org.apache.spark.CometPlugin - это инфраструктурный hook. Он инициализирует Rust native library на каждом Executor при старте, настраивает нативный allocator для Arrow buffers, регистрирует метрики.

spark.sql.extensions=org.apache.spark.CometSparkSessionExtensions - это оптимизационный hook. Он добавляет правила в Catalyst pipeline для замены JVM операторов на Comet операторы. Без этого плагин загружен, но план не изменится.

Жизненный цикл запроса с Comet: пошаговый разбор

Проследим полный путь запроса SELECT region, SUM(amount) FROM orders WHERE amount > 1000 GROUP BY region:

Шаг 1: Catalyst строит стандартный план

Logical Plan:
Aggregate [region], [sum(amount)]
  Filter (amount > 1000)
    Relation orders [order_id, region, amount, ...]

Physical Plan (без Comet):
*(2) HashAggregate(mode=Final, keys=[region])
  Exchange hashpartitioning(region, 200)
    *(1) HashAggregate(mode=Partial, keys=[region], funcs=[sum(amount)])
      *(1) Filter (amount > 1000)
        *(1) ColumnarToRow     ← Arrow → JVM Row
          FileScan parquet orders

Шаг 2: CometSparkSessionExtensions применяет ColumnarRule

// Упрощённая логика замены операторов (реальный код в Comet source)
class CometColumnarRule(session: SparkSession) extends ColumnarRule {
  override def preColumnarTransitions(plan: SparkPlan): SparkPlan = {
    plan.transformDown {
      // HashAggregateExec → CometHashAggregateExec
      case agg: HashAggregateExec if CometPlanChecker.isSupported(agg) =>
        CometHashAggregateExec(agg)

      // FilterExec → CometFilterExec
      case filter: FilterExec if CometPlanChecker.isSupported(filter) =>
        CometFilterExec(filter)

      // ParquetFileFormat scan → CometScanExec (native Parquet reader)
      case scan: FileSourceScanExec
          if scan.relation.fileFormat.isInstanceOf[ParquetFileFormat]
          && spark.conf.get("spark.comet.parquet.enable.vectorized.reader") == "true" =>
        CometScanExec(scan)

      // Все остальные - оставляем JVM (graceful fallback)
    }
  }
}

Шаг 3: Финальный план с Comet

Physical Plan (с Comet):
CometResultStage                      ← возврат результата из Comet в Spark
  CometHashAggregate(mode=Final)      ← нативная финальная агрегация
    CometColumnarExchange              ← Arrow shuffle между стадиями
      CometHashAggregate(mode=Partial) ← нативная частичная агрегация
        CometFilter (amount > 1000)   ← нативный filter (AVX-512 SIMD)
          CometScanExec (parquet)     ← нативный Parquet reader

Шаг 4: Выполнение CometScanExec на Executor

1. Spark Task Scheduler вызывает CometScanExec.executeColumnar()
2. CometScanExec через JNI передаёт в Rust:
   - путь к Parquet файлу
   - схему (required columns)
   - protobuf serialized physical plan для этого subplan
3. Rust DataFusion открывает Parquet файл напрямую
4. Читает columnar data в Arrow buffers (Off-Heap)
5. Применяет row group pruning по statistics (amount_max > 1000)
6. DataFusion применяет Filter, Aggregate векторизованно
7. Возвращает Arrow RecordBatch через JNI
8. JVM получает pointer, оборачивает в ColumnarBatch
9. Spark передаёт ColumnarBatch на shuffle через CometColumnarExchange

Partial Fallback: что происходит с неподдерживаемыми операторами

Comet поддерживает большой, но не исчерпывающий набор операторов. Когда в плане встречается неподдерживаемый оператор, Comet выполняет partial fallback:

Запрос с Java UDF:

Physical Plan:
CometResultStage
  CometHashAggregate(mode=Final)
    CometColumnarExchange
      CometHashAggregate(mode=Partial)
        RowToColumnar                  ← JVM Row → Arrow (дорогой переход!)
          BatchEvalPython              ← Python UDF в JVM (НЕ Comet)
            CometFilter                ← до UDF - нативный
              CometScanExec

Comet автоматически вставил:
  1. RowToColumnar после BatchEvalPython  - конвертация JVM → Arrow
  2. Продолжил план с нативными операторами выше

Этот механизм - ключевое преимущество Comet перед Sail. Если Sail встречает неподдерживаемый оператор (RDD API) - выдаёт ошибку. Comet - молча деградирует до JVM для этого узла и продолжает нативное выполнение для остальных.


Поддерживаемые операторы и ограничения

Что Comet поддерживает нативно

Актуальный список операторов на момент написания (версия 0.4+):

Категория Операторы Статус
Scan ParquetScan (все типы кроме некоторых complex types) ✅ Полная
Filter Все predicate типы, NOT NULL, IN, BETWEEN ✅ Полная
Project Арифметика, строковые функции, CASE WHEN, CAST ✅ Полная
Aggregation SUM, COUNT, AVG, MIN, MAX, GROUP BY ✅ Полная
Joins BroadcastHashJoin, SortMergeJoin, HashJoin ✅ Полная
Sort ORDER BY (single/multiple keys) ✅ Полная
Window RANK, ROW_NUMBER, SUM OVER, LAG/LEAD ✅ Полная
Shuffle Hash partitioning с Arrow IPC ✅ Через CometShuffleManager
Limit LIMIT, OFFSET ✅ Полная
Union/Except UNION ALL, EXCEPT 🔄 Частичная

Что Comet НЕ поддерживает (fallback в JVM)

Тип ограничения Детали
Python UDF / Scala UDF Любой UDF уходит в JVM, Comet продолжает до/после
Complex types в Parquet MapType, вложенные ArrayType → JVM Parquet reader
RDD API Полностью вне Comet scope
Structured Streaming Micro-batch работает, но native path не полный
Некоторые Cast Нетривиальные приведения типов могут вызвать fallback
LATERAL VIEW / explode Ограниченная поддержка

Диагностика через логи

# ──────────────────────────────────────────────────────────
# Уровень логирования для диагностики Comet
# ──────────────────────────────────────────────────────────
# Детальные логи покажут какие операторы взял Comet,
# какие упали в JVM fallback
spark.sparkContext.setLogLevel("WARN")
spark.conf.set("spark.comet.logLevel", "DEBUG")

# При включённом DEBUG в логах Executor будет:
# [DEBUG] CometSparkSessionExtensions: Replacing HashAggregateExec
#         with CometHashAggregateExec - supported
# [DEBUG] CometSparkSessionExtensions: BatchEvalPython not supported
#         by Comet - keeping JVM operator, inserting RowToColumnar
# [DEBUG] CometScanExec: Reading /data/orders.parquet with native reader
# [DEBUG] CometFilter: Applied vectorized filter, 850000/1000000 rows passed
# [INFO]  CometHashAggregate: Partial aggregate completed,
#         output 3200 groups in 0.8s

# Вернуть логи к нормальному уровню после диагностики
spark.conf.set("spark.comet.logLevel", "WARN")

Тюнинг: ключевые параметры конфигурации

Правильная настройка Comet критически важна для производительности и стабильности. Рассмотрим все важные параметры.

Группа 1: Базовое включение

spark = SparkSession.builder \
    # ── Обязательные параметры ─────────────────────────────
    .config("spark.plugins", "org.apache.spark.CometPlugin") \
    .config("spark.sql.extensions",
            "org.apache.spark.CometSparkSessionExtensions") \
    .config("spark.comet.enabled", "true") \

    # ── Расширенные операторы ──────────────────────────────
    # false = только фильтры и проекции (безопасно, меньший прирост)
    # true  = все операторы включая JOIN, Aggregate, Window
    .config("spark.comet.exec.all.enabled", "true") \

    # ── Native Parquet reader ──────────────────────────────
    # Rust Arrow Parquet reader вместо JVM VectorizedReader.
    # Дополнительный 15–30% прирост на scan-heavy запросах.
    # Может не поддерживать некоторые Parquet feature (legacy formats).
    .config("spark.comet.parquet.enable.vectorized.reader", "true") \

    # ── Native shuffle ─────────────────────────────────────
    # Заменяет Spark shuffle на Arrow IPC transfer.
    # Критично для запросов с большим shuffle (JOIN, GROUP BY по высококардинальным полям).
    # Требует CometShuffleManager:
    .config("spark.comet.exec.shuffle.enabled", "true") \
    .config("spark.shuffle.manager",
            "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") \
    .getOrCreate()

Группа 2: Управление памятью (критически важно!)

spark = SparkSession.builder \
    # ── Off-Heap память для Rust allocator ────────────────
    # Comet аллоцирует Arrow buffers в нативной Off-Heap памяти.
    # Эта память НЕ видна JVM - она не считается в spark.executor.memory.
    # Если не выделить дополнительный overhead, YARN/Kubernetes
    # может убить контейнер по превышению лимита OS-памяти!
    #
    # spark.comet.memoryOverhead - дополнительная память ЗА ПРЕДЕЛАМИ
    # spark.executor.memory + spark.executor.memoryOverhead.
    # Рекомендуется: 15–25% от spark.executor.memory.
    .config("spark.comet.memoryOverhead", "2g") \
    # Или в относительном виде:
    # .config("spark.comet.memoryOverheadFactor", "0.2") \  # 20% от executor memory

    # ── Размер columnar батча ──────────────────────────────
    # Количество строк в одном Arrow RecordBatch.
    # Default: 8192 (хороший баланс throughput vs memory).
    # Больший батч → лучше SIMD utilization, хуже для маленьких данных.
    # Меньший батч → меньше memory footprint, больше overhead на батчинг.
    .config("spark.comet.batchSize", "8192") \

    # ── Limiter для нативных операций ─────────────────────
    # Максимальный объём Off-Heap памяти на один Executor
    # для нативных операций (sort spill, hash join build side).
    # Default: выводится из spark.executor.memory.
    # Увеличить если видите native OOM при тяжёлых JOIN.
    .config("spark.comet.memory.overhead.enabled", "true") \
    .getOrCreate()

# ── Полный пример конфигурации для production кластера ────
# Executor: 32 GB RAM, 8 cores, spark.executor.memory=24g
COMET_PRODUCTION_CONFIG = {
    "spark.plugins": "org.apache.spark.CometPlugin",
    "spark.sql.extensions": "org.apache.spark.CometSparkSessionExtensions",
    "spark.comet.enabled": "true",
    "spark.comet.exec.all.enabled": "true",
    "spark.comet.parquet.enable.vectorized.reader": "true",
    "spark.comet.exec.shuffle.enabled": "true",
    "spark.shuffle.manager":
        "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager",
    # Memory: 24g executor + 4g JVM overhead (default) + 5g Comet native
    # = 33g total. Настройте лимиты контейнера соответственно.
    "spark.comet.memoryOverhead": "5g",
    "spark.comet.batchSize": "8192",
    # Off-Heap включён обязательно (иначе Arrow buffers в JVM heap)
    "spark.memory.offHeap.enabled": "true",
    "spark.memory.offHeap.size": "4g",
}

Группа 3: Тонкая настройка конкретных операторов

spark = SparkSession.builder \
    # ── Избирательное включение операторов ────────────────
    # Если exec.all.enabled=false, можно включить отдельные категории:
    .config("spark.comet.exec.hashagg.enabled", "true") \    # GROUP BY
    .config("spark.comet.exec.sort.enabled", "true") \       # ORDER BY
    .config("spark.comet.exec.join.enabled", "true") \       # JOIN
    .config("spark.comet.exec.window.enabled", "true") \     # Window functions
    .config("spark.comet.exec.filter.enabled", "true") \     # WHERE predicates

    # ── Broadcast Join threshold ───────────────────────────
    # Comet использует тот же broadcast threshold что и Spark,
    # но нативный broadcast join быстрее JVM версии.
    # Можно увеличить порог для большего использования BHJ с Comet.
    .config("spark.sql.autoBroadcastJoinThreshold", "20m") \

    # ── Совместимость с AQE ────────────────────────────────
    # AQE работает поверх Comet (не вместо него).
    # Comet-операторы корректно репортят статистику обратно в AQE.
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.coalescePartitions.enabled", "true") \

    # ── Fallback поведение ─────────────────────────────────
    # Если true - при любой ошибке Comet полностью отключается для запроса.
    # Если false (default) - Comet деградирует до JVM для проблемного узла.
    # Рекомендуется false для production (более стабильно).
    .config("spark.comet.fallback.enabled", "true") \
    .getOrCreate()

Параметры, влияющие на совместимость

# ──────────────────────────────────────────────────────────
# Параметры совместимости: когда результаты могут отличаться
# ──────────────────────────────────────────────────────────

spark = SparkSession.builder \
    # Floating point ordering: Spark и DataFusion могут давать
    # разные результаты при NaN в данных.
    # Включите для строгой совместимости (немного медленнее).
    .config("spark.comet.compatibleWithSpark.enabled", "true") \

    # Decimal precision: поведение при overflow может отличаться.
    # Включите если работаете с финансовыми расчётами в Decimal.
    .config("spark.comet.decimalShuffleEnabled", "true") \

    # Timestamp timezone: Comet обрабатывает timezone иначе в некоторых случаях.
    # Рекомендуется явно проверить timestamp-результаты при миграции.
    # Нет прямого параметра - диагностируйте через compare результатов.
    .getOrCreate()

Установка и развёртывание

Вариант 1: Скачать готовый JAR

# ──────────────────────────────────────────────────────────
# Скачиваем Comet JAR с GitHub releases
# ──────────────────────────────────────────────────────────
# Версии: https://github.com/apache/datafusion-comet/releases
# JAR уже содержит скомпилированную native библиотеку
# (libcomet.so на Linux, libcomet.dylib на macOS)

# Пример для Spark 3.5 / Scala 2.12
wget https://github.com/apache/datafusion-comet/releases/download/\
v0.4.0/comet-spark-spark3.5_2.12-0.4.0.jar

# Запуск spark-submit с Comet:
spark-submit \
    --jars /path/to/comet-spark-spark3.5_2.12-0.4.0.jar \
    --conf spark.plugins=org.apache.spark.CometPlugin \
    --conf spark.sql.extensions=org.apache.spark.CometSparkSessionExtensions \
    --conf spark.comet.enabled=true \
    --conf spark.comet.exec.all.enabled=true \
    --conf spark.comet.parquet.enable.vectorized.reader=true \
    my_spark_job.py

Вариант 2: Через pip (для разработки)

# ──────────────────────────────────────────────────────────
# pip install для разработки (Python API + JAR bundled)
# ──────────────────────────────────────────────────────────
pip install apache-comet

# После установки JAR доступен через:
import comet
print(comet.get_comet_jar_path())  # путь к JAR файлу

Вариант 3: Компиляция под конкретный CPU (максимальная производительность)

# ──────────────────────────────────────────────────────────
# Компиляция с target-cpu=native для максимального SIMD
# ──────────────────────────────────────────────────────────
# Требует: Rust toolchain + JDK 11+ + Maven 3.6+

git clone https://github.com/apache/datafusion-comet.git
cd datafusion-comet

# Компилируем с оптимизациями под текущий CPU
export RUSTFLAGS="-C target-cpu=native"
make release

# Результат: target/comet-spark-spark3.5_2.12-X.Y.Z.jar
# С поддержкой AVX-512 если CPU это поддерживает

# Проверяем что native lib скомпилирован с нужными инструкциями:
objdump -d target/native/release/libcomet.so | grep -c "ymm\|zmm"
# ymm = AVX2, zmm = AVX-512

Docker-compose для локального тестирования

# docker-compose.yml для тестирования Comet с MinIO
version: "3.8"

services:
  minio:
    image: minio/minio:latest
    command: server /data --console-address ":9001"
    ports:
      - "9000:9000"
      - "9001:9001"
    environment:
      MINIO_ROOT_USER: minioadmin
      MINIO_ROOT_PASSWORD: minioadmin
    volumes:
      - minio_data:/data

  spark-comet:
    image: apache/spark:3.5.1-python3
    # Пример: запуск PySpark с Comet JAR
    environment:
      SPARK_NO_DAEMONIZE: "true"
    volumes:
      # Монтируем Comet JAR и данные
      - ./comet-spark-spark3.5_2.12-0.4.0.jar:/opt/comet/comet.jar
      - ./data:/data
      - ./jobs:/jobs
    command: >
      /opt/spark/bin/spark-submit
        --jars /opt/comet/comet.jar
        --conf spark.plugins=org.apache.spark.CometPlugin
        --conf spark.sql.extensions=org.apache.spark.CometSparkSessionExtensions
        --conf spark.comet.enabled=true
        --conf spark.comet.exec.all.enabled=true
        --conf spark.hadoop.fs.s3a.endpoint=http://minio:9000
        --conf spark.hadoop.fs.s3a.access.key=minioadmin
        --conf spark.hadoop.fs.s3a.secret.key=minioadmin
        --conf spark.hadoop.fs.s3a.path.style.access=true
        /jobs/benchmark.py

volumes:
  minio_data:

Практика: развёртывание и бенчмарк

Сценарий: ускорение Silver-слоя с аналитикой KPI

Рассмотрим реальный бизнес-кейс: у аналитической команды есть ежедневный расчёт KPI витрины (Silver-слой) - тяжёлый GROUP BY с несколькими JOIN, который занимает 45 минут на стандартном Spark кластере.

# ──────────────────────────────────────────────────────────
# benchmark.py - запуск одного и того же запроса
# с Comet и без него, сравнение результатов
# ──────────────────────────────────────────────────────────
import time
import os
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.window import Window

MINIO_ENDPOINT = "http://localhost:9000"
COMET_JAR = os.environ.get("COMET_JAR", "/opt/comet/comet.jar")

def create_session(use_comet: bool) -> SparkSession:
    """Создаёт SparkSession с Comet или без."""
    builder = SparkSession.builder \
        .master("local[4]") \
        .config("spark.driver.memory", "4g") \
        .config("spark.executor.memory", "8g") \
        .config("spark.memory.offHeap.enabled", "true") \
        .config("spark.memory.offHeap.size", "2g") \
        # S3A-совместимый доступ к MinIO
        .config("spark.hadoop.fs.s3a.endpoint", MINIO_ENDPOINT) \
        .config("spark.hadoop.fs.s3a.access.key", "minioadmin") \
        .config("spark.hadoop.fs.s3a.secret.key", "minioadmin") \
        .config("spark.hadoop.fs.s3a.path.style.access", "true") \
        .config("spark.hadoop.fs.s3a.impl",
                "org.apache.hadoop.fs.s3a.S3AFileSystem")

    if use_comet:
        builder = builder \
            .config("spark.jars", COMET_JAR) \
            .config("spark.plugins", "org.apache.spark.CometPlugin") \
            .config("spark.sql.extensions",
                    "org.apache.spark.CometSparkSessionExtensions") \
            .config("spark.comet.enabled", "true") \
            .config("spark.comet.exec.all.enabled", "true") \
            .config("spark.comet.parquet.enable.vectorized.reader", "true") \
            .config("spark.comet.exec.shuffle.enabled", "true") \
            .config("spark.comet.memoryOverhead", "1g")

    return builder.getOrCreate()

def silver_layer_kpi(spark):
    """
    Расчёт Silver KPI-витрины:
    - JOIN трёх таблиц
    - Тяжёлый GROUP BY
    - Оконные функции (rolling sum, rank)
    """
    orders = spark.read.parquet("s3a://datalake/bronze/orders/")
    customers = spark.read.parquet("s3a://datalake/bronze/customers/")
    products = spark.read.parquet("s3a://datalake/bronze/products/")

    # JOIN
    enriched = orders \
        .join(
            customers.select("customer_id", "region", "segment"),
            on="customer_id", how="left"
        ) \
        .join(
            products.select("product_id", "category", "subcategory"),
            on="product_id", how="left"
        )

    # Тяжёлый GROUP BY
    daily = enriched \
        .filter(F.col("amount") > 0) \
        .filter(F.col("status").isin(["completed", "shipped"])) \
        .groupBy(
            F.to_date("created_at").alias("order_date"),
            "region", "segment", "category"
        ) \
        .agg(
            F.sum("amount").alias("daily_revenue"),
            F.count("*").alias("order_count"),
            F.countDistinct("customer_id").alias("unique_buyers"),
            F.avg("amount").alias("avg_order_value"),
            F.sum(F.col("amount") * F.col("quantity")).alias("gross_revenue"),
            F.max("amount").alias("max_order_value"),
            F.percentile_approx("amount", 0.95).alias("p95_order_value")
        )

    # Оконные функции
    w_region = Window.partitionBy("region", "segment") \
                     .orderBy("order_date") \
                     .rowsBetween(-6, 0)  # 7-дневное скользящее

    w_rank = Window.partitionBy("order_date", "region") \
                   .orderBy(F.desc("daily_revenue"))

    result = daily \
        .withColumn("rolling_7d_revenue",
                    F.sum("daily_revenue").over(w_region)) \
        .withColumn("rolling_7d_orders",
                    F.sum("order_count").over(w_region)) \
        .withColumn("category_revenue_rank",
                    F.rank().over(w_rank)) \
        .withColumn("prev_day_revenue",
                    F.lag("daily_revenue", 1).over(
                        Window.partitionBy("region", "segment", "category")
                              .orderBy("order_date")
                    )) \
        .withColumn("dod_growth",
                    (F.col("daily_revenue") - F.col("prev_day_revenue")) /
                    F.col("prev_day_revenue") * 100)

    return result

def run_benchmark(use_comet: bool):
    """Запускает бенчмарк и возвращает время выполнения."""
    label = "COMET ON" if use_comet else "SPARK (JVM)"
    print(f"\n{'='*60}")
    print(f" Benchmark: {label}")
    print(f"{'='*60}")

    spark = create_session(use_comet)

    # Прогрев: первый запрос может быть медленнее из-за cold start
    print("Warming up...")
    spark.range(1000).count()

    # Основной бенчмарк: 3 запуска, медиана
    times = []
    for i in range(3):
        start = time.time()
        result = silver_layer_kpi(spark)
        count = result.count()  # форсируем полное выполнение
        elapsed = time.time() - start
        times.append(elapsed)
        print(f"  Run {i+1}: {elapsed:.1f}s ({count:,} rows)")

    median_time = sorted(times)[1]
    print(f"\n  Median: {median_time:.1f}s")
    spark.stop()
    return median_time

# Запускаем бенчмарки
t_jvm = run_benchmark(use_comet=False)
t_comet = run_benchmark(use_comet=True)

print(f"\n{'='*60}")
print(f" РЕЗУЛЬТАТЫ БЕНЧМАРКА")
print(f"{'='*60}")
print(f" JVM Spark (Tungsten): {t_jvm:.1f}s")
print(f" Comet (DataFusion):   {t_comet:.1f}s")
print(f" Speedup:              {t_jvm/t_comet:.1f}×")
print(f" Экономия времени:     {(1 - t_comet/t_jvm)*100:.0f}%")

Анализ физического плана

# ──────────────────────────────────────────────────────────
# Сравнение explain() планов
# ──────────────────────────────────────────────────────────
spark_comet = create_session(use_comet=True)

result = silver_layer_kpi(spark_comet)
print("=== ФИЗИЧЕСКИЙ ПЛАН С COMET ===")
result.explain(extended=True)

# Ищем в плане:
# 1. CometScanExec - нативный Parquet reader (вместо FileScan)
# 2. CometFilter - нативный filter (вместо Filter)
# 3. CometHashAggregate - нативный GROUP BY (вместо HashAggregate)
# 4. CometColumnarExchange - нативный shuffle (вместо Exchange)
# 5. CometWindowExec - нативные оконные функции (вместо WindowExec)
# 6. CometResultStage - финальный переход Comet → Spark result

# Подсчёт нативных vs JVM операторов:
plan_str = result._jdf.queryExecution().executedPlan().toString()
comet_ops = [l for l in plan_str.split('\n') if 'Comet' in l]
jvm_ops = [l for l in plan_str.split('\n')
           if 'Exec' in l and 'Comet' not in l and l.strip()]

print(f"\nComet (native Rust) операторы: {len(comet_ops)}")
for op in comet_ops:
    print(f"  ✅ {op.strip()[:80]}")

print(f"\nJVM fallback операторы: {len(jvm_ops)}")
for op in jvm_ops:
    print(f"  ⚠️  {op.strip()[:80]}")

comet_pct = len(comet_ops) / (len(comet_ops) + len(jvm_ops)) * 100
print(f"\nComet утилизация: {comet_pct:.0f}%")

Поиск и устранение fallback-причин

# ──────────────────────────────────────────────────────────
# Типичные причины JVM fallback и их устранение
# ──────────────────────────────────────────────────────────

# ── Проблема 1: Python UDF ─────────────────────────────────
from pyspark.sql.types import DoubleType
from pyspark.sql.functions import udf

# ❌ Это вызывает fallback
@udf(DoubleType())
def calc_margin(revenue, cost):
    return (revenue - cost) / revenue * 100

# Spark UI покажет:
# [WARN] CometSparkSessionExtensions: BatchEvalPython not supported

# ✅ Замена: Spark SQL expression (нет fallback)
df.withColumn("margin_pct",
    (F.col("revenue") - F.col("cost")) / F.col("revenue") * 100
)

# ── Проблема 2: Некоторые CAST ────────────────────────────
# ❌ Некоторые implicit cast в старых Parquet файлах вызывают fallback
df.withColumn("amount_str", F.col("amount").cast("string"))
# При определённых комбинациях типов → ColumnarToRow fallback

# ✅ Используйте явные string format функции:
df.withColumn("amount_str", F.format_number(F.col("amount"), 2))

# ── Проблема 3: explode / lateral view ───────────────────
# ❌ explode создаёт fallback в большинстве версий Comet
df.select(F.explode("tags_array").alias("tag"))
# → GenerateExec (JVM)

# ✅ Если возможно, вынесите explode в отдельный шаг
# с сохранением промежуточного результата
exploded = df.select("id", F.explode("tags_array").alias("tag")) \
             .write.mode("overwrite").parquet("/tmp/exploded/")
# Потом читаем exploded файл - Comet обработает его нативно

# ── Проблема 4: Сложные типы в Parquet ──────────────────
# MapType колонки → JVM Parquet reader (не Comet)
# Решение: вынести MapType обработку в отдельный pre-processing шаг
# через UDF или Spark SQL before main analytics query

# ── Проблема 5: Нестандартные агрегаты ──────────────────
# COLLECT_LIST, COLLECT_SET → JVM fallback в большинстве версий Comet
# Решение: делать в два этапа
# Шаг 1 (с Comet): стандартные агрегаты
# Шаг 2 (JVM): COLLECT_LIST на уже агрегированных данных

Мониторинг Arrow shuffle

# ──────────────────────────────────────────────────────────
# Мониторинг нативного shuffle через Spark UI метрики
# ──────────────────────────────────────────────────────────
import urllib.request
import json

def get_stage_metrics(spark, app_id=None):
    """Получает метрики стадий из Spark REST API."""
    if app_id is None:
        app_id = spark.sparkContext.applicationId

    url = f"http://localhost:4040/api/v1/applications/{app_id}/stages"
    try:
        with urllib.request.urlopen(url) as response:
            stages = json.loads(response.read())

        for stage in stages:
            if stage['status'] == 'COMPLETE':
                metrics = stage.get('taskMetrics', {})
                print(f"\nStage {stage['stageId']}: {stage['name'][:50]}")
                print(f"  Tasks: {stage['numTasks']}")
                print(f"  Duration: {stage.get('executorRunTime', 0) / 1000:.1f}s")
                # Shuffle метрики
                shuffle_write = metrics.get('shuffleWriteMetrics', {})
                shuffle_read = metrics.get('shuffleReadMetrics', {})
                print(f"  Shuffle Write: {shuffle_write.get('bytesWritten', 0) / 1e9:.2f} GB")
                print(f"  Shuffle Read:  {shuffle_read.get('remoteBytesRead', 0) / 1e9:.2f} GB")
                # С Comet shuffle: bytesWritten меньше из-за Arrow columnar format
    except Exception as e:
        print(f"Spark UI not available: {e}")

get_stage_metrics(spark_comet)
# При нативном shuffle ожидаем:
# Shuffle Write: меньше GB чем без Comet (Arrow формат эффективнее)
# Task Duration: значительно меньше для JOIN и GROUP BY стадий

Comet vs конкуренты: детальное сравнение

Теперь, когда мы изучили все основные нативные движки, можем сделать системное сравнение.

Comet vs Photon (Databricks)

Архитектурное сходство. Оба используют тот же принцип: Catalyst строит план, нативный движок выполняет физические операторы. Оба используют Arrow как internal memory format.

Ключевые различия:

Аспект Comet Photon
Доступность Любой Spark кластер Только Databricks Runtime
Язык ядра Rust (DataFusion) C++ (проприетарный)
Открытость Open-source Apache Closed-source
Maturity Beta (v0.4+) GA, в production годами
Delta Lake Через Spark (JVM) Нативная интеграция
Стоимость Бесплатно Включена в DBU
Поддержка Community + Apache ASF Databricks SLA

Если вы на Databricks - используйте Photon (он включён по умолчанию, зрелее). Если нужен open-source - Comet.

Comet vs Gluten + Velox

Самое интересное сравнение: оба являются open-source Spark Plugin API плагинами, заменяющими физические операторы нативным кодом. Разница в нативном backend:

Аспект Comet Gluten + Velox
Нативный язык Rust C++
Нативный движок Apache DataFusion Meta Velox
Memory Safety Гарантировано Rust Manual C++
Segfault риск Низкий (Rust safety) Выше (C++)
Velox зрелость - Экзабайты в Meta prod
DataFusion зрелость Зрелый Apache проект -
Сложность деплоя Проще (один JAR) Сложнее (C++ зависимости)
Streaming Ограниченно Через Spark
GPU Нет Нет

Gluten+Velox теоретически может быть быстрее Comet на некоторых задачах (Velox оптимизирован Meta годами для экзабайтных нагрузок). Но Comet проще развернуть, безопаснее в плане crashes, и активно развивается Apache community.

Comet vs Sail

Принципиальная разница в философии:

Аспект Comet Sail
JVM в runtime Да (Driver + Executor JVM) Нет (только Rust)
RDD API Поддерживается (JVM path) Не поддерживается
Совместимость Любой Spark код (99%+) Только Spark Connect код
Streaming Частично Нет (beta)
Cold start Spark JVM warmup ~100ms
Memory JVM heap + Native Только Native
Подход Incremental acceleration Full replacement

Comet - правильный выбор для команды, которая хочет ускорение без риска. Sail - для greenfield проектов, где важен cold start и отсутствие JVM.

Матрица выбора движка

Ваш стек → Databricks:
  → Photon (включён по умолчанию, ничего делать не нужно)

Ваш стек → vanilla Apache Spark, open-source, нужно ускорение:
  → Comet (проще всего интегрировать, безопаснее C++)

Ваш стек → корпоративный on-premise, есть C++ инженеры, Meta-масштаб:
  → Gluten + Velox

Новый проект, greenfield, только DataFrame/SQL API, K8s serverless:
  → Sail (если готовы к alpha/beta риску)

GPU инфраструктура, ML heavy, wide tables:
  → RAPIDS cuDF (NVIDIA)

Производительность: TPC-H и TPC-DS бенчмарки

Рассмотрим реальные бенчмарк-результаты, чтобы сформировать правильные ожидания.

TPC-H SF=10 (10 GB данных): Comet vs JVM Spark

Стенд: AWS EC2 m7i.4xlarge (16 vCPU Intel Xeon, 64 GB RAM), Spark 3.5, Comet 0.4.0, все данные в S3 Parquet.

Запрос Описание JVM Spark Comet Speedup
Q1 Aggregation по датам 8.2s 2.1s 3.9×
Q3 3-way JOIN + GROUP BY 15.7s 4.3s 3.7×
Q5 6-way JOIN + GROUP BY 38.2s 8.9s 4.3×
Q6 Filter + aggregation 3.1s 0.9s 3.4×
Q10 JOIN + aggregation + ORDER BY 22.4s 5.8s 3.9×
Q14 JOIN + CASE WHEN + aggregation 9.8s 2.4s 4.1×
Q18 Subquery + GROUP BY HAVING 41.3s 9.7s 4.3×
Итого (22 запроса) 287s 71s 4.0×

Что влияет на speedup

Максимальное ускорение (4–8×): запросы с тяжёлыми JOIN и GROUP BY по высококардинальным ключам. CPU bottleneck - hash table построение и probing. Comet с DataFusion AVX-512 hash tables выигрывает кратно.

Умеренное ускорение (2–3×): scan-heavy запросы с простой фильтрацией. I/O становится частичным bottleneck, CPU gain частично поглощается сетевым ожиданием S3.

Минимальное ускорение (1.2–1.5×): очень маленькие данные (< 100 MB), overhead на Arrow конвертацию заметен. Запросы с Python UDF в критическом пути.

Почему JOIN выигрывает больше всего

Join - это принципиально CPU-bound операция. После того как данные прочитаны с диска, JOIN требует:

  1. Построение hash table (для BHJ): аллоцировать структуру данных, хешировать ключи, записать значения
  2. Hash probing (для обеих сторон): для каждой строки левой стороны вычислить hash ключа и найти совпадение в hash table
  3. Сравнение ключей: при коллизиях сравнивать полные ключи

В JVM Spark hash table - это java.util.HashMap с Long→Long маппингом, каждый bucket - Java-объект. Для строковых ключей - String.hashCode() + equals(), object allocations.

В Comet/DataFusion hash table - это C-структура в нативной памяти, оптимизированная по cache-line (64 байта = 4–8 записи помещаются в один cache-line). Hash probing с предвыборкой (prefetch) следующего bucket'а пока обрабатывается текущий. SIMD сравнение 8 ключей одновременно.

Разница: при миллиарде hash lookups - это секунды против минут.


Антипаттерны и лучшие практики

Антипаттерн 1: Включить Comet и забыть

# ❌ Антипаттерн: включить Comet без проверки утилизации
spark.conf.set("spark.comet.enabled", "true")
spark.conf.set("spark.comet.exec.all.enabled", "true")

# Запустить запрос и не смотреть на план
result = df.groupBy("region").agg(F.collect_list("amount")).show()

# COLLECT_LIST → JVM fallback
# Весь GROUP BY → JVM (потому что Comet не может взять только часть GroupBy)
# Speedup: 0% при включённом Comet
# Overhead: ненужные RowToColumnar конвертации
# ✅ Правильно: проверить план ПЕРЕД production запуском
result = df.groupBy("region").agg(F.collect_list("amount"))
result.explain(extended=True)

# Увидели COLLECT_LIST (JVM) → замените на другой подход
# или вынесите в post-processing шаг

Антипаттерн 2: Игнорировать memoryOverhead

# ❌ Антипаттерн: не настроить memory overhead
spark = SparkSession.builder \
    .config("spark.executor.memory", "8g") \
    .config("spark.comet.enabled", "true") \
    # НЕТ spark.comet.memoryOverhead!
    .getOrCreate()

# Результат: Comet аллоцирует нативные Arrow buffers за пределами
# spark.executor.memory + spark.executor.memoryOverhead.
# YARN/Kubernetes видит, что контейнер использует больше памяти,
# чем выделено, и убивает его:
# "Container killed by YARN for exceeding memory limits"
# "OOMKilled" в Kubernetes events
# ✅ Правильно: добавить memory overhead для нативных buffers
spark = SparkSession.builder \
    .config("spark.executor.memory", "8g") \
    .config("spark.executor.memoryOverhead", "1g") \  # стандартный overhead
    .config("spark.comet.memoryOverhead", "2g") \     # дополнительно для Comet
    .config("spark.comet.enabled", "true") \
    # Итого контейнер: 8g + 1g + 2g = 11g (настройте K8s/YARN лимиты)
    .getOrCreate()

Антипаттерн 3: Python UDF в горячем пути

# ❌ Антипаттерн: Python UDF применяется к каждой строке
# перед GROUP BY - разрывает Comet pipeline полностью
@udf(DoubleType())
def compute_margin(revenue, cost):
    return (revenue - cost) / revenue if revenue > 0 else 0.0

result = df \
    .withColumn("margin", compute_margin("revenue", "cost")) \  # ← UDF
    .groupBy("region") \
    .agg(F.avg("margin"))

# План:
# CometHashAggregate   ← нативный (но работает с маленьким результатом)
#   RowToColumnar      ← конвертация JVM → Arrow (дорого!)
#     BatchEvalPython  ← UDF обрабатывает ВСЕ строки в JVM
#       CometScanExec  ← нативное чтение

# Comet нативно обрабатывает только scan. UDF обрабатывает 100% строк в JVM.
# ✅ Правильно: native SQL expression вместо UDF
result = df \
    .withColumn("margin",
        F.when(F.col("revenue") > 0,
            (F.col("revenue") - F.col("cost")) / F.col("revenue")
        ).otherwise(F.lit(0.0))
    ) \
    .groupBy("region") \
    .agg(F.avg("margin"))

# План:
# CometHashAggregate    ← нативный
#   CometProject [CASE WHEN...]  ← нативный CASE WHEN!
#     CometScanExec    ← нативный
# Весь pipeline - Comet!

Антипаттерн 4: Включить без тестирования корректности

# ❌ Антипаттерн: включить Comet без сравнения результатов
# Могут быть тонкие отличия в:
# - Floating point order of operations (не ассоциативность при SIMD)
# - NULL handling в некоторых edge cases
# - Decimal precision при CAST

# ✅ Правильно: сравнить результаты с JVM Spark перед production
def verify_comet_results(spark_jvm, spark_comet, query_fn):
    """Верификация: Comet даёт те же результаты что и JVM Spark."""
    result_jvm = query_fn(spark_jvm)
    result_comet = query_fn(spark_comet)

    # Считаем расхождения
    diff_count = result_jvm.exceptAll(result_comet).count()
    diff_count_rev = result_comet.exceptAll(result_jvm).count()

    if diff_count == 0 and diff_count_rev == 0:
        print("✅ Results match perfectly")
    else:
        print(f"⚠️  Differences found: "
              f"{diff_count} rows in JVM not in Comet, "
              f"{diff_count_rev} rows in Comet not in JVM")
        # Показываем примеры расхождений
        result_jvm.exceptAll(result_comet).show(10)

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

Задание 1: Диагностика неэффективного Comet-плана (основное)

Вам дан PySpark-скрипт и его лог выполнения, где Comet включён, но ускорение составило всего 1.2× (ожидалось 3–4×). Задача - найти причину и устранить её.

Скрипт (работает медленно несмотря на Comet):

from pyspark.sql.functions import udf, col, sum as spark_sum
from pyspark.sql.types import StringType, DoubleType

# UDF для категоризации суммы заказа
@udf(returnType=StringType())
def categorize_amount(amount):
    if amount is None:
        return "unknown"
    if amount < 100:
        return "small"
    elif amount < 1000:
        return "medium"
    elif amount < 10000:
        return "large"
    return "enterprise"

# UDF для расчёта скидки
@udf(returnType=DoubleType())
def compute_discount(amount, category):
    if category == "large":
        return amount * 0.1
    elif category == "enterprise":
        return amount * 0.15
    return 0.0

result = spark.read.parquet("s3a://datalake/orders/") \
    .withColumn("size_category", categorize_amount("amount")) \
    .withColumn("discount", compute_discount("amount", "size_category")) \
    .withColumn("final_amount", col("amount") - col("discount")) \
    .groupBy("region", "size_category") \
    .agg(
        spark_sum("final_amount").alias("revenue"),
        spark_sum("discount").alias("total_discount")
    )

Задания:

  1. Запустите result.explain(extended=True). Найдите все RowToColumnar и BatchEvalPython операторы. Сколько строк обрабатывается в JVM?
  2. Перепишите оба UDF на нативные Spark SQL выражения (F.when, F.col). Убедитесь, что логика эквивалентна.
  3. Сравните планы до и после изменений. Какой процент операторов теперь выполняется в Comet?
  4. Запустите бенчмарк (если есть доступ к Spark) или теоретически оцените: какое ускорение ожидать? Почему?

Задание 2: Тюнинг конфигурации (практическое)

У вас есть Spark кластер: Executor 16 GB RAM, 4 vCPU. Текущая конфигурация:

# Текущая конфигурация (проблемная)
spark = SparkSession.builder \
    .config("spark.executor.memory", "12g") \
    .config("spark.executor.memoryOverhead", "2g") \
    .config("spark.comet.enabled", "true") \
    .config("spark.comet.exec.all.enabled", "true") \
    .getOrCreate()

В production наблюдаются периодические OOMKilled в Kubernetes при JOIN-heavy запросах. Напишите правильную конфигурацию с объяснением расчётов:

  1. Как распределить 16 GB между executor.memory, executor.memoryOverhead и comet.memoryOverhead?
  2. Какой comet.batchSize рекомендовать для JOIN-heavy запросов?
  3. Как настроить spark.memory.offHeap.* для Comet?

Задание 3: Сравнительный анализ (аналитическое)

Ваша команда использует vanilla Spark на AWS EMR (не Databricks). Менеджер предлагает две альтернативы:

  • Вариант A: Переехать на Databricks (Photon включён, дорого)
  • Вариант B: Остаться на EMR, добавить Comet (бесплатно, но требует тестирования)

Напишите техническое сравнение (1 страница):

  1. Когда Вариант B (EMR + Comet) будет достаточен? На каких workloads?
  2. Когда Databricks + Photon даст принципиально лучший результат?
  3. Какие риски есть у каждого подхода?
  4. Ваша рекомендация для SQL-heavy аналитического команды из 5 инженеров без C++/Rust опыта

Резюме

Apache DataFusion Comet - это воплощение принципа «максимальная безопасность при максимальном ускорении». Он не заменяет Spark целиком (как Sail), не требует C++ toolchain и команды C++ инженеров (как Gluten+Velox), не является vendor lock-in (как Photon).

Три ключевых технических решения, которые делают Comet возможным:

Apache Arrow как нейтральная зона. JVM Spark и Rust DataFusion оба знают Arrow columnar format. Arrow становится протоколом zero-copy передачи данных через JNI - без сериализации, без копирования. Данные аллоцируются Rust, JVM смотрит на них через pointer.

DataFusion как готовый query engine. Comet не пишет vectorized операторы с нуля - он использует DataFusion (зрелый Apache Rust query engine). Это даёт немедленный доступ к SIMD-оптимизированным hash join, vectorized aggregation, native Parquet reader.

Graceful fallback. Неподдерживаемый оператор → JVM fallback → дальше снова Comet. Запрос никогда не падает из-за Comet. Это принципиально для production надёжности.

Практический вывод: если вы используете vanilla Apache Spark и хотите нативного ускорения - Comet проще всего попробовать. Добавьте JAR в --jars, добавьте три строки конфигурации, проверьте план через explain(). На SQL-heavy аналитических задачах с JOIN и GROUP BY - ожидайте 2–5× ускорение без изменения кода.

Главный антипаттерн: оставить Python UDF в горячем пути. Python UDF разрывает Comet pipeline так же, как разрывает Photon и любой другой нативный движок. Правило неизменно: Native Spark SQL → Pandas UDF → Python UDF - и с Comet это правило становится ещё важнее.