RAPIDS: когда GPU не даёт прироста - малые данные, IO-bound запросы

RAPIDS: когда GPU не даёт прироста - малые данные, IO-bound запросы

core

Парадокс GPU в Big Data: почему железо для ИИ иногда проигрывает CPU

Предыдущий урок демонстрировал впечатляющие цифры: RAPIDS на H100 даёт до 10–12× ускорения на TPC-H-запросах. Это правда. Но TPC-H - это специально подобранный набор сложных аналитических запросов, который выглядит как идеальный кейс для GPU: огромные таблицы, тяжёлые join, много агрегаций, никаких Python UDF, данные заранее в Parquet. В реальном мире Data Engineering пайплайн выглядит иначе.

Этот урок - намеренно другой. Мы разберём сценарии, где GPU не только не помогает, но и активно мешает: замедляет задачу, увеличивает стоимость инфраструктуры и создаёт операционную сложность без измеримой отдачи. Понимать эти сценарии - не менее важно, чем уметь запустить RAPIDS на кластере.

Главный миф индустрии

«GPU всегда быстрее CPU» - это маркетинговое утверждение, верное в контексте Deep Learning, где задача состоит из миллиардов умножений матриц с идеально регулярной структурой. В этой задаче GPU действительно убивает CPU: NVIDIA H100 выполняет матричное умножение с производительностью до 3958 TFLOPS (BF16), тогда как лучшие CPU дают 10–20 TFLOPS в той же операции - разница в 200–400 раз.

Но Data Engineering - это не матричное умножение. Это:

  • Чтение данных из S3/HDFS по сети (IO-bound)
  • Десериализация JSON/CSV/Parquet (CPU-bound, но несложно параллелизуется)
  • Множество маленьких трансформаций (узкие преобразования, не требующие параллелизма)
  • Python UDF и бизнес-логика (строго последовательно, CPU-only)
  • Shuffle между Executor-ами (сеть - IO-bound)
  • Агрегации (да, GPU помогает - но это часто лишь 20–30% пайплайна)

Закон Амдала: почему 30% выигрыша в одном месте не спасают всё

Закон Амдала формулирует ограничение на максимальное теоретическое ускорение системы при ускорении одной её части:

Speedup = 1 / ((1 - p) + p/s)

где:
  p = доля задачи, которую можно ускорить (0.0–1.0)
  s = коэффициент ускорения этой части

Пример: пайплайн состоит из:

  • Чтение из S3 по сети: 40% времени (IO-bound, GPU не поможет)
  • JSON-парсинг: 20% (GPU частично помогает)
  • Python UDF для бизнес-логики: 25% (GPU не поможет вообще)
  • GROUP BY агрегация: 15% (GPU даёт 10× ускорение)

Считаем: только 15% задачи реально ускоряется на GPU в 10 раз. Подставляем:

p = 0.15  (агрегация)
s = 10    (GPU ускорение)

Speedup = 1 / ((1 - 0.15) + 0.15/10)
        = 1 / (0.85 + 0.015)
        = 1 / 0.865
        ≈ 1.16

Итог: GPU ускоряет весь пайплайн лишь на 16%. При этом GPU-инстанс стоит в 3–5 раз дороже CPU-инстанса. Результат - затраты выросли в 4 раза при росте производительности на 16%. Это провал.

Именно поэтому перед внедрением GPU критически важно измерить: какой процент времени пайплайна реально можно ускорить. Если ответ < 50% - GPU, скорее всего, не окупится.


Ловушка №1: малые данные и оверхед на инициализацию

Что такое «малые данные» в контексте GPU

GPU - это устройство с высокой пропускной способностью (throughput), но высокой задержкой (latency) при старте работы. CPU - наоборот: умеренная пропускная способность, но низкая латентность для первой операции.

Говоря о «малых данных», мы имеем в виду не абсолютный размер, а размер относительно overhead-а GPU. Эмпирическое правило:

  • Датасет < 1 GB на Executor: PCIe transfer overhead + kernel launch overhead > время вычислений на GPU
  • Датасет < 10 GB на Executor: выигрыш есть, но скромный (< 2×)
  • Датасет > 50 GB на Executor при вычислительно-тяжёлых запросах: GPU начинает заметно выигрывать

Холодный старт CUDA: что происходит при первом запуске

Когда Executor запускается на узле с GPU, RAPIDS должен инициализировать нативную инфраструктуру. Это происходит один раз при старте, но занимает ощутимое время:

1. Загрузка нативной библиотеки cuDF (libcudf.so): ~0.5–1s
2. Инициализация CUDA-контекста: ~0.5–2s
   (CUDA context содержит JIT-скомпилированные kernels, 
    таблицы символов, менеджер памяти)
