Comet: установка, операторы, ограничения и TPC-H бенчмарк
Установка jar и конфигурация Comet, поддерживаемые операторы, что падает на JVM, запуск TPC-H бенчмарка.
Практический старт: Установка 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 считают всю память процесса - включая нативную.
Если не выделить дополнительную память, произойдёт следующее:
- Executor JVM запросил 4 GB (
spark.executor.memory) - Comet выделил ещё 2 GB нативной памяти для Arrow-буферов
- YARN видит процесс, использующий 6 GB, при лимите 4.5 GB
- YARN убивает Executor как превысивший лимит
- 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:
- Первый проход (partial aggregate): каждый Executor обрабатывает свою партицию данных. Для каждого уникального ключа группировки строится запись в нативной хэш-таблице (Rust
HashMapс открытой адресацией для кэш-дружественного доступа). - Shuffle: промежуточные результаты перемешиваются по ключу группировки между Executor-ами.
- Второй проход (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) и конвертирует их в ArrowRecordBatchдля передачи в нативный рантайм. Эта конвертация не бесплатна - она требует копирования данных.ColumnarToRow: принимает Arrow RecordBatch из нативного рантайма и конвертирует их обратно в JVMInternalRow. Также требует копирования.
Каждая такая "граница" между 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[...]
Здесь видны все компоненты смешанного плана:
CometBatchScan+CometFilterExec- нативное начало (Comet читает и фильтрует)ColumnarToRow- граница выхода из native: Arrow RecordBatch конвертируются в JVM InternalRowPythonEvalUDF- Python-процесс получает строки, обрабатывает, возвращаетHashAggregateExec(без префикса Comet) - JVM-агрегация на JVM-строкахCometRowToColumnarExec- граница возврата в native: JVM InternalRow конвертируются в ArrowCometHashAggregateExec- финальная нативная агрегация
Два перехода (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()
Методология проведения бенчмарка¶
Корректный бенчмарк требует соблюдения нескольких правил:
- Прогревочный запуск: первый запуск любого Spark-приложения медленнее из-за JIT-компиляции JVM, прогрева OS page cache и инициализации Comet runtime. Не считайте первый результат.
- Сброс кэша: между запусками Baseline и Comet необходимо сбросить OS page cache (
sudo echo 3 > /proc/sys/vm/drop_caches) и перезапустить Spark-кластер, чтобы данные не оставались в памяти. - Изоляция переменных: запускайте Baseline и Comet на одном и том же кластере с одинаковыми ресурсами. Разница должна быть только в наличии/отсутствии Comet.
- Медиана, не среднее: запустите каждый запрос минимум 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/")
Задание:
-
Включите
spark.comet.logFallbackReasons.enabled = trueиspark.comet.explainFallback.enabled = true. Запустите скрипт и изучите выводdf.explain()и логи Driver. Определите, какой оператор вызывает fallback на JVM и почему. -
Исправьте скрипт так, чтобы весь план выполнялся нативно в Comet. Подсказка: замените проблемный оператор встроенной функцией Spark SQL.
-
Проведите сравнительный бенчмарк: запустите исходный скрипт (с проблемой) и исправленный скрипт на одном датасете объёмом не менее 5 GB. Зафиксируйте время выполнения.
-
Проверьте
df.explain()после исправления: убедитесь, что в плане нет ни одного оператора без префиксаComet(кромеCometColumnarExchangeпри необходимости shuffle). -
Приложите: итоговый скрипт, вывод
df.explain()до и после исправления, таблицу сравнения времени (Baseline JVM / Comet с UDF / Comet нативный). Ожидаемое ускорение после исправления - не менее 2.5× по сравнению с Baseline JVM.