Apache DataFusion Comet: архитектура Rust + Arrow, как подключается к Spark
Apache Comet как open-source Rust-ускоритель Spark: DataFusion внутри, Arrow zero-copy, Plugin API интеграция, JVM fallback и TPC-H бенчмарки.
Концепция и предпосылки: зачем ещё один нативный движок¶
К этому моменту курса мы изучили три подхода к нативному ускорению 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. Рассмотрим, что именно происходит:
- Spark Executor получает Task с физическим планом
CometExecоператор (JVM Scala-класс) вызывается Spark'ом- Внутри
CometExec.executeColumnar()происходит JNI-вызов в Rust - Rust-код создаёт DataFusion physical plan из protobuf
- DataFusion читает Parquet напрямую в Arrow buffers - минуя JVM heap полностью
- DataFusion выполняет операторы (Filter, Aggregate, Join) над Arrow buffers
- Результат - Arrow RecordBatch в нативной памяти
- JNI возвращает в JVM
long-указатель на этот RecordBatch - 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:
CometSparkSessionExtensionsобходит физический план Catalyst- Поддерживаемые узлы сериализуются в protobuf схему (
.protoфайлы описывают операторы:CometFilter,CometAggregate,CometJoin, ...) - Protobuf байты передаются через JNI в Rust как
byte[]массив - Rust десериализует protobuf в DataFusion physical plan объекты (используя crate
prost) - 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 требует:
- Построение hash table (для BHJ): аллоцировать структуру данных, хешировать ключи, записать значения
- Hash probing (для обеих сторон): для каждой строки левой стороны вычислить hash ключа и найти совпадение в hash table
- Сравнение ключей: при коллизиях сравнивать полные ключи
В 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")
)
Задания:
- Запустите
result.explain(extended=True). Найдите всеRowToColumnarиBatchEvalPythonоператоры. Сколько строк обрабатывается в JVM? - Перепишите оба UDF на нативные Spark SQL выражения (
F.when,F.col). Убедитесь, что логика эквивалентна. - Сравните планы до и после изменений. Какой процент операторов теперь выполняется в Comet?
- Запустите бенчмарк (если есть доступ к 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 запросах. Напишите правильную конфигурацию с объяснением расчётов:
- Как распределить 16 GB между
executor.memory,executor.memoryOverheadиcomet.memoryOverhead? - Какой
comet.batchSizeрекомендовать для JOIN-heavy запросов? - Как настроить
spark.memory.offHeap.*для Comet?
Задание 3: Сравнительный анализ (аналитическое)¶
Ваша команда использует vanilla Spark на AWS EMR (не Databricks). Менеджер предлагает две альтернативы:
- Вариант A: Переехать на Databricks (Photon включён, дорого)
- Вариант B: Остаться на EMR, добавить Comet (бесплатно, но требует тестирования)
Напишите техническое сравнение (1 страница):
- Когда Вариант B (EMR + Comet) будет достаточен? На каких workloads?
- Когда Databricks + Photon даст принципиально лучший результат?
- Какие риски есть у каждого подхода?
- Ваша рекомендация для 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 это правило становится ещё важнее.