3. Аллокация RMM Memory Pool (резервирование 80% VRAM): ~0.2–0.5s
   (cudaMalloc для большого блока медленнее, чем маленькие аллокации)
4. Прогрев (warm-up): первый CUDA kernel выполняется медленно —
   нужно JIT-скомпилировать PTX → SASS под конкретный GPU: ~1–3s

Суммарный cold-start: 2–6 секунд на каждый Executor. Для задачи, которая сама по себе выполняется 10 секунд - это 20–60% дополнительного overhead. Для задачи длиной в 2 часа - менее 0.1%.

Вывод: короткоживущие задачи (< 1–2 минут) теряют значительную долю потенциального выигрыша на cold-start.

Физика PCIe: почему передача данных не бесплатна

Рассмотрим конкретный расчёт для малого датасета:

# Сценарий: фильтрация датасета 500 MB на T4 GPU (16 GB VRAM)
# T4 в Google Cloud: ~$0.35/час

dataset_size_mb = 500
pcie_bandwidth_gbs = 16   # PCIe 3.0 x16 для T4 в облачной VM

# Время передачи CPU RAM → GPU VRAM:
transfer_time_s = (dataset_size_mb / 1024) / pcie_bandwidth_gbs
print(f"PCIe transfer time: {transfer_time_s * 1000:.1f} ms")
# PCIe transfer time: 30.5 ms

# Время вычислений на GPU (фильтрация, простой предикат):
# T4 имеет пропускную память ~320 GB/s
# 500 MB / 320 GB/s = 1.56 ms
gpu_compute_time_ms = 1.56
print(f"GPU compute time: {gpu_compute_time_ms:.1f} ms")

# Аналогичный расчёт для CPU (AVX-512 SIMD):
# CPU RAM bandwidth ~50 GB/s, фильтрация за ~10 ms
cpu_compute_time_ms = 10

# На GPU: 30.5 (туда) + 1.56 (вычисление) + 30.5 (обратно) = 62.6 ms
gpu_total_ms = transfer_time_s * 1000 * 2 + gpu_compute_time_ms
print(f"Total GPU path: {gpu_total_ms:.1f} ms")

# На CPU: 10 ms (нет передачи, данные уже в RAM)
print(f"Total CPU path: {cpu_compute_time_ms:.1f} ms")

# GPU в 6.26 раза МЕДЛЕННЕЕ для этого сценария!
print(f"GPU/CPU ratio: {gpu_total_ms / cpu_compute_time_ms:.2f}x slower on GPU")
PCIe transfer time: 30.5 ms
GPU compute time: 1.6 ms
Total GPU path: 62.6 ms
Total CPU path: 10.0 ms
GPU/CPU ratio: 6.26x slower on GPU

Эта математика объясняет: для маленьких партиций GPU медленнее, потому что 97% времени тратится на перекладывание данных через PCIe, а вычисления занимают 3%. CPU за те же 62.6 мс успел бы обработать не 500 MB, а все 3 GB.

Проблема мелких партиций и малого параллелизма

GPU эффективен, когда у него много работы одновременно: тысячи CUDA-потоков должны быть заняты. Если батч содержит 10000 строк, а GPU имеет 16896 CUDA-ядер - более 6000 ядер простаивают. Это называется низкая occupancy (загруженность GPU).

Идеальный размер батча для полной загрузки GPU:

  • Минимально: 100000–500000 строк (зависит от ширины строки)
  • Оптимально: 1–10 миллионов строк на батч
  • Недостаточно: < 10000 строк - большинство CUDA-ядер простаивает

Когда Spark создаёт слишком много маленьких партиций (например, из-за AQE с маленькими файлами или настройки spark.sql.shuffle.partitions=200 при маленьких данных), каждая партиция обрабатывается отдельным батчем - и GPU никогда не загружается полностью.


Ловушка №2: IO-bound запросы

Compute-bound vs IO-bound: фундаментальное различие

Любой вычислительный пайплайн можно классифицировать по типу узкого места:

Compute-bound: CPU/GPU ядра работают на полную мощность, не успевая обрабатывать поступающие данные. Ускорение процессора непосредственно ускоряет задачу. Типичные примеры: матричное умножение, хэш-агрегация очень широких таблиц, JOIN с большой cardinality.

IO-bound: CPU/GPU ядра простаивают, ожидая данных из сети или с диска. Ускорение процессора не меняет ничего - данные всё равно приходят с той же скоростью. Типичные примеры: чтение из S3 с throttling, работа с тысячами мелких файлов, запросы к медленным JDBC-источникам.

Ключевой инструмент диагностики: время задачи (Task) в Spark UI. Если в метриках задачи доля Fetch Wait Time или HDFS Read Time > 50% - задача IO-bound, и GPU не поможет.

