RAPIDS cuDF: GPU acceleration, CUDA требования и профиль данных

RAPIDS cuDF: GPU acceleration, CUDA требования и профиль данных

optimization

Архитектурный сдвиг: когда 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-аналогами: FilterExecGpuFilterExec, HashAggregateExecGpuHashAggregateExec, SortMergeJoinExecGpuSortMergeJoinExec. Данные при этом перемещаются из 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. Он работает следующим образом:

  1. CPU читает метаданные: footer Parquet-файла (схема, row group statistics, column chunk offsets) читается CPU в host-память. Это лёгкая операция.
  2. CPU передаёт задачу GPU: список compressed column chunks (байтовые диапазоны файла) отправляется в cuDF.
  3. GPU декомпрессирует параллельно: каждый column chunk (Snappy/Zstd-блок) декомпрессируется отдельным CUDA-потоком. На GPU одновременно могут работать тысячи декомпрессоров.
  4. GPU декодирует encoding: RLE (Run-Length Encoding), Delta encoding, Dictionary encoding - всё декодируется параллельными CUDA-kernels.
  5. Результат: готовая 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/")

Задание:

  1. Анализ fallback: запустите с spark.rapids.sql.explain=ALL и изучите explain()-вывод. Для каждого из трёх UDF и функции format_number определите: почему RAPIDS не может исполнить их на GPU?

  2. Рефакторинг:

  3. extract_sku → замените на F.regexp_extract(F.col("product_name"), r'SKU-(\d{4,8})', 1). Проверьте: поддерживает ли cuDF этот regexp-паттерн?
  4. top_categories → замените на комбинацию F.split(), F.array_distinct(), F.slice(), F.sort_array()
  5. format_number → замените на F.round(F.col("price"), 2).cast("string")

  6. Верификация: убедитесь, что результаты скрипта до и после рефакторинга идентичны (используйте .subtract() для проверки).

  7. Профилирование: запустите RAPIDS Qualification Tool на event logs до и после рефакторинга. Сравните Estimated Speedup и GPU Opportunity Score.

  8. Отчёт: приложите:

  9. Таблицу: какой оператор заменили, на что, и почему оригинал не поддерживался GPU
  10. Вывод explain() до и после (только нерелевантные части можно сократить)
  11. Скриншот nvidia-smi во время выполнения исправленного скрипта (GPU utilization > 80%)
  12. Сравнение времени выполнения: оригинал vs рефакторинг