RAPIDS: когда GPU не даёт прироста - малые данные, IO-bound запросы
RAPIDS: когда GPU не даёт прироста - малые данные, IO-bound запросы
Парадокс 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 дороже
Правильная последовательность действий:
- Конвертировать JSON → Parquet (бесплатно, разовая задача)
- Заменить Python UDF на
F.when/ встроенные функции (разработка, 2–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: что сделать сначала¶
-
Профилирование без GPU: запустите задачу на CPU с включённым
spark.rapids.sql.explain=ALL. Посмотрите, какой процент операторов мог бы быть на GPU (без реального GPU запустить). -
RAPIDS Qualification Tool на существующих event logs: если у вас есть Spark History Server с логами CPU-запусков, Qualification Tool может оценить потенциальный speedup без GPU-кластера.
-
Исправьте IO-проблемы: плохой format, мелкие файлы и медленные источники данных должны быть устранены до GPU. Иначе GPU будет ускорять 5% задачи.
-
Устраните Python UDF:
F.when,F.expr, встроенные функции - всё это GPU-совместимо. Python UDF - нет. -
Запустите 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 должен содержать:
-
Анализ Закона Амдала: рассчитайте теоретический максимальный speedup для данного пайплайна, если только агрегации (15% времени) выполняются на GPU в 10×. Покажите расчёт формулой.
-
Экономический расчёт: посчитайте стоимость одного запуска на предложенном GPU-кластере (
p4d.24xlarge × 4), учитывая, что GPU даёт не 10×, а только теоретически рассчитанный speedup из пункта 1. Сравните с текущей стоимостью ($145.8). -
Список корневых причин низкой GPU-эффективности: для каждого источника (JDBC, Python UDF, CSV write) объясните физическую причину, почему GPU не поможет.
-
Альтернативный план (без GPU): что нужно сделать, чтобы ускорить пайплайн на существующем CPU-кластере?
- Конвертация из JDBC в Parquet (разовая операция)
- Замена Python UDF на
F.when/F.expr -
Оптимизация write (Parquet вместо CSV) Рассчитайте ожидаемое ускорение и новую стоимость.
-
Вывод ADR: "Принять" или "Отклонить" предложение о GPU. Обоснуйте одним параграфом.
Ожидаемый результат: ADR с математически обоснованным отказом от GPU в данном конкретном кейсе и альтернативным планом оптимизации, который даёт сопоставимое ускорение за меньшую стоимость.