Сценарий: чтение неоптимизированных данных

from pyspark.sql import SparkSession, functions as F
import time

spark = SparkSession.builder.appName("io-bound-demo").getOrCreate()

# ПЛОХОЙ сценарий 1: терабайт несжатых JSON
# JSON - строковый формат: нет предиката pushdown, нет column pruning,
# нет нативного GPU-декодера (cuDF плохо работает с JSON)
t0 = time.time()
df_json = (
    spark.read.json("/data/events/2024/01/*/")   # тысячи мелких файлов
    .filter(F.col("event_type") == "purchase")
    .select("user_id", "amount", "product_id")
    .groupBy("product_id")
    .agg(F.sum("amount").alias("revenue"))
)
df_json.count()
print(f"JSON pipeline: {time.time()-t0:.1f}s")
# Типичный результат:
# Stage 1 (Read + Filter): 4800s total / 200 tasks = 24s/task average
# Из которых: 21s - чтение JSON (IO), 3s - парсинг (CPU), 0.3s - фильтр (GPU)
# GPU утилизация: ~1.2% от времени task

# ПЛОХОЙ сценарий 2: медленный JDBC с ограниченным параллелизмом
df_jdbc = (
    spark.read
    .format("jdbc")
    .option("url", "jdbc:postgresql://prod-db:5432/orders")
    .option("dbtable", "large_orders_table")
    .option("numPartitions", "20")      # максимум 20 параллельных соединений
    .option("fetchsize", "1000")        # PostgreSQL fetchsize
    .load()
    .groupBy("region").agg(F.sum("amount"))
)
# JDBC читает данные через сетевой стрим с PostgreSQL
# Скорость ограничена JDBC connection и сетью: ~100 MB/s max
# При таблице 200 GB: 200 GB / 100 MB/s = 2000 секунд = 33 минуты только на чтение
# GPU будет простаивать всё это время

В этих сценариях nvidia-smi покажет:

watch -n 1 nvidia-smi --query-gpu=utilization.gpu,memory.used --format=csv,noheader
# utilization.gpu [%], memory.used [MiB]
# 2 %, 1024 MiB      ← GPU почти не используется
# 1 %, 1024 MiB      ← данные ещё читаются из сети
# 3 %, 1024 MiB      ← GPU обработал крошечный батч
# 1 %, 1024 MiB      ← снова ждём IO

GPU утилизован на 1–3% при цене GPU-инстанса в 5–10 раз дороже CPU-инстанса.

Сценарий: S3 throttling

AWS S3 (и любое облачное объектное хранилище) имеет ограничения на скорость:

  • S3 Standard: 5500 GET-запросов/секунду на префикс, ~100 GB/s суммарный throughput
  • Но у каждого конкретного bucket и prefix есть soft limits
  • При высокой нагрузке S3 отвечает 503 SlowDown - Spark делает retry с backoff

При чтении 10000 мелких файлов по 10 MB (= 100 GB суммарно):

# Каждый файл = один S3 GetObject запрос
# 10000 запросов / 5500 req/s = ~1.8 секунды только на открытие файлов
# Плюс латентность каждого запроса: 20–100 ms (time-to-first-byte)
# 10000 файлов × 50 ms TTFB = 500 секунд дополнительного ожидания!

# RAPIDS может ускорить фильтрацию в 10 раз:
# Фильтрация 100 GB: CPU 300s → GPU 30s
# Но IO overhead: 500s не изменится от GPU

# Итог с GPU:  500 (IO) + 30 (compute) = 530s
# Итог с CPU:  500 (IO) + 300 (compute) = 800s
# Ускорение: 800/530 = 1.51x
# Стоимость GPU: в 5x дороже
# TCO: GPU дороже в 5x/1.51x ≈ в 3.3x дороже за единицу вычислений!

Решение проблемы мелких файлов - не GPU, а компакция данных: объединить 10000 файлов по 10 MB в 1000 файлов по 100 MB через Spark compaction job. После этого S3-запросов станет в 10 раз меньше, и CPU-кластер ускорится в 2–3 раза - без дорогостоящего GPU.

Правило: исправляйте архитектуру, а не железо

IO-bound проблемы решаются оптимизацией данных и архитектуры, а не добавлением GPU:

IO-проблема Решение Инструмент
Мелкие файлы Компакция Spark, Delta OPTIMIZE
Несжатый JSON/CSV Конвертация в Parquet Spark ETL
Медленный JDBC Репликация в колоночный формат Spark + Parquet
S3 throttling Оптимизация партиционирования Delta Z-order, Iceberg
Горячие shuffle-партиции AQE skew join spark.sql.adaptive

