Comet: установка, операторы, ограничения и TPC-H бенчмарк

Установка jar и конфигурация Comet, поддерживаемые операторы, что падает на JVM, запуск TPC-H бенчмарка.

optimization

Практический старт: Установка JAR и конфигурация Comet

Прежде чем переходить к изучению операторов и бенчмаркам, необходимо понять, каким образом Apache DataFusion Comet вообще подключается к Spark. Это знание не сводится к копированию трёх строчек конфигурации - здесь важно понять, почему каждый параметр существует и что произойдёт, если его пропустить.

Как Comet попадает в Spark

Comet - это внешний плагин для Spark. Он не встроен в дистрибутив Apache Spark и не является его частью. Механизм подключения реализован через Spark Plugin API (появился в Spark 3.0): специальный интерфейс, который позволяет сторонним библиотекам встраиваться в жизненный цикл Driver и Executor без модификации исходного кода самого Spark.

Физически Comet поставляется в виде FAT JAR - единого архива, внутри которого упакованы:

  • JVM-классы (Scala/Java): реализации физических операторов (CometScanExec, CometFilterExec, CometHashAggregateExec и другие), классы Plugin API, утилиты конфигурации.
  • Нативная shared library (.so на Linux, .dylib на macOS): скомпилированный Rust-код самого DataFusion - векторизованные операторы, аллокатор памяти на базе Arrow, нативный Parquet-ридер.

При старте Executor JVM автоматически распаковывает нативную библиотеку во временную директорию и загружает её через System.loadLibrary(). С этого момента JVM и Rust-рантайм существуют в одном OS-процессе и общаются через JNI.

Выбор артефакта: Maven или сборка из исходников

Comet публикуется в Maven Central. Координаты артефакта зависят от трёх переменных:

  • Версия Spark (3.4, 3.5, 4.0)
  • Версия Scala (2.12, 2.13)
  • Версия самого Comet

Это принципиально важно: Comet компилирует JVM-классы под конкретную версию Spark API, а несовпадение версий приводит к ошибкам NoSuchMethodError или ClassNotFoundException при старте.

# Проверяем версию текущего Spark
spark-submit --version

# Подключение через --packages (для Spark 3.5 + Scala 2.12, Comet 0.15)
spark-submit \
  --packages org.apache.datafusion:datafusion-comet-spark-3.5_2.12:0.15.0 \
  your_job.py

# Альтернатива: указать путь к скачанному JAR
spark-submit \
  --jars /opt/jars/datafusion-comet-spark-3.5_2.12-0.15.0.jar \
  your_job.py

Если в вашей организации используется внутренний Maven-репозиторий (Nexus, Artifactory), рекомендуется скачать JAR туда и раздавать через --packages с настроенным зеркалом, а не полагаться на прямой доступ в Maven Central из Executor-нод.

Регистрация плагина: два обязательных параметра

Для активации Comet необходимо передать два конфигурационных ключа - и они не взаимозаменяемы:

from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .appName("comet-demo")
    # 1. Регистрация нативного плагина через Plugin API
    .config("spark.plugins", "org.apache.spark.CometPlugin")
    # 2. Регистрация расширений Catalyst для перехвата планов
    .config("spark.sql.extensions", "org.apache.spark.sql.comet.CometSparkSessionExtensions")
    # 3. Активация нативного выполнения
    .config("spark.comet.enabled", "true")
    .config("spark.comet.exec.all.enabled", "true")
    .getOrCreate()
)

Почему нужны оба?

spark.plugins регистрирует CometPlugin - класс, который инициализируется при запуске Driver и Executor. Он отвечает за загрузку нативной библиотеки, инициализацию Arrow-аллокатора и передачу конфигурации в Rust-рантайм. Без этого параметра нативный код вообще не запустится.

spark.sql.extensions регистрирует CometSparkSessionExtensions - набор правил, которые встраиваются в Catalyst optimizer. Именно они преобразуют стандартные физические операторы Spark (HashAggregateExec, FilterExec, SortMergeJoinExec) в их нативные аналоги (CometHashAggregateExec, CometFilterExec и т.д.). Без этого параметра Comet загрузится, но не будет трансформировать план - Spark продолжит работать полностью на JVM.

Иными словами: spark.plugins запускает нативный движок, а spark.sql.extensions говорит Spark, какие именно части плана передать этому движку.

Ключевые параметры spark.comet.*

После включения базовой активации следует настроить дополнительные параметры:

