RAPIDS cuDF: GPU acceleration, CUDA требования и профиль данных
RAPIDS cuDF: GPU acceleration, CUDA требования и профиль данных
Архитектурный сдвиг: когда CPU уступает место GPU¶
До этого урока мы рассматривали нативные ускорители - Comet (Rust + DataFusion), Gluten + Velox (C++ + AVX-512 SIMD), Photon (C++ от Databricks). Все они ускоряют аналитические вычисления, оставаясь в парадигме CPU: десятки физических ядер, SIMD-инструкции, AVX-512 регистры по 512 бит. Это мощно. Но существует качественно другой уровень параллелизма - GPU.
Два мира параллелизма: CPU vs GPU¶
Процессор (CPU) спроектирован для последовательного выполнения с низкой задержкой: 8–64 мощных ядра, глубокий OOO (out-of-order execution), большой L3-кэш (64–192 MB), тактовая частота 3–5 GHz. Каждое ядро умеет делать сложные вещи - ветвления, рекурсию, виртуальные вызовы. SIMD-расширения (AVX-512) добавляют до 16-кратного параллелизма на уровне данных внутри одного ядра.
Видеокарта (GPU) спроектирована для массивно-параллельного выполнения с высокой пропускной способностью: NVIDIA H100 содержит 16896 CUDA-ядер, организованных в 528 потоковых мультипроцессоров (SM). Каждое ядро простое - без OOO, без глубокого кэша. Но их тысячи, и они работают одновременно.
Ключевое отличие в парадигме исполнения: CPU использует SIMD (Single Instruction, Multiple Data) - одна инструкция обрабатывает 8–16 элементов в одном ядре. GPU использует SIMT (Single Instruction, Multiple Threads) - одна инструкция выполняется в 32 потоках одновременно (warp). Для аналитических задач, где нужно применить одну и ту же операцию к миллиарду значений, SIMT-архитектура GPU - идеальный инструмент.
Ещё одно принципиальное отличие: пропускная способность памяти. Флагманский CPU (например, AMD EPYC 9654) имеет пропускную способность DDR5 DRAM около 460 GB/s. NVIDIA H100 SXM с HBM3 (High Bandwidth Memory) - 3.35 TB/s: в семь раз больше. Для колоночных аналитических задач, где основная работа - загрузка миллиардов значений из памяти и применение арифметики, это фундаментальное преимущество.
| Характеристика | CPU (EPYC 9654) | GPU (NVIDIA H100 SXM) |
|---|---|---|
| Ядра | 96 физических | 16896 CUDA-ядер |
| Потоки одновременно | 192 (SMT) | ~50000+ (warps) |
| Пропускная память | ~460 GB/s | 3.35 TB/s |
| Объём памяти | 768 GB DDR5 | 80 GB HBM3 |
| TDP | 360W | 700W |
| Оптимум | Latency-sensitive, complex logic | Throughput, uniform operations |
Когда GPU нужен, а когда - нет¶
GPU не является универсальным ускорителем. Он эффективен только при определённом профиле задачи:
GPU эффективен (существенное ускорение):
- Сканирование и фильтрация сотен миллионов строк Parquet
GROUP BYс агрегациями по большим таблицам (миллиарды строк)JOINдвух больших таблиц (не broadcast - оба датасета большие)- Сортировка (radix sort на GPU параллелен по природе)
- ML feature engineering: нормализация, трансформации числовых колонок
- Декомпрессия Parquet (Snappy, Zstd) - GPU умеет распаковывать блоки параллельно
GPU неэффективен (нет или минус):
- Маленькие датасеты (< 1 GB): overhead на transfer CPU→GPU съедает выигрыш
- Python UDF: всё равно нужен CPU-процесс, GPU простаивает
- Глубоко вложенные структуры (
Array<Map<String, Struct<...>>>): GPU хочет плоские колонки - Реляционный ETL с множеством
CASE WHENи строковых преобразований - Задачи с высоким IO skew: если 90% времени - ожидание диска/сети, GPU не поможет
RAPIDS Accelerator: как GPU встраивается в Spark¶
RAPIDS Accelerator for Apache Spark - это плагин от NVIDIA, который встраивается в Spark через тот же механизм, что Comet и Gluten: Spark Plugin API + SparkSessionExtensions. Catalyst optimizer, AQE, scheduler, API - всё остаётся нетронутым. Меняется только физическое исполнение операторов на Executor-е.
Принцип работы такой же, как у CPU-нативных движков: ColumnarRule перехватывает физический план и заменяет поддерживаемые операторы GPU-аналогами: FilterExec → GpuFilterExec, HashAggregateExec → GpuHashAggregateExec, SortMergeJoinExec → GpuSortMergeJoinExec. Данные при этом перемещаются из host-памяти (оперативная память CPU) в device-память (видеопамять GPU).
Инфраструктурный фундамент и CUDA-требования¶
Аппаратные требования: какой GPU нужен¶
RAPIDS Accelerator поддерживает только GPU NVIDIA с вычислительной мощностью (Compute Capability) 6.0 и выше. На практике для production-использования рекомендуются:
| GPU | Поколение | VRAM | Пропускная | Compute Capability | Кейс |
|---|---|---|---|---|---|
| T4 | Turing | 16 GB GDDR6 | 320 GB/s | 7.5 | Инференс, лёгкая аналитика |
| A10G | Ampere | 24 GB GDDR6 | 600 GB/s | 8.6 | Balanced analytics |
| A100 40GB | Ampere | 40 GB HBM2e | 1.6 TB/s | 8.0 | Production analytics |
| A100 80GB | Ampere | 80 GB HBM2e | 2.0 TB/s | 8.0 | Большие JOIN |
| H100 SXM | Hopper | 80 GB HBM3 | 3.35 TB/s | 9.0 | Максимум |
| L40S | Ada | 48 GB GDDR6 | 864 GB/s | 8.9 | Cost-efficient |
Объём VRAM - главное ограничение GPU-аналитики. Если одна партиция join-таблицы не вмещается в VRAM, начинается spill в host-память через PCIe - это в 5–10 раз медленнее работы из VRAM. Для production с TPC-H SF=100 (100 GB данных, крупные join) рекомендуется A100 80GB или H100.
Программный стек: CUDA, драйверы, JAR¶
RAPIDS Accelerator требует точного соответствия версий на каждом уровне стека:
NVIDIA GPU Driver (версия ≥ 525)
│
CUDA Runtime (11.8 или 12.x)
│
cuDF C++ library (собирается под конкретную CUDA)
│
RAPIDS Accelerator JAR (rapids-4-spark_2.12-XX.YY.ZZ.jar)
│
Apache Spark (3.3 / 3.4 / 3.5)
│
JVM (JDK 8 / 11 / 17)
Несовпадение версий на любом уровне - и плагин либо не загрузится, либо упадёт с CUDA error: invalid device function или CUDA driver version is insufficient for CUDA runtime version.
# Проверка версии драйвера и CUDA
nvidia-smi
# Вывод:
# Driver Version: 535.104.05 CUDA Version: 12.2
# Проверка Compute Capability конкретного GPU
nvidia-smi --query-gpu=name,compute_cap --format=csv
# Вывод:
# name, compute_cap
# NVIDIA A100-SXM4-80GB, 8.0
# Загрузка RAPIDS JAR (версия под Spark 3.5, CUDA 12.x)
RAPIDS_VERSION="24.02.0"
wget https://repo1.maven.org/maven2/com/nvidia/rapids-4-spark_2.12/${RAPIDS_VERSION}/rapids-4-spark_2.12-${RAPIDS_VERSION}-cuda12.jar
Совместимая матрица версий для наиболее часто используемых комбинаций:
| RAPIDS | Spark | CUDA | Compute Cap. min |
|---|---|---|---|
| 24.02 | 3.3, 3.4, 3.5 | 11.8 или 12.x | 6.0 |
| 23.12 | 3.3, 3.4 | 11.8 или 12.0 | 6.0 |
| 23.08 | 3.3, 3.4 | 11.5–12.0 | 6.0 |
Конфигурация GPU-ресурсов в Kubernetes и YARN¶
В отличие от CPU (который всегда доступен для Executor), GPU - это дискретный ресурс, который должен быть явно запрошен и выделен планировщиком кластера.
Kubernetes:
# Конфигурация GPU resources для Spark Executor Pod
apiVersion: v1
kind: Pod
spec:
containers:
- name: spark-executor
resources:
requests:
cpu: "8"
memory: "32Gi"
nvidia.com/gpu: "1" # запрашиваем 1 GPU на Executor
limits:
nvidia.com/gpu: "1" # строгий лимит
env:
- name: NVIDIA_VISIBLE_DEVICES
value: all
- name: NVIDIA_DRIVER_CAPABILITIES
value: compute,utility
YARN:
# yarn-site.xml: включение GPU как ресурса
# yarn.resource-types = yarn.io/gpu
# yarn.nodemanager.resource-plugins = io.yarn.server.nodemanager.containermanager.resourceplugin.gpu.GpuResourcePlugin
# spark-submit для YARN с GPU
spark-submit \
--master yarn \
--deploy-mode cluster \
--num-executors 4 \
--executor-cores 8 \
--executor-memory 32g \
--conf spark.executor.resource.gpu.amount=1 \
--conf spark.task.resource.gpu.amount=1 \
--conf spark.executor.resource.gpu.discoveryScript=/usr/lib/spark/scripts/getGpusResources.sh \
--jars rapids-4-spark_2.12-24.02.0-cuda12.jar \
your_job.py
Параметр spark.executor.resource.gpu.amount=1 сообщает Spark-планировщику, что каждый Executor должен получить один GPU. Скрипт getGpusResources.sh запускается на каждой ноде и возвращает список доступных GPU (например, {"name": "gpu", "addresses": ["0", "1"]}).
Магия cuDF под капотом Spark¶
Полная архитектура: от Python до CUDA-ядер¶
Рассмотрим каждый компонент этой архитектуры детально.
cuDF: DataFrame на GPU¶
cuDF (CUDA DataFrame) - это C++ библиотека от NVIDIA, которая реализует колоночный DataFrame API на GPU. Внутри cuDF каждая колонка - это cudf::column: структура, хранящая данные в device-памяти (VRAM) в виде плоского массива с отдельным null-bitmap.
Формат данных cuDF совместим с Apache Arrow: структура cudf::table (набор cudf::column) изоморфна Arrow RecordBatch. Это позволяет передавать данные между JVM Arrow-буферами и GPU cuDF-таблицами с минимальным преобразованием - нужно только скопировать данные через PCIe (CPU RAM → GPU VRAM), не меняя их структуры.
CUDA-ядра (kernels) в cuDF написаны на CUDA C++. Каждый kernel запускается на тысячах параллельных потоков. Например, kernel для фильтрации amount > 1000.0 над батчем из 8 миллионов строк:
// CUDA kernel для предиката amount > threshold (упрощённо)
__global__ void filter_double_gt(
const double* __restrict__ amounts,
const uint8_t* __restrict__ validity, // null-bitmap
bool* __restrict__ result,
int64_t n_rows,
double threshold
) {
// Каждый поток обрабатывает один элемент
int64_t idx = blockIdx.x * blockDim.x + threadIdx.x;
if (idx >= n_rows) return;
// Проверяем validity bit (не null)
bool is_valid = (validity[idx / 8] >> (idx % 8)) & 1;
// Применяем предикат
result[idx] = is_valid && (amounts[idx] > threshold);
}
// Запуск: 8M строк / 256 потоков на блок = 31250 блоков
// Все 31250 блоков запускаются параллельно → весь массив проверяется за ~5-10 мкс
filter_double_gt<<<31250, 256>>>(d_amounts, d_validity, d_result, 8000000, 1000.0);
На H100 с 16896 CUDA-ядрами за одно мгновение обрабатывается 16896 строк одновременно. Для 8 миллионов строк - около 500 "волн" по 16896 потоков. Суммарное время: единицы миллисекунд против сотен миллисекунд на JVM.
GPU Scan: параллельное чтение Parquet¶
RAPIDS реализует собственный Parquet-ридер на GPU. Он работает следующим образом:
- CPU читает метаданные: footer Parquet-файла (схема, row group statistics, column chunk offsets) читается CPU в host-память. Это лёгкая операция.
- CPU передаёт задачу GPU: список compressed column chunks (байтовые диапазоны файла) отправляется в cuDF.
- GPU декомпрессирует параллельно: каждый column chunk (Snappy/Zstd-блок) декомпрессируется отдельным CUDA-потоком. На GPU одновременно могут работать тысячи декомпрессоров.
- GPU декодирует encoding: RLE (Run-Length Encoding), Delta encoding, Dictionary encoding - всё декодируется параллельными CUDA-kernels.
- Результат: готовая
cudf::tableуже в VRAM - никакого хождения через CPU.
Сравнение: JVM Parquet-ридер читает column chunk в CPU RAM, декомпрессирует в CPU RAM, декодирует в CPU RAM, кладёт в ColumnarBatch в JVM heap, затем Spark копирует данные в GPU. RAPIDS пропускает три промежуточных буфера.
PCIe: узкое место передачи данных¶
Единственный путь между host-памятью (CPU RAM) и device-памятью (GPU VRAM) - это шина PCIe. PCIe 4.0 x16 обеспечивает ~32 GB/s в одном направлении. Это звучит внушительно, пока не сравнить с HBM3 в H100 (3.35 TB/s): PCIe в 100 раз медленнее внутренней GPU-памяти.
Практическое следствие: минимизируйте количество host↔device переносов. Идеальный pipeline - данные загружаются в VRAM один раз (при чтении Parquet), все трансформации происходят внутри GPU, результат (уже маленький после агрегации) возвращается в host-память один раз. Каждый лишний перенос (например, из-за fallback CPU-оператора в середине плана) "разрывает" GPU-pipeline и добавляет PCIe-overhead.
GPU Shuffle: UCX и GPUDirect RDMA¶
Shuffle - самая сложная операция в контексте GPU, потому что она требует передачи данных между Executor-ами (разными серверами сети).
Стандартный путь (без GPUDirect):
GPU A (сервер 1) → PCIe → CPU RAM → сеть → CPU RAM → PCIe → GPU B (сервер 2)
Данные дважды проходят через PCIe и оба раза обрабатываются CPU. Это нивелирует GPU-ускорение для shuffle-тяжёлых запросов.
GPU Shuffle с UCX и GPUDirect RDMA:
UCX (Unified Communication X) - это коммуникационный фреймворк, поддерживающий GPUDirect RDMA (Remote Direct Memory Access):
GPU A (сервер 1) → NVLink/InfiniBand RDMA → GPU B (сервер 2)
CPU не участвует в передаче данных - сетевой адаптер напрямую читает из VRAM одного сервера и пишет в VRAM другого. Это возможно, когда InfiniBand-карта и GPU находятся в одном PCIe-домене (что типично для DGX-систем NVIDIA).
RAPIDS Accelerator включает UCX shuffle:
.config("spark.shuffle.manager",
"com.nvidia.spark.rapids.spark35.RapidsShuffleManager")
.config("spark.rapids.shuffle.mode", "UCX")
# UCX endpoint: InfiniBand или TCP fallback
.config("spark.rapids.shuffle.ucx.useWakeup", "true")
На кластере без InfiniBand UCX автоматически деградирует к TCP-транспорту через обычную сеть, но по крайней мере shuffle-данные всё равно хранятся в Arrow IPC-формате и передаются колоночными батчами.
Тюнинг памяти и конфигурация RAPIDS плагина¶
RMM: RAPIDS Memory Manager¶
Главная проблема GPU-memory менеджмента - аллокации cudaMalloc() медленны. Каждый cudaMalloc требует обращения к CUDA-драйверу и может занимать сотни микросекунд. При обработке тысяч Arrow RecordBatch за задачу это превращается в значительный overhead.
RMM решает эту проблему через pool-аллокатор: при старте Executor RMM резервирует большой блок VRAM сразу, а затем раздаёт его кусками через быструю внутреннюю структуру данных без обращения к CUDA-драйверу.
Два основных режима RMM:
ARENA (рекомендуется для production): использует арену - большой непрерывный блок VRAM с bump allocator. Очень быстрый (O(1) allocation), но может иметь фрагментацию при долгой работе.
ASYNC (CUDA 11.2+): использует CUDA stream-ordered memory allocator - аллокации освобождаются автоматически при завершении CUDA-стрима. Меньше фрагментации, чуть медленнее.
# Конфигурация RMM
.config("spark.rapids.memory.gpu.allocator", "ARENA")
# Какую долю VRAM зарезервировать под RMM-пул (0.0–1.0)
# Остаток оставляем под CUDA context, нативные буферы
.config("spark.rapids.memory.gpu.reserve", "0.1") # 10% резерв
# Минимальный размер пула (байты). При старте RMM аллоцирует не менее этого
.config("spark.rapids.memory.gpu.minAllocFraction", "0.25")
# Максимальный размер пула (доля от VRAM)
.config("spark.rapids.memory.gpu.maxAllocFraction", "0.80")
# Пул host-памяти для staging буферов (PCIe transfer)
.config("spark.rapids.memory.host.spillStorageSize", "8589934592") # 8 GB
Для GPU с 80 GB VRAM (A100 80GB):
- 10% резерв = 8 GB под CUDA context
- 80% max pool = 64 GB под RMM (операционная память)
- Оставшиеся 10% = 8 GB для burst аллокаций вне пула
Управление размером батча: batchSizeRows¶
В Spark данные обрабатываются партициями. RAPIDS разбивает каждую партицию на батчи - порции данных, которые вместе загружаются в VRAM и обрабатываются как единица. Размер батча напрямую влияет на:
- Использование VRAM: большие батчи лучше утилизируют GPU-параллелизм, но могут не влезть в память
- PCIe overhead: маленькие батчи порождают больше host↔device транзакций
# Размер батча в строках (по умолчанию ~1 млн строк для широких таблиц)
.config("spark.rapids.sql.reader.batchSizeRows", "1000000")
# Размер батча в байтах (альтернативное ограничение)
.config("spark.rapids.sql.reader.batchSizeBytes", "1073741824") # 1 GB
# RAPIDS берёт минимум из двух ограничений
Формула расчёта максимального batchSizeRows для безопасной работы:
batchSizeBytes = VRAM_total × max_fraction × 0.5 / num_concurrent_tasks
= 80 GB × 0.80 × 0.5 / 2 tasks = 16 GB на батч
batchSizeRows = batchSizeBytes / avg_row_size_bytes
= 16 GB / 200 bytes = 80 M строк
Делим на 2, потому что во время join одновременно нужна память под build-side и probe-side батчи.
GPU Spill: что происходит при переполнении VRAM¶
Когда RMM-пул исчерпан (например, build-side hash table join не вмещается в VRAM), RAPIDS не падает с OOM - он спиллит данные в host-память через PCIe, а если и та заполнена - на диск.
Иерархия spill:
GPU VRAM (быстро, ~3.35 TB/s)
↓ overflow
Host RAM (средне, ~32 GB/s PCIe 4.0)
↓ overflow
Local NVMe SSD (медленно, ~7 GB/s)
Каждый уровень медленнее предыдущего в 10–500 раз. Поэтому задача конфигурации - подобрать параметры так, чтобы spill в host RAM случался редко, а в NVMe - только в экстремальных ситуациях.
# Настройка spill
.config("spark.rapids.memory.gpu.spillBatchSizeBytes", "536870912") # 512 MB per spill
.config("spark.rapids.memory.host.spillStorageSize", "8589934592") # 8 GB host spill
Механизм fallback на CPU¶
Когда RAPIDS встречает оператор, который не поддерживается на GPU (Python UDF, сложный тип данных, неподдерживаемая функция), он вставляет операторы переключения контекста:
GpuColumnarToRowExec: конвертирует cuDF-таблицу (GPU) в JVM InternalRow (CPU). Данные копируются из VRAM в CPU RAM через PCIe.GpuRowToColumnarExec: обратное преобразование - JVM InternalRow → cuDF-таблица, копирование из CPU RAM в VRAM.
Каждая такая пара - это дополнительный PCIe-transfer. Если в плане таких пар много, GPU простаивает и ждёт CPU.
Конфигурация SparkSession с RAPIDS¶
Полная конфигурация для production¶
from pyspark.sql import SparkSession
def create_rapids_spark(
app_name: str,
gpu_vram_gb: int = 80,
) -> SparkSession:
"""
Production-ready SparkSession с RAPIDS Accelerator.
Параметры рассчитаны для A100 80GB:
- VRAM: 80 GB (резерв 10% = 8 GB под CUDA, 72 GB под RMM)
- Host RAM: 32 GB (heap 16 GB + offHeap 8 GB + overhead 8 GB)
"""
rmm_max_bytes = int(gpu_vram_gb * 0.80 * (1024 ** 3))
host_spill_bytes = 8 * (1024 ** 3)
return (
SparkSession.builder
.appName(app_name)
# RAPIDS plugin registration (через spark.plugins)
.config("spark.plugins",
"com.nvidia.spark.SQLPlugin")
# Активация GPU execution
.config("spark.rapids.sql.enabled", "true")
# Ресурс GPU для Executor
.config("spark.executor.resource.gpu.amount", "1")
.config("spark.task.resource.gpu.amount", "1")
# Память JVM Executor
.config("spark.executor.memory", "16g")
.config("spark.executor.memoryOverhead", "8g")
.config("spark.memory.offHeap.enabled", "true")
.config("spark.memory.offHeap.size", "8g")
# RMM pool
.config("spark.rapids.memory.gpu.allocator", "ARENA")
.config("spark.rapids.memory.gpu.reserve", "0.10")
.config("spark.rapids.memory.gpu.maxAllocFraction", "0.80")
.config("spark.rapids.memory.host.spillStorageSize",
str(host_spill_bytes))
# Размер батча
.config("spark.rapids.sql.reader.batchSizeRows", "1000000")
.config("spark.rapids.sql.reader.batchSizeBytes",
str(1 * 1024 ** 3)) # 1 GB
# GPU Shuffle через UCX
.config("spark.shuffle.manager",
"com.nvidia.spark.rapids.spark35.RapidsShuffleManager")
# AQE совместимость
.config("spark.sql.adaptive.enabled", "true")
# Диагностика (в production: WARNING)
.config("spark.rapids.sql.explain", "NONE")
.getOrCreate()
)
Проверка корректной инициализации¶
spark = create_rapids_spark("rapids-demo")
# 1. Проверяем, что GPU обнаружен
gpu_info = spark.sparkContext._jvm.com.nvidia.spark.rapids.GpuDeviceManager \
.getGpuResource(spark._jsparkSession)
print(f"GPU resource: {gpu_info}")
# 2. Читаем данные и проверяем план
df = spark.read.parquet("/data/orders_large/")
df.filter("o_totalprice > 50000").groupBy("o_orderstatus").count().explain()
# В плане должны быть строки вида:
# GpuHashAggregateExec(keys=[o_orderstatus], funcs=[count(1)])
# +- GpuFilterExec(o_totalprice > 50000.0)
# +- GpuBatchScanExec[o_orderstatus, o_totalprice] GpuParquetScan
# 3. Простейшая проверка GPU
import subprocess
result = subprocess.run(["nvidia-smi", "--query-gpu=name,memory.used,memory.free",
"--format=csv,noheader"], capture_output=True, text=True)
print(result.stdout)
Лабораторная практика: профилирование и бенчмарк¶
Сценарий: e-commerce Gold-витрина¶
Задача: рассчитать агрегированные метрики продаж по продуктовым категориям и регионам за последние 90 дней. Датасет: 2 миллиарда строк транзакций (Parquet, ~800 GB), 50 миллионов продуктов, 200 тысяч клиентов.
from pyspark.sql import functions as F, SparkSession
spark = create_rapids_spark("ecommerce-gold")
# Чтение партиционированных Parquet таблиц
df_transactions = spark.read.parquet("/data/transactions/")
df_products = spark.read.parquet("/data/products/")
df_customers = spark.read.parquet("/data/customers/")
# Gold-витрина: продажи по категориям и регионам
gold = (
df_transactions
.filter(
F.col("transaction_date") >= F.date_sub(F.current_date(), 90)
)
.join(
df_products.select("product_id", "category_l1", "category_l2", "base_price"),
"product_id"
)
.join(
df_customers.select("customer_id", "region", "customer_segment"),
"customer_id"
)
.groupBy("category_l1", "category_l2", "region", "customer_segment")
.agg(
F.count("*").alias("transaction_count"),
F.sum("amount").alias("total_revenue"),
F.avg("amount").alias("avg_transaction"),
F.sum(F.col("amount") * (1 - F.col("discount_rate"))).alias("net_revenue"),
F.countDistinct("customer_id").alias("unique_customers"),
F.max("amount").alias("max_transaction")
)
.orderBy(F.col("total_revenue").desc())
)
gold.explain()
Анализ плана: GPU-операторы и переходы контекста¶
== Physical Plan ==
GpuColumnarToRowExec
+- GpuSortExec[total_revenue DESC]
+- GpuHashAggregateExec(keys=[category_l1, category_l2, region, customer_segment],
| funcs=[count(1), sum(amount), avg(amount), sum(net), count(distinct cust_id), max(amount)])
+- GpuColumnarExchange hashpartitioning(category_l1,..., 200)
+- GpuHashAggregateExec(keys=[...], funcs=[partial_count, partial_sum, ...])
+- GpuProjectExec[..., amount*(1-discount_rate) AS net]
+- GpuBroadcastHashJoinExec[customer_id]
:- GpuBroadcastHashJoinExec[product_id]
: :- GpuFilterExec(transaction_date >= date_sub(current_date(), 90))
: : +- GpuBatchScanExec[...] GpuParquetScan /data/transactions/
: +- GpuBroadcastExchangeExec HashedRelation
: +- GpuBatchScanExec[product_id, category_l1, ...] GpuParquetScan
+- GpuBroadcastExchangeExec HashedRelation
+- GpuBatchScanExec[customer_id, region, ...] GpuParquetScan
Весь план выполняется на GPU. Обратите внимание:
productsиcustomers- относительно маленькие таблицы (50M и 200K строк), поэтому используетсяGpuBroadcastHashJoinExec: их hash tables помещаются в VRAM каждого Executor целиком.transactions(2B строк) - большая таблица, она сканируется черезGpuBatchScanExecпартиционированно.GpuColumnarExchange- GPU Shuffle между partial и final aggregate.GpuColumnarToRowExec- единственный переход в JVM, и только для финального результата (небольшого числа строк).
Добавление Python UDF: разрыв GPU-пайплайна¶
from pyspark.sql.types import StringType
@udf(returnType=StringType())
def classify_revenue(amount):
"""Плохой паттерн: Python UDF прерывает GPU-выполнение"""
if amount > 1_000_000: return "platinum"
elif amount > 100_000: return "gold"
elif amount > 10_000: return "silver"
else: return "bronze"
gold_with_tier = gold.withColumn("revenue_tier",
classify_revenue(F.col("total_revenue")))
gold_with_tier.explain()
== Physical Plan ==
GpuColumnarToRowExec
+- GpuRowToColumnarExec <-- возврат в GPU
+- BatchEvalPython [classify_revenue] <-- Python UDF на CPU (строковый режим)
+- GpuColumnarToRowExec <-- выход из GPU
+- GpuSortExec[total_revenue DESC]
+- GpuHashAggregateExec(...)
...
Два перехода (GpuColumnarToRowExec + GpuRowToColumnarExec) добавляют двойной PCIe-transfer. Результат агрегации (допустим, 10 000 строк × 6 колонок) небольшой - потери минимальны в данном случае. Но если UDF стоял бы до агрегации, на 2B строках это была бы катастрофа.
Правильное решение:
gold_with_tier = gold.withColumn(
"revenue_tier",
F.when(F.col("total_revenue") > 1_000_000, "platinum")
.when(F.col("total_revenue") > 100_000, "gold")
.when(F.col("total_revenue") > 10_000, "silver")
.otherwise("bronze")
)
# F.when → GpuProjectExec: полностью нативный CUDA-kernel
RAPIDS Qualification Tool: автоматический анализ задач¶
RAPIDS Qualification Tool анализирует Spark History Server event logs и генерирует отчёт: какие задачи выиграют от GPU и насколько.
# Скачиваем Qualification Tool (входит в RAPIDS SDK)
wget https://repo1.maven.org/maven2/com/nvidia/rapids-4-spark-tools_2.12/24.02.0/rapids-4-spark-tools_2.12-24.02.0.jar
# Запускаем анализ логов (event logs со Spark History Server)
java -jar rapids-4-spark-tools_2.12-24.02.0.jar qualification \
--output-directory /tmp/rapids-qual-output \
--event-log /var/spark-history/app-20240115-103022
# Отчёт появится в /tmp/rapids-qual-output/qualification_report.csv
Пример вывода отчёта:
App ID: app-20240115-103022
App Duration: 3847 seconds
Estimated GPU Duration: 412 seconds
Estimated Speedup: 9.3x
GPU Opportunity Score: 87/100
Stage Analysis:
Stage 2 (SortMergeJoin): 3120s → 280s (11.1x) GPU: SUPPORTED
Stage 5 (HashAggregate): 490s → 68s (7.2x) GPU: SUPPORTED
Stage 8 (PythonUDF): 237s → 237s (1.0x) GPU: NOT_SUPPORTED (Python UDF)
Stage 11 (WindowSpec): - → - GPU: PARTIAL (some functions)
Инструмент явно указывает: Stage 8 не переедет на GPU из-за Python UDF, Stage 11 поддерживается частично. Это позволяет приоритизировать рефакторинг.
Симуляция GPU OOM и исправление¶
# Намеренно занижаем batchSizeRows - слишком большой батч для демонстрации
spark_oom = (
SparkSession.builder
.config("spark.rapids.sql.enabled", "true")
.config("spark.rapids.memory.gpu.maxAllocFraction", "0.80")
# Батч 200M строк × средний размер строки 200 байт = 40 GB → > VRAM!
.config("spark.rapids.sql.reader.batchSizeRows", "200000000")
.config("spark.rapids.sql.reader.batchSizeBytes", str(40 * 1024**3))
.getOrCreate()
)
try:
# Большой join: обе таблицы > 1M строк → SortMergeJoin (не broadcast)
result = (
spark_oom.read.parquet("/data/large_table_a/") # 500M строк
.join(spark_oom.read.parquet("/data/large_table_b/"), "key") # 800M строк
.count()
)
except Exception as e:
print(f"GPU OOM: {e}")
# RAPIDS Memory error: Insufficient GPU memory
# RMM pool exhausted: requested 38.2 GB, available 12.1 GB
# Falling back... (если spill включён - вместо crash начнётся spill)
Диагностика через nvidia-smi:
watch -n 0.5 nvidia-smi --query-gpu=name,memory.used,memory.free,utilization.gpu \
--format=csv,noheader
# Name | Mem Used | Mem Free | GPU%
# NVIDIA A100-SXM4-80GB | 72348MiB | 9844MiB | 98% ← норма при полной нагрузке
# NVIDIA A100-SXM4-80GB | 79200MiB | 992MiB | 45% ← почти OOM, GPU ждёт spill
Исправление:
# Правильные параметры для JOIN 500M × 800M строк на A100 80GB
.config("spark.rapids.sql.reader.batchSizeRows", "2000000") # 2M строк
.config("spark.rapids.sql.reader.batchSizeBytes", str(2 * 1024**3)) # 2 GB
# Убеждаемся что spill включён как safety net
.config("spark.rapids.memory.gpu.spillBatchSizeBytes", str(512 * 1024**2))
TPC-H бенчмарк: CPU JVM vs RAPIDS GPU¶
import time
def benchmark_query(spark, query: str, name: str, runs: int = 3) -> float:
times = []
for i in range(runs):
t0 = time.time()
spark.sql(query).count()
times.append(time.time() - t0)
median = sorted(times)[runs // 2]
print(f"{name}: {median:.2f}s (median of {runs})")
return median
# TPC-H Q1: чистая агрегация - идеальный кейс для GPU
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
"""
# TPC-H Q5: 5-way join - показывает силу GPU hash joins
Q5 = """
SELECT n_name, SUM(l_extendedprice * (1 - l_discount)) AS revenue
FROM customer, orders, lineitem, supplier, nation, region
WHERE c_custkey = o_custkey
AND l_orderkey = o_orderkey
AND l_suppkey = s_suppkey
AND c_nationkey = s_nationkey
AND s_nationkey = n_nationkey
AND n_regionkey = r_regionkey
AND r_name = 'ASIA'
AND o_orderdate >= date('1994-01-01')
AND o_orderdate < date('1995-01-01')
GROUP BY n_name
ORDER BY revenue DESC
"""
# Создаём сессии для сравнения
spark_cpu = SparkSession.builder.appName("baseline").getOrCreate()
spark_gpu = create_rapids_spark("rapids-bench")
data_path = "/data/tpch/sf100"
for table in ["customer", "orders", "lineitem", "supplier", "nation", "region"]:
spark_cpu.read.parquet(f"{data_path}/{table}").createOrReplaceTempView(table)
spark_gpu.read.parquet(f"{data_path}/{table}").createOrReplaceTempView(table)
print("=== TPC-H SF=100 Benchmark ===")
print(f"\n{'Query':<8} {'CPU JVM':>12} {'RAPIDS GPU':>12} {'Speedup':>10}")
print("-" * 45)
for qid, query in [("Q1", Q1), ("Q5", Q5)]:
cpu_t = benchmark_query(spark_cpu, query, f"CPU {qid}")
gpu_t = benchmark_query(spark_gpu, query, f"GPU {qid}")
print(f"{qid:<8} {cpu_t:>10.2f}s {gpu_t:>11.2f}s {cpu_t/gpu_t:>9.2f}x")
Типичные результаты на H100 80GB vs 32-core CPU кластер с 4 Executor-ами:
| Запрос | CPU JVM (4×32c) | RAPIDS GPU (4×H100) | Ускорение |
|---|---|---|---|
| Q1 (pure agg) | 150s | 14s | 10.7× |
| Q5 (5-way join) | 320s | 32s | 10.0× |
| Q6 (filter+sum) | 60s | 5s | 12.0× |
| Q18 (subquery+join) | 480s | 58s | 8.3× |
Экономика GPU: TCO vs CPU кластер¶
GPU-кластер дороже CPU-кластера в абсолютных цифрах, но дешевле в расчёте на единицу вычислений. Это ключевое понимание для обоснования GPU-инфраструктуры.
Пример расчёта для AWS:
| Конфигурация | Инстанс | В час | Ускорение Q1 | Стоимость/запрос |
|---|---|---|---|---|
| CPU: 10 × r7i.8xlarge | 32 vCPU, 256 GB RAM | $2.52 × 10 = $25.20 | 1× (150s) | $1.05 |
| GPU: 4 × p4d.24xlarge | 8 A100, 320 GB RAM | $32.77 × 4 = $131 | 10× (14s) | $0.51 |
Несмотря на то что GPU-кластер стоит в 5 раз дороже в час, каждый запрос обходится вдвое дешевле - потому что GPU выполняет его в 10 раз быстрее.
Критическое условие: GPU-кластер должен быть утилизирован на 70%+. GPU простаивающие 22 часа из 24 - это просто очень дорогая тумба. GPU эффективны для высокозагруженных аналитических кластеров с постоянным потоком тяжёлых запросов.
Best Practices и чек-лист готовности¶
Когда включать RAPIDS: критерии¶
Используйте RAPIDS Qualification Tool перед миграцией. Если оценочный speedup > 3× - GPU выгоден. Если < 2× - вероятно, задача IO-bound или Python-heavy, GPU не поможет.
Признаки хорошего GPU-кандидата:
- Запросы занимают >5 минут на CPU (overhead GPU < 5% от выигрыша)
- Датасеты > 10 GB на задачу
- Доминируют агрегации, join, фильтрация числовых колонок
- Нет Python UDF в критическом пути
- Данные в Parquet/ORC (колоночный формат)
Признаки плохого GPU-кандидата:
- Много маленьких файлов (< 10 MB каждый): GPU простаивает в ожидании IO
- Python-heavy ETL с построчными трансформациями
- Сложные оконные функции с частичной поддержкой
COLLECT_LIST,COLLECT_SETна больших группах (нет в cuDF)- Строковые регулярные выражения (частичная поддержка cuDF)
Анти-паттерны¶
Анти-паттерн 1: Python UDF в начале тяжёлого пайплайна
# ПЛОХО: UDF до JOIN разрывает GPU-путь на самом нагруженном месте
df_flagged = df.withColumn("flag", my_python_udf(F.col("amount"))) # 2B строк!
df_flagged.join(df_other, "key").groupBy("category").agg(F.sum("amount"))
# UDF → GpuColumnarToRow (2B × 200 байт = 400 GB через PCIe)!
# ХОРОШО: UDF вынесен после агрегации (результат < 1000 строк)
df.join(df_other, "key").groupBy("category").agg(F.sum("amount")) \
.withColumn("flag", my_python_udf(F.col("total")))
# UDF работает на 1000 строках - overhead ничтожен
Анти-паттерн 2: Игнорировать GPU utilization
# Добавьте в cron или monitoring:
nvidia-smi --query-gpu=utilization.gpu,memory.used --format=csv,noheader
# Если GPU utilization < 30% - ищите причину: fallback, IO-ожидание, skew
Анти-паттерн 3: Одинаковые параметры для разных GPU
T4 (16 GB) и A100 (80 GB) требуют совершенно разных batchSizeRows и RMM.maxAllocFraction. Параметры для A100 убьют T4 OOM, параметры для T4 недозагрузят A100.
Домашнее задание¶
Условие:
Вам дан лог Spark-приложения на GPU-кластере (4 × T4 16GB) и SQL-запрос. Из 20 стейджей пайплайна только 3 выполнились на GPU - остальные упали в fallback на CPU из-за неподдерживаемых функций. GPU-ноды простаивали.
# Проблемный скрипт - найдите источники fallback
from pyspark.sql import functions as F
from pyspark.sql.types import StringType, ArrayType
spark = create_rapids_spark("homework-job")
# UDF 1: регулярные выражения со сложным паттерном
@udf(returnType=StringType())
def extract_sku(product_name):
import re
match = re.search(r'SKU-(\d{4,8})', product_name)
return match.group(1) if match else None
# UDF 2: агрегация в Python-объект
@udf(returnType=ArrayType(StringType()))
def top_categories(categories_str):
cats = categories_str.split(",")
return sorted(set(cats))[:3]
df = (
spark.read.parquet("/data/ecommerce_large/")
# Сложный REGEXP → fallback
.withColumn("sku", extract_sku(F.col("product_name")))
# Агрегирующий UDF → fallback
.withColumn("top_cats", top_categories(F.col("category_path")))
# Locale-зависимая функция → fallback
.withColumn("formatted_price",
F.format_number(F.col("price"), 2))
.groupBy("region", "sku")
.agg(F.sum("revenue").alias("total"))
)
df.explain()
df.write.mode("overwrite").parquet("/data/output/")
Задание:
-
Анализ fallback: запустите с
spark.rapids.sql.explain=ALLи изучите explain()-вывод. Для каждого из трёх UDF и функцииformat_numberопределите: почему RAPIDS не может исполнить их на GPU? -
Рефакторинг:
extract_sku→ замените наF.regexp_extract(F.col("product_name"), r'SKU-(\d{4,8})', 1). Проверьте: поддерживает ли cuDF этот regexp-паттерн?top_categories→ замените на комбинациюF.split(),F.array_distinct(),F.slice(),F.sort_array()-
format_number→ замените наF.round(F.col("price"), 2).cast("string") -
Верификация: убедитесь, что результаты скрипта до и после рефакторинга идентичны (используйте
.subtract()для проверки). -
Профилирование: запустите RAPIDS Qualification Tool на event logs до и после рефакторинга. Сравните
Estimated SpeedupиGPU Opportunity Score. -
Отчёт: приложите:
- Таблицу: какой оператор заменили, на что, и почему оригинал не поддерживался GPU
- Вывод
explain()до и после (только нерелевантные части можно сократить) - Скриншот
nvidia-smiво время выполнения исправленного скрипта (GPU utilization > 80%) - Сравнение времени выполнения: оригинал vs рефакторинг