После устранения IO-проблем CPU-кластер может ускориться на 5–20×. Только после этого имеет смысл оценивать, даст ли GPU дополнительное ускорение.


Ловушка №3: постоянный пинг-понг между средами

Самая коварная ловушка GPU-ускорения - это частичный fallback: когда часть операторов работает на GPU, а часть возвращается в JVM. Каждый переход между GPU и CPU - это дорогая операция копирования данных через PCIe.

Анатомия разрушительного пинг-понга

На диаграмме видно: данные объёмом 1 TB проходят через PCIe дважды (туда и обратно). При пропускной способности PCIe 30 GB/s:

Transfer #1 (GPU → CPU): 1000 GB / 30 GB/s = 33.3 секунды
Python UDF выполнение:                        60 секунд
Transfer #2 (CPU → GPU): 1000 GB / 30 GB/s = 33.3 секунды
GpuHashAggregateExec:                          8 секунд
----
Total с GPU:                                 134.6 секунды

Для сравнения - чистый CPU (без GPU):
CPU HashAgg:                                  80 секунд
----
Total без GPU:                                80 секунд

GPU МЕДЛЕННЕЕ в 1.68 раза при цене в 5x дороже!

Типичные источники пинг-понга

Python UDF - самый частый виновник:

from pyspark.sql.types import StringType, DoubleType

# Любой Python UDF → немедленный GpuColumnarToRow
@udf(returnType=StringType())
def classify_customer(spend, region):
    # Сложная бизнес-логика на Python
    tiers = {"EMEA": [1000, 10000], "APAC": [500, 5000]}
    thresholds = tiers.get(region, [2000, 20000])
    if spend >= thresholds[1]: return "platinum"
    if spend >= thresholds[0]: return "gold"
    return "standard"

# Этот код разрывает GPU-план в точке вызова UDF
df.withColumn("tier", classify_customer(F.col("spend"), F.col("region")))

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

# TimestampNTZ (без timezone) - не поддерживается в старых версиях RAPIDS
df.withColumn("ts_str",
    F.col("created_at_ntz").cast("string"))  # TimestampNTZ → String: fallback!

# Вложенный Map → неполная поддержка
df.withColumn("keys", F.map_keys(F.col("metadata")))  # MapType: fallback

# Решение: конвертировать во время подготовки данных
df_clean = df.withColumn(
    "created_at",
    F.col("created_at_ntz").cast("timestamp")  # → обычный Timestamp: поддерживается
)

Сложные оконные функции:

from pyspark.sql.window import Window

w = Window.partitionBy("customer_id").orderBy("transaction_date")

# Эти функции поддерживаются RAPIDS:
df.withColumn("row_num", F.row_number().over(w))  # GpuWindowExec
df.withColumn("running_sum", F.sum("amount").over(w))  # GpuWindowExec

# Эти - НЕТ (fallback):
df.withColumn("lag_2", F.lag("amount", 2).over(w))  # Lag с offset > 1: fallback
df.withColumn("pct", F.percent_rank().over(w))        # percent_rank: fallback

COLLECT_LIST и COLLECT_SET - не реализованы в cuDF:

# Этот план разрывает GPU-выполнение на финальной агрегации
df.groupBy("user_id").agg(
    F.collect_list("product_id").alias("products"),  # → CPU
    F.sum("amount").alias("total")                   # → GPU (но вынужден CPU из-за collect_list)
)

Как обнаружить пинг-понг

# Включаем детальный вывод оператора offloading
spark.conf.set("spark.rapids.sql.explain", "ALL")

# Запускаем запрос с UDF
df_result = (
    df.withColumn("tier", classify_customer(F.col("spend"), F.col("region")))
      .groupBy("region", "tier")
      .agg(F.sum("spend").alias("total"))
)

df_result.explain()

В выводе explain() ищите чередование GPU и CPU операторов:

GpuHashAggregateExec(...)            ← GPU (финальный)
+- GpuColumnarExchange               ← GPU shuffle
   +- GpuHashAggregateExec(partial)  ← GPU
      +- GpuRowToColumnarExec        ← ПЕРЕХОД CPU→GPU (красный флаг!)
         +- BatchEvalPython          ← CPU Python UDF
            +- GpuColumnarToRowExec  ← ПЕРЕХОД GPU→CPU (красный флаг!)
               +- GpuProjectExec     ← GPU
                  +- GpuBatchScanExec ← GPU

Два красных флага (GpuColumnarToRowExec + GpuRowToColumnarExec) - это один пинг-понг цикл. Если их несколько в плане - суммарный PCIe-overhead может быть катастрофическим.


Лабораторная практика: анти-бенчмарк

Сценарий: провальная GPU-оптимизация