spark = (
    SparkSession.builder
    # ... базовые настройки выше ...

    # Нативный Shuffle-менеджер (Arrow IPC вместо Java serialization)
    .config("spark.shuffle.manager", "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
    .config("spark.comet.exec.shuffle.enabled", "true")

    # Управление памятью: КРИТИЧЕСКИ важно для предотвращения OOM
    .config("spark.memory.offHeap.enabled", "true")
    .config("spark.memory.offHeap.size", "4g")
    .config("spark.executor.memoryOverhead", "2g")

    # Размер Arrow-батча (по умолчанию 8192 строк)
    .config("spark.comet.batchSize", "8192")

    # Логирование причин деградации (для траблшутинга)
    .config("spark.comet.logFallbackReasons.enabled", "true")
    .config("spark.comet.explainFallback.enabled", "true")

    .getOrCreate()
)

Управление памятью: самая частая причина OOM

Это наиболее критичная часть настройки, которую часто игнорируют - и получают OutOfMemoryError в продакшене.

Rust-рантайм DataFusion аллоцирует память вне JVM heap: он использует нативную память OS через Arrow-аллокатор. JVM не знает об этой памяти и не учитывает её при сборке мусора. Но YARN и Kubernetes считают всю память процесса - включая нативную.

Если не выделить дополнительную память, произойдёт следующее:

  1. Executor JVM запросил 4 GB (spark.executor.memory)
  2. Comet выделил ещё 2 GB нативной памяти для Arrow-буферов
  3. YARN видит процесс, использующий 6 GB, при лимите 4.5 GB
  4. YARN убивает Executor как превысивший лимит
  5. Spark регистрирует это как потерю Executor и пытается перезапустить задачу

Правильная схема расчёта памяти для Executor:

spark.executor.memory        = JVM heap для Spark
spark.memory.offHeap.size    = нативная память для Arrow-буферов Comet
spark.executor.memoryOverhead = буфер для JVM накладных расходов + нативные библиотеки

Итого физической памяти на контейнер:
  executor.memory + executor.memoryOverhead + offHeap.size

Пример для Executor с 16 GB RAM:

.config("spark.executor.memory", "10g")
.config("spark.memory.offHeap.enabled", "true")
.config("spark.memory.offHeap.size", "4g")
.config("spark.executor.memoryOverhead", "2g")
# Итого: 10 + 4 + 2 = 16 GB - точно вписываемся в контейнер

Проверка успешной активации

После запуска Spark-сессии с Comet проверьте логи:

# В логах Driver при старте должны появиться строки:
INFO CometPlugin: Initializing Comet plugin
INFO CometPlugin: Comet native library loaded successfully
INFO CometPlugin: DataFusion version: X.Y.Z

# В логах Executor:
INFO CometPlugin: Executor plugin initialized
INFO NativeBase: Native library loaded from: /tmp/comet-native-XXXXX.so

Если нативная библиотека не загрузилась, вы увидите предупреждение и Spark продолжит работу в режиме полного JVM-fallback (Comet не будет использоваться вообще, но и не упадёт с ошибкой).


Карта покрытия: поддерживаемые физические операторы

Понимание того, какие операторы Comet умеет выполнять нативно, а какие нет - ключ к правильной оценке применимости инструмента и интерпретации планов выполнения.

Принцип трансформации операторов

Когда Catalyst строит физический план, Comet через CometSparkSessionExtensions применяет специальные правила замены (ColumnarRule). Каждый поддерживаемый оператор JVM заменяется на его нативный аналог с префиксом Comet. Важно понимать, что это замена физических операторов - логический план Catalyst и AQE при этом остаются нетронутыми.

Решение о том, может ли оператор быть выполнен нативно, принимается на основе трёх факторов: тип оператора, типы данных всех вовлечённых колонок, набор выражений (expressions) внутри оператора.

CometScan / CometBatchScan - чтение данных

Это точка входа данных в нативный рантайм. Когда Comet берёт на себя чтение Parquet-файлов, он использует собственный Rust-реализованный Parquet-ридер (из библиотеки parquet-rs), который:

  • Читает колоночные данные Parquet непосредственно в Arrow RecordBatch-структуры в нативной памяти, минуя JVM heap
  • Применяет column pruning (выбор только нужных колонок) и predicate pushdown (вынос фильтров на уровень чтения файла) уже на уровне Rust
  • Декодирует сжатые данные (Snappy, Zstd, LZ4) нативными декомпрессорами без Java-обёрток

Разница с JVM Parquet-ридером: стандартный VectorizedParquetRecordReader читает данные в JVM ColumnarBatch (Java-объект на куче), тогда как CometBatchScan читает напрямую в Arrow-буферы в native memory. Это устраняет один лишний цикл копирования данных.

-- SQL, который активирует CometBatchScan:
SELECT order_id, amount FROM orders
WHERE order_date >= '2024-01-01'

В плане выполнения вы увидите:

CometBatchScan[order_id, amount]
  PushedFilters: [GreaterThanOrEqual(order_date, 2024-01-01)]
  ReadSchema: struct<order_id:bigint, amount:decimal(18,2)>

CometFilterExec - фильтрация данных

Comet выполняет предикаты фильтрации на уровне SIMD-инструкций процессора. Вместо того чтобы проверять каждую строку в JVM-цикле, DataFusion применяет векторные инструкции AVX2/AVX-512, которые обрабатывают 8 или 16 значений типа double за одну инструкцию CPU.

Пример: предикат amount > 1000.0 применяется к батчу из 8192 строк. JVM-реализация выполнит 8192 сравнения последовательно. Нативная SIMD-реализация выполнит 8192 / 16 = 512 векторных операций - теоретически в 16 раз быстрее (на практике 3–8×, с учётом накладных расходов).

Поддерживаемые типы предикатов в CometFilterExec:

  • Сравнения (>, <, >=, <=, =, !=)
  • Логические комбинации (AND, OR, NOT)
  • Проверки на NULL (IS NULL, IS NOT NULL)
  • Диапазоны (BETWEEN)
  • Списки значений (IN (...) с фиксированными значениями)
  • Регулярные выражения (LIKE, RLIKE) - при совместимости паттерна

CometProjectExec - проекция колонок

Проекция - это операция выбора и преобразования колонок: SELECT col1, col1 + col2 AS sum_val, ROUND(price, 2) AS rounded_price. В нативном исполнении каждое выражение вычисляется над целым Arrow RecordBatch (8192 элементов) за одну операцию, а не построчно.

Поддерживаемые арифметические операции: +, -, *, /, %, ABS, ROUND, FLOOR, CEIL. Поддерживаемые строковые функции: UPPER, LOWER, TRIM, SUBSTR, CONCAT, LENGTH, REPLACE. Большинство стандартных Spark SQL функций поддерживается - полный список обновляется с каждой версией Comet.

CometHashAggregateExec - агрегация

Агрегация - один из наиболее выигрывающих от нативного выполнения операторов. Причина проста: GROUP BY с SUM, COUNT, AVG, MIN, MAX - это классически CPU-bound операция, которая идеально ложится на векторизованное выполнение.

Как работает нативная хэш-агрегация DataFusion:

  1. Первый проход (partial aggregate): каждый Executor обрабатывает свою партицию данных. Для каждого уникального ключа группировки строится запись в нативной хэш-таблице (Rust HashMap с открытой адресацией для кэш-дружественного доступа).
  2. Shuffle: промежуточные результаты перемешиваются по ключу группировки между Executor-ами.
  3. Второй проход (final aggregate): каждый Executor объединяет промежуточные результаты для своих ключей.

Всё это происходит в нативной памяти, без аллокации Java-объектов для каждой строки. Результат: для запросов вида GROUP BY ... HAVING ... ORDER BY Comet даёт наибольшее ускорение - типично 3–6× по сравнению с JVM HashAggregateExec.

-- Этот запрос целиком выполняется нативно в Comet:
SELECT
    region,
    product_category,
    COUNT(*) AS cnt,
    SUM(revenue) AS total_revenue,
    AVG(discount_pct) AS avg_discount
FROM sales
WHERE sale_date >= '2024-01-01'
GROUP BY region, product_category
HAVING COUNT(*) > 100
ORDER BY total_revenue DESC

CometBroadcastHashJoinExec и CometSortMergeJoinExec - объединение таблиц

Join - самая сложная операция с точки зрения нативной реализации, так как она требует взаимодействия данных с разных Executor-ов.

CometBroadcastHashJoinExec: Меньшая таблица целиком помещается в нативную хэш-таблицу в памяти каждого Executor. При обработке большой таблицы каждая строка сравнивается с этой нативной хэш-таблицей - без JVM-объектов, с SIMD-ускоренным хэшированием.

CometSortMergeJoinExec: Обе таблицы предварительно отсортированы по ключам join. Затем выполняется линейное слияние двух отсортированных потоков Arrow RecordBatch. Эффективен для больших таблиц, когда broadcast нецелесообразен.

Поддерживаемые join-типы: INNER JOIN, LEFT JOIN (LEFT OUTER), RIGHT JOIN, FULL OUTER JOIN (с ограничениями), LEFT SEMI JOIN, LEFT ANTI JOIN.

CometColumnarExchange - нативный Shuffle

Shuffle - это передача данных между Executor-ами по сети и их запись/чтение с диска. В стандартном Spark shuffle-данные сериализуются в Java-формат. Comet заменяет это на Arrow IPC (Inter-Process Communication format) - бинарный формат сериализации Arrow RecordBatch.

Преимущества Arrow IPC shuffle:

  • Нулевое копирование при чтении: если данные уже в Arrow-формате, они передаются напрямую без де/сериализации
  • Более эффективное сжатие: Arrow IPC поддерживает LZ4 и Zstd на уровне батчей
  • Меньше CPU: нет преобразования JVM-объектов → байты и обратно

Границы применимости: что падает обратно на JVM

Понимание fallback-механизма критично для production-использования. Comet не требует 100% покрытия всех операторов - он разработан с принципом "partial acceleration is better than no acceleration". Части плана, которые Comet не поддерживает, автоматически выполняются в обычном JVM-режиме.

Механика деградации

Когда ColumnarRule встречает оператор, который не может быть выполнен нативно, он вставляет операторы переключения контекста:

  • CometRowToColumnarExec: принимает строковые данные из JVM-мира (объекты InternalRow) и конвертирует их в Arrow RecordBatch для передачи в нативный рантайм. Эта конвертация не бесплатна - она требует копирования данных.
  • ColumnarToRow: принимает Arrow RecordBatch из нативного рантайма и конвертирует их обратно в JVM InternalRow. Также требует копирования.

Каждая такая "граница" между JVM и нативным миром - это overhead. Если план содержит много таких переходов, суммарный overhead может нивелировать выигрыш от нативного выполнения отдельных операторов.

Неподдерживаемые типы данных

Comet работает с типами данных, которые имеют прямое представление в Apache Arrow:

Тип Spark Поддержка Comet Причина
IntegerType Полная Прямой Arrow int32
LongType Полная Arrow int64
DoubleType Полная Arrow float64
StringType Полная Arrow utf8
BinaryType Полная Arrow binary
DecimalType(p, s) Полная (p ≤ 38) Arrow decimal128
TimestampType Полная Arrow timestamp[us]
DateType Полная Arrow date32
ArrayType Частичная Arrow list - поддержка растёт
MapType Ограниченная Arrow map - не все операторы
StructType (вложенный) Частичная Зависит от глубины вложенности
UserDefinedType Нет Нет Arrow-представления
CalendarIntervalType Нет Сложная семантика

Когда Comet встречает колонку неподдерживаемого типа, весь оператор, работающий с этой колонкой, переходит в JVM-режим. Даже если остальные 100 колонок поддерживаются - одна "плохая" тянет весь оператор обратно.

Несовместимые выражения и функции

Не все SQL-функции реализованы в DataFusion с точностью, совместимой со Spark. Некоторые функции зависят от JVM-специфичной семантики (локали, часовые пояса, форматы дат):

-- Эти функции ВЫЗОВУТ fallback:
SELECT
    to_date(date_str, 'dd-MMM-yyyy'),     -- locale-зависимый формат
    format_number(amount, 2),              -- locale-зависимое форматирование
    translate(name, 'абвгд', 'АБВГД'),    -- Unicode collation
    regexp_replace(text, pattern, repl)    -- сложные Spark-специфичные regex
FROM table

-- Безопасные аналоги:
SELECT
    to_date(date_str, 'yyyy-MM-dd'),      -- ISO-формат: поддерживается
    ROUND(amount, 2),                      -- поддерживается
    UPPER(name),                           -- поддерживается
    regexp_like(text, '[0-9]+')            -- простые паттерны: поддерживается
FROM table

Python UDF и Pandas UDF

Python UDF - это абсолютный блокировщик нативного выполнения. Когда Spark встречает Python UDF в плане, он обязан вернуться к строковому (row-by-row) выполнению, отправить данные в Python-процесс через сокет/Arrow IPC и получить результат обратно. Comet в этот момент не может ничего сделать.

from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType

# ПЛОХО: этот UDF разрывает нативное выполнение
@udf(returnType=DoubleType())
def calculate_tax(amount, rate):
    return amount * rate / 100

df = spark.table("sales")
result = df.withColumn("tax", calculate_tax(df.amount, df.tax_rate))
# В плане: ... -> ColumnarToRow -> PythonEvalUDF -> CometRowToColumnarExec -> ...
# Весь батч данных "выходит" из нативного мира, обрабатывается в Python, возвращается

# ХОРОШО: замена на встроенную функцию
from pyspark.sql import functions as F
result = df.withColumn("tax", df.amount * df.tax_rate / 100)
# В плане: CometProjectExec (полностью нативный)

Pandas UDF (@pandas_udf) ведут себя лучше с точки зрения производительности, но также разрывают нативный план Comet в точке своего вызова. Разница в том, что Pandas UDF передают Arrow-батчи напрямую (без построчной передачи), но Comet всё равно не может включить Pandas UDF в свой нативный план - они выполняются в Python-процессе.

Ограничения нативных Join-паттернов

Некоторые сложные join-паттерны не полностью поддерживаются:

  • Full Outer Join с non-equi условиями: a JOIN b ON a.id = b.id AND a.date BETWEEN b.start AND b.end - падает в JVM
  • Cross Join (CROSS JOIN): поддерживается частично
  • Multiple join conditions с OR: ON (a.x = b.x OR a.y = b.y) - как правило JVM
  • Self-join со сложными коррелированными подзапросами: зависит от того, как Catalyst разворачивает план

Join-паттерны и AQE

Adaptive Query Execution (AQE) может изменять тип join во время выполнения (например, превращать Sort Merge Join в Broadcast Hash Join, если одна сторона оказалась меньше ожидаемого). Comet поддерживает AQE и адаптируется к его решениям, но в момент адаптации возможна кратковременная пересборка плана, включая вставку переходных операторов.


Чтение планов выполнения: разбор explain()

Умение читать df.explain() при включённом Comet - ключевой навык для диагностики производительности. Это не просто "смотрим на имена операторов", а понимание того, где начинается и заканчивается нативное выполнение, и где данные пересекают границу JVM ↔ нативный мир.

Анатомия плана с Comet

Рассмотрим конкретный запрос и его план:

from pyspark.sql import functions as F

df = (
    spark.read.parquet("/data/sales/")
    .filter(F.col("sale_date") >= "2024-01-01")
    .filter(F.col("amount") > 100)
    .groupBy("region", "product_category")
    .agg(
        F.sum("amount").alias("total"),
        F.count("*").alias("cnt")
    )
)

df.explain(True)  # True = показать все фазы плана

Вывод (физический план с Comet):

== Physical Plan ==
CometHashAggregateExec(keys=[region, product_category], funcs=[sum(amount), count(1)])
+- CometHashAggregateExec(keys=[region, product_category], funcs=[partial_sum(amount), partial_count(1)])
   +- CometProjectExec[region, product_category, amount]
      +- CometFilterExec(condition=(sale_date >= 18262) AND (amount > 100))
         +- CometBatchScan[sale_date, amount, region, product_category]
            parquet.../data/sales/

Все операторы имеют префикс Comet - это означает, что весь план выполняется нативно, без единого перехода в JVM. Обратите внимание:

  • Два CometHashAggregateExec: первый - partial aggregate (на каждом Executor), второй - final aggregate (после shuffle). Это стандартный двухфазный паттерн агрегации в Spark.
  • sale_date >= 18262: Catalyst автоматически конвертировал строку '2024-01-01' в число дней от эпохи Unix. Это нормально - Comet работает с нативным представлением.

Идентификация границ нативного выполнения

Теперь рассмотрим план с fallback:

# Python UDF добавлен в середину пайплайна
@udf(returnType=DoubleType())
def apply_discount(amount, region):
    discounts = {"North": 0.10, "South": 0.05}
    return amount * (1 - discounts.get(region, 0))

df_with_discount = (
    spark.read.parquet("/data/sales/")
    .filter(F.col("amount") > 100)
    .withColumn("discounted", apply_discount(F.col("amount"), F.col("region")))
    .groupBy("region")
    .agg(F.sum("discounted").alias("total_discounted"))
)

df_with_discount.explain()

Вывод:

== Physical Plan ==
CometHashAggregateExec(keys=[region], funcs=[sum(discounted)])
+- CometRowToColumnarExec               <-- переход JVM → native
   +- HashAggregateExec                 <-- JVM агрегация
      +- PythonEvalUDF apply_discount   <-- Python UDF (строковый)
         +- ColumnarToRow               <-- переход native → JVM
            +- CometFilterExec(amount > 100)
               +- CometBatchScan[...]

Здесь видны все компоненты смешанного плана:

  1. CometBatchScan + CometFilterExec - нативное начало (Comet читает и фильтрует)
  2. ColumnarToRow - граница выхода из native: Arrow RecordBatch конвертируются в JVM InternalRow
  3. PythonEvalUDF - Python-процесс получает строки, обрабатывает, возвращает
  4. HashAggregateExec (без префикса Comet) - JVM-агрегация на JVM-строках
  5. CometRowToColumnarExec - граница возврата в native: JVM InternalRow конвертируются в Arrow
  6. CometHashAggregateExec - финальная нативная агрегация

Два перехода (ColumnarToRow и CometRowToColumnarExec) - это двойной overhead. Данные копируются дважды: туда и обратно.

Включение логирования причин деградации

Иногда неочевидно, почему именно оператор "упал" в JVM. Для диагностики:

spark.conf.set("spark.comet.logFallbackReasons.enabled", "true")
spark.conf.set("spark.comet.explainFallback.enabled", "true")

После включения в логах Driver при компиляции плана появятся строки вида:

[COMET] Operator HashAggregateExec cannot be converted to native:
  Reason: Expression 'apply_discount' is a Python UDF which is not supported in native execution.
  Suggestion: Replace Python UDF with equivalent built-in Spark SQL functions.

[COMET] Operator ProjectExec cannot be converted to native:
  Reason: Function 'to_date' with format 'dd-MMM-yyyy' uses locale-dependent parsing,
  not supported. Use ISO format 'yyyy-MM-dd' or cast from standard date string.

Эти сообщения дают точную диагностику: какой оператор не поддерживается и почему. В production рекомендуется включать это логирование на этапе разработки и проверки, но отключать в эксплуатации - оно генерирует значительный объём логов.

Использование Spark UI для анализа

В Spark UI вкладка "SQL / DataFrame" показывает физический план в виде графа. При включённом Comet нативные операторы отображаются с особым маркером. Вкладка "Stages" показывает распределение времени по задачам: если задачи, выполняемые нативно, занимают значительно меньше времени, чем JVM-задачи при аналогичном объёме данных - это верный признак работающего Comet.

Метрики в Spark UI для нативных стейджей:

  • Task Time снижается - меньше CPU на выполнение одной задачи
  • GC Time снижается до ~0% - нативная память не создаёт GC-давление
  • Peak Execution Memory может быть выше - Arrow-буферы активнее используют offHeap

Лабораторная практика: запуск TPC-H бенчмарка

TPC-H (Transaction Processing Performance Council - Benchmark H) - индустриальный стандарт оценки производительности аналитических систем. Он состоит из 22 SQL-запросов, имитирующих типичные BI-задачи: многотабличные JOIN, агрегации с фильтрацией, подзапросы, оконные функции. TPC-H специально разработан так, чтобы быть нечувствительным к кэшированию - каждый запрос работает с разными подмножествами данных.

Почему TPC-H - правильный бенчмарк для Comet

Comet ускоряет CPU-bound аналитические операции: агрегации, фильтрации, join. Именно это TPC-H и нагружает. Если бы мы тестировали на workload с большим количеством Python UDF или сложных строковых трансформаций - результаты были бы хуже. TPC-H показывает Comet в его сильнейшей стороне - "чистый" аналитический SQL.

Подготовка данных

# Генерация TPC-H данных с помощью tpch-dbgen
# Устанавливаем утилиту:
# sudo apt-get install -y cmake git
# git clone https://github.com/databricks/tpch-dbgen && cd tpch-dbgen && make

# Scale Factor 10 = ~10 GB данных (для локального тестирования)
# Scale Factor 100 = ~100 GB (для кластера)
import subprocess
import os

SF = 10  # изменить на 100 для кластера

# Генерация: для SF=10 займёт ~2-3 минуты
subprocess.run(["./dbgen", "-s", str(SF), "-f"], check=True)

# Конвертация TBL → Parquet через PySpark
from pyspark.sql import SparkSession

spark_prep = SparkSession.builder.appName("tpch-data-prep").getOrCreate()

tables = {
    "lineitem": "l_orderkey,l_partkey,l_suppkey,l_linenumber,l_quantity,l_extendedprice,"
                "l_discount,l_tax,l_returnflag,l_linestatus,l_shipdate,l_commitdate,"
                "l_receiptdate,l_shipinstruct,l_shipmode,l_comment",
    "orders":   "o_orderkey,o_custkey,o_orderstatus,o_totalprice,o_orderdate,"
                "o_orderpriority,o_clerk,o_shippriority,o_comment",
    "customer": "c_custkey,c_name,c_address,c_nationkey,c_phone,c_acctbal,"
                "c_mktsegment,c_comment",
    "part":     "p_partkey,p_name,p_mfgr,p_brand,p_type,p_size,p_container,"
                "p_retailprice,p_comment",
    "supplier": "s_suppkey,s_name,s_address,s_nationkey,s_phone,s_acctbal,s_comment",
    "partsupp": "ps_partkey,ps_suppkey,ps_availqty,ps_supplycost,ps_comment",
    "nation":   "n_nationkey,n_name,n_regionkey,n_comment",
    "region":   "r_regionkey,r_name,r_comment",
}

data_path = f"/data/tpch/sf{SF}"
os.makedirs(data_path, exist_ok=True)

for table_name, schema_str in tables.items():
    columns = schema_str.split(",")
    tbl_file = f"./{table_name}.tbl"
    print(f"Converting {table_name}...")
    df = spark_prep.read.csv(tbl_file, sep="|", inferSchema=True)
    # Берём только нужные колонки (dbgen добавляет лишний разделитель в конце)
    df = df.toDF(*columns)
    df.write.mode("overwrite").parquet(f"{data_path}/{table_name}")

print("Data prepared!")
spark_prep.stop()

Методология проведения бенчмарка

Корректный бенчмарк требует соблюдения нескольких правил:

  1. Прогревочный запуск: первый запуск любого Spark-приложения медленнее из-за JIT-компиляции JVM, прогрева OS page cache и инициализации Comet runtime. Не считайте первый результат.
  2. Сброс кэша: между запусками Baseline и Comet необходимо сбросить OS page cache (sudo echo 3 > /proc/sys/vm/drop_caches) и перезапустить Spark-кластер, чтобы данные не оставались в памяти.
  3. Изоляция переменных: запускайте Baseline и Comet на одном и том же кластере с одинаковыми ресурсами. Разница должна быть только в наличии/отсутствии Comet.
  4. Медиана, не среднее: запустите каждый запрос минимум 3 раза, берите медиану. Это устраняет outlier от GC-паузы или сетевого затора.
import time
from pyspark.sql import SparkSession

def create_spark(comet_enabled: bool) -> SparkSession:
    """Создаём SparkSession с Comet или без."""
    builder = (
        SparkSession.builder
        .appName(f"tpch-bench-{'comet' if comet_enabled else 'baseline'}")
        .config("spark.executor.memory", "8g")
        .config("spark.executor.cores", "4")
        .config("spark.memory.offHeap.enabled", "true")
        .config("spark.memory.offHeap.size", "4g")
        .config("spark.executor.memoryOverhead", "2g")
        .config("spark.sql.adaptive.enabled", "true")
    )

    if comet_enabled:
        builder = (
            builder
            .config("spark.plugins", "org.apache.spark.CometPlugin")
            .config("spark.sql.extensions",
                    "org.apache.spark.sql.comet.CometSparkSessionExtensions")
            .config("spark.comet.enabled", "true")
            .config("spark.comet.exec.all.enabled", "true")
            .config("spark.shuffle.manager",
                    "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager")
        )

    return builder.getOrCreate()


def run_tpch_queries(spark: SparkSession, data_path: str) -> dict:
    """Запускаем все TPC-H запросы и возвращаем таймингу."""

    # Регистрируем таблицы
    for table in ["lineitem", "orders", "customer", "part", "supplier",
                  "partsupp", "nation", "region"]:
        spark.read.parquet(f"{data_path}/{table}").createOrReplaceTempView(table)

    queries = {
        "Q1":  """
            SELECT l_returnflag, l_linestatus,
                   SUM(l_quantity) AS sum_qty,
                   SUM(l_extendedprice) AS sum_base_price,
                   SUM(l_extendedprice * (1 - l_discount)) AS sum_disc_price,
                   SUM(l_extendedprice * (1 - l_discount) * (1 + l_tax)) AS sum_charge,
                   AVG(l_quantity) AS avg_qty,
                   AVG(l_extendedprice) AS avg_price,
                   AVG(l_discount) AS avg_disc,
                   COUNT(*) AS count_order
            FROM lineitem
            WHERE l_shipdate <= date_sub(date('1998-12-01'), 90)
            GROUP BY l_returnflag, l_linestatus
            ORDER BY l_returnflag, l_linestatus
        """,
        "Q3": """
            SELECT l_orderkey,
                   SUM(l_extendedprice * (1 - l_discount)) AS revenue,
                   o_orderdate, o_shippriority
            FROM customer
            JOIN orders ON c_custkey = o_custkey
            JOIN lineitem ON l_orderkey = o_orderkey
            WHERE c_mktsegment = 'BUILDING'
              AND o_orderdate < date('1995-03-15')
              AND l_shipdate > date('1995-03-15')
            GROUP BY l_orderkey, o_orderdate, o_shippriority
            ORDER BY revenue DESC, o_orderdate
            LIMIT 10
        """,
        "Q6": """
            SELECT SUM(l_extendedprice * l_discount) AS revenue
            FROM lineitem
            WHERE l_shipdate >= date('1994-01-01')
              AND l_shipdate < date('1995-01-01')
              AND l_discount BETWEEN 0.05 AND 0.07
              AND l_quantity < 24
        """,
        "Q10": """
            SELECT c_custkey, c_name,
                   SUM(l_extendedprice * (1 - l_discount)) AS revenue,
                   c_acctbal, n_name, c_address, c_phone, c_comment
            FROM customer
            JOIN orders ON c_custkey = o_custkey
            JOIN lineitem ON l_orderkey = o_orderkey
            JOIN nation ON c_nationkey = n_nationkey
            WHERE o_orderdate >= date('1993-10-01')
              AND o_orderdate < date('1994-01-01')
              AND l_returnflag = 'R'
            GROUP BY c_custkey, c_name, c_acctbal, c_phone, n_name, c_address, c_comment
            ORDER BY revenue DESC
            LIMIT 20
        """,
        "Q18": """
            SELECT c_name, c_custkey, o_orderkey, o_orderdate, o_totalprice,
                   SUM(l_quantity) AS total_qty
            FROM customer
            JOIN orders ON o_custkey = c_custkey
            JOIN lineitem ON o_orderkey = l_orderkey
            WHERE o_orderkey IN (
                SELECT l_orderkey FROM lineitem GROUP BY l_orderkey
                HAVING SUM(l_quantity) > 300
            )
            GROUP BY c_name, c_custkey, o_orderkey, o_orderdate, o_totalprice
            ORDER BY o_totalprice DESC, o_orderdate
            LIMIT 100
        """,
    }

    results = {}
    for query_id, sql in queries.items():
        runs = []
        for attempt in range(3):
            start = time.time()
            spark.sql(sql).count()  # .count() форсирует полное выполнение
            elapsed = time.time() - start
            runs.append(elapsed)
            print(f"  {query_id} attempt {attempt+1}: {elapsed:.2f}s")

        median_time = sorted(runs)[1]  # медиана из 3 запусков
        results[query_id] = median_time
        print(f"  {query_id} MEDIAN: {median_time:.2f}s\n")

    return results

Запуск и сравнение результатов

SF = 10
data_path = f"/data/tpch/sf{SF}"

# --- Baseline: обычный Spark без Comet ---
print("=== BASELINE (JVM Spark) ===")
spark_baseline = create_spark(comet_enabled=False)
baseline_results = run_tpch_queries(spark_baseline, data_path)
spark_baseline.stop()

# Сбрасываем кэш между запусками!
# os.system("sudo bash -c 'echo 3 > /proc/sys/vm/drop_caches'")

# --- Comet: Spark с нативным ускорением ---
print("=== COMET (Native Execution) ===")
spark_comet = create_spark(comet_enabled=True)
comet_results = run_tpch_queries(spark_comet, data_path)
spark_comet.stop()

# --- Анализ результатов ---
print("\n=== РЕЗУЛЬТАТЫ ===")
print(f"{'Query':<8} {'Baseline':>12} {'Comet':>10} {'Speedup':>10} {'Operator':<25}")
print("-" * 65)

operator_hints = {
    "Q1":  "Pure aggregation (best case for Comet)",
    "Q3":  "3-table join + aggregation",
    "Q6":  "Filter + simple sum (SIMD-heavy)",
    "Q10": "4-table join + aggregation",
    "Q18": "Subquery + join + aggregation",
}

for qid in sorted(baseline_results.keys()):
    base_t = baseline_results[qid]
    comet_t = comet_results[qid]
    speedup = base_t / comet_t
    hint = operator_hints.get(qid, "")
    print(f"{qid:<8} {base_t:>10.2f}s {comet_t:>9.2f}s {speedup:>9.2f}x  {hint}")

Интерпретация результатов

Типичные результаты на TPC-H SF=10 при правильной конфигурации:

Запрос Baseline Comet Ускорение Доминирующая операция
Q1 45s 12s 3.7× Aggregation только на lineitem
Q3 78s 28s 2.8× 3-way join + GROUP BY
Q6 18s 5s 3.6× Фильтр + SUM (SIMD-heavy)
Q10 95s 31s 3.1× 4-way join + aggregation
Q18 120s 48s 2.5× Subquery + join

Почему Q1 и Q6 ускоряются больше всего?

Q1 и Q6 - CPU-bound в чистом виде: они читают один файл (lineitem) и выполняют агрегацию или фильтрацию. Никаких shuffle или complex join. Весь план умещается в одну цепочку нативных операторов без ни одного fallback. SIMD-векторизация фильтров и агрегаций даёт максимальный эффект.

Q18 ускоряется скромнее, потому что содержит коррелированный подзапрос (IN (SELECT ...)). Catalyst разворачивает его в semi-join, который Comet поддерживает, но с некоторыми ограничениями. Часть плана может содержать переходы JVM ↔ native.

Когда Comet даёт мало или ничего:

  • IO-bound запросы, где время тратится на чтение с диска, а не на CPU. Comet ускоряет CPU-фазу, но не сеть и диск.
  • Запросы с высокой долей Python UDF - вся векторизация нивелируется overhead'ом Python IPC.
  • Запросы, где Catalyst генерирует план с операторами, которые Comet не поддерживает, что приводит к многократным JVM ↔ native переходам.

Best Practices и анти-паттерны

Правила написания SQL-кода для максимального покрытия Comet

Следующие практики позволяют минимизировать fallback и обеспечить максимальное нативное выполнение:

Используйте встроенные функции Spark вместо UDF:

# ПЛОХО: Python UDF блокирует native execution
@udf(returnType=StringType())
def categorize(amount):
    if amount > 10000: return "large"
    elif amount > 1000: return "medium"
    else: return "small"

df.withColumn("category", categorize(df.amount))

# ХОРОШО: F.when полностью нативный
df.withColumn("category",
    F.when(F.col("amount") > 10000, "large")
     .when(F.col("amount") > 1000, "medium")
     .otherwise("small")
)

Используйте ISO-форматы дат:

# ПЛОХО: locale-зависимый формат → fallback
F.to_date(F.col("date_str"), "dd-MMM-yyyy")   # "15-Jan-2024"

# ХОРОШО: ISO → нативный Comet
F.to_date(F.col("date_str"), "yyyy-MM-dd")    # "2024-01-15"

Избегайте смешивания несовместимых типов в одном операторе:

# ПЛОХО: MapType в одной таблице с числовыми колонками
# тянет весь оператор в JVM
df_with_map = df.withColumn("metadata", F.create_map(...))  # MapType
df_with_map.groupBy("region").agg(F.sum("amount"))
# sum(amount) может упасть в JVM из-за MapType в schema!

# ХОРОШО: вынесите map-операции отдельно
df_clean = df.select("region", "amount")
df_clean.groupBy("region").agg(F.sum("amount"))  # чисто нативно

Проверяйте покрытие плана перед деплоем:

def check_comet_coverage(df) -> float:
    """Возвращает долю нативных операторов в плане (0.0 - 1.0)."""
    plan_str = df._jdf.queryExecution().executedPlan().toString()
    comet_ops = plan_str.count("Comet")
    total_ops = sum(plan_str.count(op) for op in [
        "Exec", "Join", "Aggregate", "Project", "Filter", "Scan", "Exchange"
    ])
    return comet_ops / max(total_ops, 1)

df = spark.sql("SELECT region, SUM(amount) FROM sales GROUP BY region")
coverage = check_comet_coverage(df)
print(f"Native coverage: {coverage:.0%}")
# Если < 70% - ищем причину fallback через explainFallback

Анти-паттерны при внедрении Comet

Анти-паттерн 1: Включить Comet и считать задачу решённой

Ошибка думать, что spark.comet.enabled=true автоматически всё ускоряет. Необходимо проверить df.explain() для ключевых запросов и убедиться, что операторы действительно нативные. Без проверки легко оказаться в ситуации, когда Comet включён, но каждый запрос падает в JVM из-за одного Python UDF.

Анти-паттерн 2: Пропустить настройку памяти

Не настроить spark.memory.offHeap.size и spark.executor.memoryOverhead - верный путь к нестабильным OOM в продакшене. YARN или Kubernetes убьют Executor, когда нативная память Arrow-аллокатора превысит лимит контейнера.

Анти-паттерн 3: Сравнивать результаты без сброса кэша

Запустить Baseline, затем сразу Comet на тех же данных - и обнаружить, что Comet "в 10 раз быстрее". На самом деле это OS page cache: второй запуск работает с данными в RAM, а не читает с диска. Всегда сбрасывайте кэш между сравнениями.

Анти-паттерн 4: Использовать Comet для row-oriented ETL

Comet спроектирован под аналитические (OLAP) запросы: большие объёмы данных, агрегации, join. Для ETL-пайплайнов с построчными трансформациями, форматами без колоночного хранения (JSON, CSV row-by-row) или большим количеством Python-логики Comet даст минимальный выигрыш.


Производственное внедрение

Стратегия canary rollout

Не включайте Comet сразу для всех продакшн-задач. Используйте постепенный подход:

# Этап 1: отдельная Spark-конфигурация для тестовых джоб
spark_comet_test = create_spark(comet_enabled=True)

# Этап 2: проверка корректности результатов
baseline_result = spark_baseline.sql(query).collect()
comet_result = spark_comet_test.sql(query).collect()

# Проверяем числовую эквивалентность с допуском (float арифметика)
import math
for b_row, c_row in zip(
    sorted(baseline_result, key=lambda r: r[0]),
    sorted(comet_result, key=lambda r: r[0])
):
    for col_idx in range(len(b_row)):
        b_val, c_val = b_row[col_idx], c_row[col_idx]
        if isinstance(b_val, float):
            assert math.isclose(b_val, c_val, rel_tol=1e-6), \
                f"Mismatch at col {col_idx}: {b_val} vs {c_val}"
        else:
            assert b_val == c_val, f"Mismatch: {b_val} vs {c_val}"

print("Results match! Safe to enable Comet in production.")

Мониторинг в production

После включения Comet в продакшне следите за:

  • GC Time в Spark UI: должен снизиться (меньше heap-аллокаций в нативных операторах)
  • Task Duration: должны сократиться для CPU-bound задач
  • Executor OOM частота: если растёт - увеличьте spark.memory.offHeap.size
  • Fallback warnings в логах: если появляются - диагностируйте и исправляйте источник

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

Условие задачи:

Вам выдан PySpark-скрипт для расчёта KPI по продажам. На обычном Spark джоба работает 15 минут. После включения Comet время сократилось лишь до 14 минут - почти никакого улучшения.

# Скрипт с "проблемой" - найдите и исправьте
from pyspark.sql import functions as F
from pyspark.sql.types import DoubleType
from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .config("spark.plugins", "org.apache.spark.CometPlugin")
    .config("spark.sql.extensions",
            "org.apache.spark.sql.comet.CometSparkSessionExtensions")
    .config("spark.comet.enabled", "true")
    .config("spark.comet.exec.all.enabled", "true")
    .config("spark.comet.logFallbackReasons.enabled", "true")
    .config("spark.comet.explainFallback.enabled", "true")
    .getOrCreate()
)