Разберём реальный сценарий: команда перевела пайплайн обработки событий на GPU-кластер (4 × T4 GPU), рассчитывая на 5× ускорение. По факту задача стала работать на 15% медленнее, а счёт за инфраструктуру вырос в 4 раза.

from pyspark.sql import SparkSession, functions as F
from pyspark.sql.types import DoubleType, StringType, MapType
import time, json

# Конфигурация GPU-кластера (как её настроила команда)
spark = (
    SparkSession.builder
    .appName("analytics-gpu-wrong")
    .config("spark.plugins", "com.nvidia.spark.SQLPlugin")
    .config("spark.rapids.sql.enabled", "true")
    .config("spark.executor.memory", "16g")
    .config("spark.executor.resource.gpu.amount", "1")
    .getOrCreate()
)

# --- Источник данных: 200 GB сырых событий в JSON ---
# Проблема 1: JSON-файлы по 1–5 MB каждый (50000 файлов)
# Проблема 2: JSON не поддерживает predicate pushdown
# Проблема 3: cuDF JSON reader имеет ограничения
df_raw = spark.read.json("/data/raw-events/2024/01/*/events_*.json")

# Проблема 4: сложный Python UDF для парсинга бизнес-логики
@udf(returnType=MapType(StringType(), DoubleType()))
def parse_pricing_rules(rules_json, product_type):
    """Сложная логика ценообразования - Python only"""
    try:
        rules = json.loads(rules_json)
        multipliers = {"premium": 1.3, "standard": 1.0, "economy": 0.8}
        mult = multipliers.get(product_type, 1.0)
        return {k: v * mult for k, v in rules.items()}
    except:
        return {}

# Проблема 5: TimestampNTZ → RAPIDS старых версий не поддерживает
# Проблема 6: MapType output от UDF → нет GPU оператора

df_processed = (
    df_raw
    .filter(F.col("event_type").isin(["purchase", "refund", "adjustment"]))
    # ↑ GpuFilterExec: работает на GPU
    .withColumn("pricing", parse_pricing_rules(
        F.col("pricing_rules_json"), F.col("product_type")
    ))
    # ↑ Python UDF: GpuColumnarToRow → BatchEvalPython → (нет возврата на GPU!)
    # MapType output → нет GPU-оператора, весь downstream в CPU
    .withColumn("base_price", F.col("pricing").getItem("base"))
    # ↑ MapType операция: CPU
    .withColumn("event_ts",
        F.to_timestamp(F.col("event_timestamp_str"), "yyyy-MM-dd'T'HH:mm:ssXXX")
    )
    # ↑ Строковый timestamp с timezone: CPU
    .groupBy(
        F.window("event_ts", "1 hour"),
        F.col("product_type"),
        F.col("region")
    )
    .agg(
        F.sum("base_price").alias("hourly_revenue"),
        F.count("*").alias("event_count")
    )
    # ↑ GpuHashAggregateExec: вернулся на GPU для финала
    # но данные уже прошли через CPU, обратный transfer обязателен
)

t0 = time.time()
df_processed.write.mode("overwrite").parquet("/data/output/hourly_kpis/")
print(f"GPU pipeline duration: {time.time()-t0:.1f}s")

Анализ провала: чтение Spark UI

После запуска открываем Spark UI и смотрим на метрики стадий:

Stage 1: Read JSON + Filter
  Duration: 4800s (80 minutes!)
  Task Metrics:
    Input size: 200 GB
    Task Deserialization Time: avg 8.2s/task ← JSON parsing overhead
    HDFS Read Time: avg 18.7s/task            ← IO bottleneck
    GPU Compute Time: avg 0.4s/task           ← GPU занят 2% времени!
  GPU Utilization: 1.8% average

Stage 2: Python UDF
  Duration: 1850s (31 minutes)
  Task Metrics:
    Python Worker Time: avg 9.3s/task
    GPU Columnar To Row: avg 3.2s/task        ← PCIe transfer!
  GPU Utilization: 0% (GPU не используется в этой стадии)

Stage 3: Shuffle + Final Aggregate
  Duration: 420s (7 minutes)
  Task Metrics:
    Row To GPU Columnar: avg 2.1s/task        ← PCIe transfer обратно!
    GPU Aggregate Time: avg 0.8s/task         ← GPU занят 28% в этой стадии
  GPU Utilization: 28% (самая лучшая стадия, но она занимает < 6% общего времени)

TOTAL: 7070s (117 minutes)

Итог: GPU утилизирован менее 5% суммарного времени. 95% времени задача либо ждёт IO, либо выполняется на CPU с накладными расходами на PCIe-transfer.

Использование RAPIDS Qualification Tool

# Скачиваем утилиту анализа
java -jar rapids-4-spark-tools_2.12-24.02.0.jar qualification \
  --output-directory /tmp/qual-report \
  --event-log /var/spark-history/app-analytics-gpu-wrong

cat /tmp/qual-report/qualification_report.csv

Вывод RAPIDS Qualification Tool:

App ID: app-analytics-gpu-wrong
App Name: analytics-gpu-wrong
App Duration: 7070 seconds
Estimated GPU Duration: 6680 seconds      ← GPU чуть медленнее!
Estimated Speedup: 1.06x                  ← почти без изменений
GPU Opportunity Score: 12/100             ← критически низкий!

Top Reasons for Low Score:
  1. Python UDF (parse_pricing_rules): 1850s unsupported
  2. JSON read: 4800s - no GPU acceleration benefit
  3. MapType operations: not supported in GPU execution
  4. TimestampNTZ cast: fallback to CPU

Recommendation: This application is NOT a good candidate for GPU acceleration.
  Consider CPU-based optimization first:
  - Convert JSON to Parquet format
  - Replace Python UDF with built-in SQL functions
  - Address IO bottleneck (small files consolidation)

Инструмент прямо говорит: GPU-ускорение для этой задачи неэффективно. Оценочный speedup - 1.06× (практически ничего), при том что GPU-инстанс стоит в 4–5× дороже.

Экономический расчёт TCO

# TCO (Total Cost of Ownership) расчёт

# CPU кластер (например, AWS r6i.4xlarge: 16 vCPU, 128 GB RAM)
cpu_instances = 10
cpu_price_per_hour = 1.01  # USD за r6i.4xlarge
cpu_duration_hours = 7070 / 3600  # 1.96 hours

# Оптимизированный CPU (после исправления проблем данных)
cpu_optimized_duration_hours = 1800 / 3600  # 30 min (после конвертации JSON→Parquet)
cpu_cost = cpu_instances * cpu_price_per_hour * cpu_duration_hours
cpu_optimized_cost = cpu_instances * cpu_price_per_hour * cpu_optimized_duration_hours

# GPU кластер (p3.2xlarge: 1 × V100, 8 vCPU, 61 GB RAM)
gpu_instances = 4
gpu_price_per_hour = 3.06  # USD за p3.2xlarge
gpu_duration_hours = 7070 / 3600  # такое же время - GPU не помог!
gpu_cost = gpu_instances * gpu_price_per_hour * gpu_duration_hours

print("=== TCO Analysis ===")
print(f"CPU cluster (неоптимизированный): ${cpu_cost:.2f}")
print(f"CPU cluster (оптимизированный):   ${cpu_optimized_cost:.2f}")
print(f"GPU cluster (неоптимизированный): ${gpu_cost:.2f}")
print(f"")
print(f"GPU vs CPU (сырой): в {gpu_cost/cpu_cost:.1f}x дороже")
print(f"GPU vs CPU (оптим): в {gpu_cost/cpu_optimized_cost:.1f}x дороже")
=== TCO Analysis ===
CPU cluster (неоптимизированный): $19.82
CPU cluster (оптимизированный):   $5.05
GPU cluster (неоптимизированный): $23.38

GPU vs CPU (сырой): в 1.2x дороже (при том же времени!)
GPU vs CPU (оптим): в 4.6x дороже

Правильная последовательность действий:

  1. Конвертировать JSON → Parquet (бесплатно, разовая задача)
  2. Заменить Python UDF на F.when / встроенные функции (разработка, 2–4 часа)
  3. Укрупнить партиции / решить проблему мелких файлов
  4. После этого оценить, даст ли GPU дополнительный выигрыш на оставшемся вычислительном узком месте

Матрица принятия решений: GPU или CPU

Чек-лист перед внедрением GPU

Ответьте на эти вопросы последовательно. Если хотя бы один ответ - «нет» на первых трёх - GPU не окупится.

Вопрос 1: Объём данных

  • Батч данных на Executor > 10 GB? → Нет → CPU, GPU не окупится из-за overhead
  • Датасет суммарно > 100 GB? → Нет → CPU, нет смысла в GPU-кластере
  • Задача выполняется > 10 минут на CPU? → Нет → CPU, cold-start GPU слишком дорог

Вопрос 2: Профиль запроса

  • В плане > 70% операторов поддерживаются RAPIDS? → Нет → CPU
  • Нет Python UDF в критическом пути? → Нет → CPU или рефакторинг сначала
  • Нет MapType / сложных nested структур? → Нет → CPU или упрощение схемы

Вопрос 3: IO-характеристики

  • Данные в Parquet/ORC (не JSON/CSV)? → Нет → конвертация сначала
  • Размер файлов > 64 MB (нет проблемы мелких файлов)? → Нет → компакция сначала
  • Задача не IO-bound (GPU utilization > 30% в RAPIDS Profiling Tool)? → Нет → CPU