@udf(returnType=DoubleType())
def apply_vat(price, vat_rate):
    return price * (1.0 + vat_rate / 100.0)

df = (
    spark.read.parquet("/data/sales_large/")
    .withColumn("price_with_vat",
                apply_vat(F.col("net_price"), F.col("vat_rate")))
    .groupBy("region", "year")
    .agg(
        F.sum("price_with_vat").alias("total_revenue_vat"),
        F.count("*").alias("transaction_count"),
        F.avg("net_price").alias("avg_net_price")
    )
)

df.explain()
df.write.mode("overwrite").parquet("/data/output/kpi_by_region/")

Задание:

  1. Включите spark.comet.logFallbackReasons.enabled = true и spark.comet.explainFallback.enabled = true. Запустите скрипт и изучите вывод df.explain() и логи Driver. Определите, какой оператор вызывает fallback на JVM и почему.

  2. Исправьте скрипт так, чтобы весь план выполнялся нативно в Comet. Подсказка: замените проблемный оператор встроенной функцией Spark SQL.

  3. Проведите сравнительный бенчмарк: запустите исходный скрипт (с проблемой) и исправленный скрипт на одном датасете объёмом не менее 5 GB. Зафиксируйте время выполнения.

  4. Проверьте df.explain() после исправления: убедитесь, что в плане нет ни одного оператора без префикса Comet (кроме CometColumnarExchange при необходимости shuffle).

  5. Приложите: итоговый скрипт, вывод df.explain() до и после исправления, таблицу сравнения времени (Baseline JVM / Comet с UDF / Comet нативный). Ожидаемое ускорение после исправления - не менее 2.5× по сравнению с Baseline JVM.