Если все «да»: оцените RAPIDS Qualification Score

  • Qualification Score > 70 → GPU даст значимое ускорение (> 3×)
  • Qualification Score 50–70 → GPU даст умеренное ускорение (2–3×), считайте TCO
  • Qualification Score < 50 → GPU сомнительно выгоден, оптимизируйте CPU первым

Таблица: типичные Spark-задачи и рекомендации

Задача GPU? Причина
TPC-H/TPC-DS аналитика на Parquet > 100 GB Да Compute-bound, полное покрытие RAPIDS
ETL JSON → Parquet (малые файлы) Нет IO-bound, JSON плохо на GPU
Gold-слой витрины: тяжёлые агрегации Parquet > 50 GB Да Compute-bound join+agg
ML feature engineering на числовых колонках > 10M строк Да Compute-heavy, нет UDF
Python UDF тяжёлой бизнес-логики Нет Python-only, GPU простаивает
Streaming micro-batch (1–10 min interval) Нет Батчи слишком малы
JDBC ETL из реляционной СУБД Нет IO-bound от начала до конца
Репартиционирование / рескейл данных Нет Shuffle-bound, не compute
Сложные строковые regex / NLP Частично cuDF поддерживает базовые regex
COLLECT_LIST / COLLECT_SET агрегации Нет Не реализованы в cuDF

Архитектурный паттерн: гибридный кластер

Оптимальная архитектура для организаций с разными типами задач - не «GPU-кластер» или «CPU-кластер», а гибридный кластер с routing по типу задачи:

Входящие задачи:
├── IO-bound ETL / JSON → Parquet → CPU Worker Pool (r7i.8xlarge)
├── Python UDF heavy workloads → CPU Worker Pool (r7i.8xlarge)
├── Streaming micro-batches → CPU Worker Pool (c7i.4xlarge)
└── Аналитика Gold: тяжёлые AGG + JOIN на Parquet > 50GB → GPU Worker Pool (p4d.24xlarge)

Routing реализуется через Kubernetes node affinity или YARN queue management:

# Для аналитических задач на GPU:
spark = create_rapids_spark("analytics-heavy")
spark.conf.set("spark.kubernetes.node.selector.workload-type", "gpu-analytics")

# Для ETL-задач на CPU:
spark_etl = SparkSession.builder \
    .config("spark.kubernetes.node.selector.workload-type", "cpu-etl") \
    .getOrCreate()

Мониторинг и диагностика: как понять, что GPU работает эффективно

nvidia-smi: базовый мониторинг

# Запускайте на GPU-нодах во время выполнения задачи
watch -n 2 nvidia-smi --query-gpu=index,name,utilization.gpu,memory.used,memory.free,temperature.gpu \
  --format=csv,noheader

# Пример здорового вывода (GPU утилизирован):
# 0, NVIDIA A100-SXM4-80GB, 94%, 68000 MiB, 13000 MiB, 72
# 1, NVIDIA A100-SXM4-80GB, 91%, 71000 MiB, 9800 MiB,  69

# Пример нездорового вывода (GPU простаивает):
# 0, NVIDIA A100-SXM4-80GB, 3%, 4100 MiB,  76900 MiB, 41
# ↑ GPU занят 3% времени - плохой кандидат для GPU acceleration

Spark UI: ключевые метрики

В Spark UI открывайте вкладку "Stages" и изучайте распределение времени каждой задачи:

  • Task Time (синий): само вычисление
  • Fetch Wait Time (жёлтый): ожидание данных от shuffle (IO-bound!)
  • GPU Columnar To Row Time (если есть в метриках): PCIe-transfer

Если Fetch Wait Time > 30% суммарного времени - задача shuffle-bound, GPU не поможет.

RAPIDS Profiling Tool

# Генерируем детальный профиль выполнения
java -jar rapids-4-spark-tools_2.12-24.02.0.jar profiling \
  --output-directory /tmp/rapids-profile \
  --event-log /var/spark-history/app-20240115-123456

# Ключевые секции отчёта:
# 1. GPU Utilization per Stage:
#    Stage 1: 2.1% avg (IO dominated)
#    Stage 2: 89.3% avg (GPU execution)
#
# 2. CPU/GPU Transitions:
#    GpuColumnarToRow: 45 occurrences, 380s total
#    GpuRowToColumnar: 45 occurrences, 412s total
#    ↑ Если transitions > 5% от общего времени - ищите причины fallback
#
# 3. Recommendations:
#    - Stage 1 is IO-bound, GPU provides no benefit
#    - Remove Python UDF to eliminate 45 GPU-CPU transitions

Best Practices: золотые правила архитектора

Перед внедрением GPU: что сделать сначала

  1. Профилирование без GPU: запустите задачу на CPU с включённым spark.rapids.sql.explain=ALL. Посмотрите, какой процент операторов мог бы быть на GPU (без реального GPU запустить).

  2. RAPIDS Qualification Tool на существующих event logs: если у вас есть Spark History Server с логами CPU-запусков, Qualification Tool может оценить потенциальный speedup без GPU-кластера.

  3. Исправьте IO-проблемы: плохой format, мелкие файлы и медленные источники данных должны быть устранены до GPU. Иначе GPU будет ускорять 5% задачи.

  4. Устраните Python UDF: F.when, F.expr, встроенные функции - всё это GPU-совместимо. Python UDF - нет.

  5. Запустите RAPIDS Qualification на исправленном пайплайне и убедитесь, что Score > 70.

Чек-лист мониторинга в production

После запуска GPU-задач регулярно проверяйте:

# Скрипт мониторинга GPU эффективности (запускать во время задачи)
#!/bin/bash
for executor_node in gpu-node-1 gpu-node-2 gpu-node-3; do
    echo "=== $executor_node ==="
    ssh $executor_node "nvidia-smi --query-gpu=utilization.gpu,memory.used \
        --format=csv,noheader" 
done

# Ожидаемый вывод при нормальной работе:
# gpu-node-1: 87%, 62000 MiB
# gpu-node-2: 91%, 68000 MiB
# gpu-node-3: 84%, 59000 MiB

# Тревожный вывод:
# gpu-node-1: 4%, 2048 MiB  ← GPU простаивает!
# gpu-node-2: 3%, 2048 MiB  ← GPU простаивает!

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

Условие:

Вы - Senior Data Engineer. Команда аналитиков хочет перевести следующий пайплайн на GPU-кластер, рассчитывая на «10-кратное ускорение как в NVIDIA benchmarks»:

Профиль пайплайна:
- Источник: PostgreSQL JDBC (24 таблицы, суммарно 800 GB)
- Формат: читается через JDBC, параллелизм ограничен 50 connections
- Трансформации: 
    * 12 Python UDF для бизнес-правил (суммарно 35% времени)
    * JOIN 5 больших таблиц (суммарно 25% времени)
    * Агрегации GROUP BY (суммарно 15% времени)
    * Shuffle / Exchange (суммарно 20% времени)
    * Запись в S3 в CSV (5% времени)
- Текущее время: 4.5 часа на CPU (r6i.16xlarge × 8 нод)
- Текущая стоимость: $4.05/час × 8 × 4.5 = $145.8/запуск
- Предлагаемый GPU-кластер: p4d.24xlarge × 4 нод = $32.77 × 4 = $131.1/час

Метрики из Spark History Server:
- Stage "JDBC Read": Fetch Wait Time = 78% от task duration
- Stage "Python UDF": GPU Utilization = 0%
- Stage "JOIN": GPU Utilization = 0% (данные в JVM после UDF)
- Stage "Aggregate": GPU Utilization = 72% (единственная хорошая стадия)
- Stage "Write CSV": GPU Utilization = 0%

Задание:

Вы должны составить Architecture Decision Record (ADR) на 1–2 страницы с ответом на вопрос: следует ли переходить на GPU-кластер?

ADR должен содержать:

  1. Анализ Закона Амдала: рассчитайте теоретический максимальный speedup для данного пайплайна, если только агрегации (15% времени) выполняются на GPU в 10×. Покажите расчёт формулой.

  2. Экономический расчёт: посчитайте стоимость одного запуска на предложенном GPU-кластере (p4d.24xlarge × 4), учитывая, что GPU даёт не 10×, а только теоретически рассчитанный speedup из пункта 1. Сравните с текущей стоимостью ($145.8).

  3. Список корневых причин низкой GPU-эффективности: для каждого источника (JDBC, Python UDF, CSV write) объясните физическую причину, почему GPU не поможет.

  4. Альтернативный план (без GPU): что нужно сделать, чтобы ускорить пайплайн на существующем CPU-кластере?

  5. Конвертация из JDBC в Parquet (разовая операция)
  6. Замена Python UDF на F.when / F.expr
  7. Оптимизация write (Parquet вместо CSV) Рассчитайте ожидаемое ускорение и новую стоимость.

  8. Вывод ADR: "Принять" или "Отклонить" предложение о GPU. Обоснуйте одним параграфом.

Ожидаемый результат: ADR с математически обоснованным отказом от GPU в данном конкретном кейсе и альтернативным планом оптимизации, который даёт сопоставимое ускорение за меньшую стоимость.