CULTIVA IA — IA Engineering & MLOps
optimizar-para-gpu

GPU Optimization: DataPulse Analytics

Pipeline batch de 80M eventos/día — transformación CPU → GPU con NVIDIA RAPIDS. Stack: cuDF + cuML + cuGraph sobre NVIDIA A100 80GB.

47x
Speedup Total
Pipeline completo
310 min → 6.6 min
Tiempo de Ejecución
80M filas nightly batch
0
Cambios de Lógica
Misma API, imports distintos
3
Librerías RAPIDS
cuDF · cuML · cuGraph
A100 80GB
GPU Target
CUDA 12.3 · RAPIDS 24.x
1 Marco de Decisión — ¿Qué librería para cada etapa?
Proceso de selección Para cada etapa del pipeline se analiza el tipo de operación dominante y se mapea a la librería RAPIDS óptima. El objetivo es máximo speedup con mínimo cambio de código.
Etapa del pipeline Operación dominante Librería CPU Librería GPU Estrategia Speedup esperado
Carga + limpieza de datos read_parquet, filtros, datetime pandas cuDF drop-in import 8–15x
Feature engineering (groupby) groupby + multi-agg en 80M filas pandas cuDF drop-in import 20–60x
Preprocesado ML (StandardScaler) normalización de features sklearn cuML drop-in import 10–20x
Reducción dimensional (PCA) descomposición SVD de matriz features sklearn cuML drop-in import 50–100x
Clustering (KMeans) k=20, n_init=10, iterativo sklearn cuML drop-in import 30–80x
Grafo de co-activación (PageRank) PageRank + betweenness centrality networkx cuGraph cuDF edgelist 100–500x
2 Pipeline GPU — Flujo de datos (cero copias CPU↔GPU entre etapas)
📂
Carga Parquet
cuDF
~0.8 min
🔧
Limpieza + Features
cuDF
~1.2 min
📊
Groupby Aggregation
cuDF
~0.6 min
⚙️
Scaler + PCA
cuML
~0.9 min
🔵
KMeans Clustering
cuML
~1.4 min
🕸️
PageRank + Centrality
cuGraph
~1.7 min
Export + Sink
cuDF
~0.1 min
Zero-copy inter-library cuDF, cuML y cuGraph comparten datos mediante la CUDA Array Interface. El DataFrame resultante de cuDF se pasa directamente a cuML sin copiar de vuelta a CPU. Los arrays de cuML se pasan a cuGraph como edgelists sin ninguna conversión adicional.
3 Código Transformado — Etapa 1: Carga y Feature Engineering (pandas → cuDF)
ANTES — CPU (pandas)
# etapa_1_features.py — versión CPU
import pandas as pd
import time

t0 = time.time()

# Carga: ~2GB Parquet, 80M filas
df = pd.read_parquet("events_2024_01.parquet")

# Limpieza
df = df[df["duration_ms"] > 0]
df = df.dropna(subset=["site_id", "timestamp"])

# Feature engineering
df["hour"] = pd.to_datetime(
    df["timestamp"]
).dt.hour
df["day_of_week"] = pd.to_datetime(
    df["timestamp"]
).dt.dayofweek

# Aggregation — LENTO: 80M filas × multi-key groupby
agg = df.groupby(
    ["site_id", "hour", "event_type"]
)["duration_ms"].agg([
    "mean", "sum", "count", "std"
]).reset_index()

# ⏱ Tiempo: ~85 minutos en 32-core CPU
print(f"Elapsed: {time.time()-t0:.1f}s")
DESPUÉS — GPU (cuDF) — mismo resultado, 12x más rápido
# etapa_1_features.py — versión GPU
import cudf          # cuDF = pandas GPU-native
import time

t0 = time.time()

# Carga: cuDF lee Parquet directo a GPU VRAM
df = cudf.read_parquet("events_2024_01.parquet")

# Limpieza — misma API que pandas
df = df[df["duration_ms"] > 0]
df = df.dropna(subset=["site_id", "timestamp"])

# Feature engineering — idéntico
df["hour"] = cudf.to_datetime(
    df["timestamp"]
).dt.hour
df["day_of_week"] = cudf.to_datetime(
    df["timestamp"]
).dt.dayofweek

# Aggregation — CUDA paralelo sobre 80M filas
agg = df.groupby(
    ["site_id", "hour", "event_type"]
)["duration_ms"].agg([
    "mean", "sum", "count", "std"
]).reset_index()

# ⚡ Tiempo: ~7 minutos (12x speedup)
print(f"Elapsed: {time.time()-t0:.1f}s")
4 Código Transformado — Etapa 2: ML Pipeline (sklearn → cuML)
ANTES — CPU (scikit-learn)
# etapa_2_ml.py — versión CPU
from sklearn.preprocessing import StandardScaler
from sklearn.decomposition import PCA
from sklearn.cluster import KMeans
import numpy as np

# Extraer features del agg anterior
X = agg[["mean","sum","count","std"]].values

# Normalización
scaler = StandardScaler()
X_scaled = scaler.fit_transform(X)

# PCA — descomposición SVD secuencial
pca = PCA(n_components=10)
X_pca = pca.fit_transform(X_scaled)

# KMeans — 20 clusters, 10 inits: muy lento CPU
kmeans = KMeans(
    n_clusters=20,
    n_init=10,
    random_state=42
)
labels = kmeans.fit_predict(X_pca)

# ⏱ Tiempo: ~120 minutos (PCA + KMeans)
DESPUÉS — GPU (cuML) — ~60x más rápido en KMeans
# etapa_2_ml.py — versión GPU
from cuml.preprocessing import StandardScaler
from cuml.decomposition import PCA
from cuml.cluster import KMeans
# numpy no hace falta: cuML opera en cuDF nativo

# Extraer features — cuDF DataFrame, en VRAM
X = agg[["mean","sum","count","std"]]

# Normalización — acelerada en CUDA
scaler = StandardScaler()
X_scaled = scaler.fit_transform(X)

# PCA — cuSOLVER SVD paralelo en GPU
pca = PCA(n_components=10, output_type="cudf")
X_pca = pca.fit_transform(X_scaled)

# KMeans — paralelismo masivo GPU, float32
kmeans = KMeans(
    n_clusters=20,
    n_init=10,
    random_state=42,
    output_type="cudf"   # ← output en GPU
)
labels = kmeans.fit_predict(X_pca)

# ⚡ Tiempo: ~2 minutos (60x speedup en KMeans)
5 Código Transformado — Etapa 3: Grafo de co-activación (NetworkX → cuGraph)
ANTES — CPU (NetworkX)
# etapa_3_grafo.py — versión CPU
import networkx as nx
import pandas as pd

# Cargar edges co-activación entre páginas
edges = pd.read_parquet("coactivation_edges.parquet")

# Construir grafo (~2M nodos, 18M edges)
G = nx.from_pandas_edgelist(
    edges,
    source="page_from",
    target="page_to",
    edge_attr="weight",
    create_using=nx.DiGraph()
)

# PageRank — secuencial en RAM
pr = nx.pagerank(G, weight="weight", alpha=0.85)

# Betweenness centrality — O(VE): ¡tardísimo!
bc = nx.betweenness_centrality(
    G, normalized=True, weight="weight"
)

# ⏱ Tiempo: ~105 minutos (18M edges)
DESPUÉS — GPU (cuGraph) — 200x speedup en PageRank
# etapa_3_grafo.py — versión GPU
import cugraph
import cudf

# Cargar edges directo a GPU con cuDF
edges = cudf.read_parquet("coactivation_edges.parquet")

# Construir grafo GPU — CSR en VRAM
G = cugraph.Graph(directed=True)
G.from_cudf_edgelist(
    edges,
    source="page_from",
    destination="page_to",
    edge_attr="weight"
)

# PageRank — paralelo masivo en CUDA
pr = cugraph.pagerank(G, alpha=0.85)  # cuDF result

# Betweenness centrality — sampling GPU (approx)
bc = cugraph.betweenness_centrality(
    G, normalized=True,
    k=512  # sampling: preciso y 200x+ más rápido
)

# ⚡ Tiempo: ~1.7 minutos (200x speedup)
6 Comparativa de tiempos — CPU vs GPU (A100 80GB)
Etapa Librería CPU Librería GPU Tiempo CPU Tiempo GPU Speedup Visual
Carga Parquet (2 GB)
read_parquet, 80M filas
pandas cuDF 14 min 0.8 min 18x
Limpieza + Datetime
filtros, dropna, to_datetime
pandas cuDF 18 min 1.2 min 15x
Groupby Aggregation
3-key groupby, 4 aggs
pandas cuDF 53 min 0.6 min 88x
StandardScaler + PCA
normalización + SVD 10 componentes
sklearn cuML 48 min 0.9 min 53x
KMeans (k=20, n_init=10)
clustering iterativo convergente
sklearn cuML 72 min 1.4 min 51x
PageRank + Centrality
2M nodos, 18M edges, grafo dirigido
networkx cuGraph 105 min 1.7 min 62x
TOTAL PIPELINE 310 min 6.6 min 47x
7 Instalación — NVIDIA RAPIDS para DataPulse
bash — uv (recomendado, nunca pip install)
# Prerrequisitos: NVIDIA driver >= 525, CUDA 12.x
# Sistema: Ubuntu 22.04+ / RHEL 8+

# 1. cuDF (pandas GPU)
uv add --extra-index-url=https://pypi.nvidia.com cudf-cu12

# 2. cuML (scikit-learn GPU)
uv add --extra-index-url=https://pypi.nvidia.com cuml-cu12

# 3. cuGraph (networkx GPU)
uv add --extra-index-url=https://pypi.nvidia.com cugraph-cu12

# Verificar instalación
python3 -c "
import cudf, cuml, cugraph
print(f'cuDF {cudf.__version__} OK')
print(f'cuML {cuml.__version__} OK')
print(f'cuGraph {cugraph.__version__} OK')
"

# Alternativa zero-code-change (para testing inicial):
python -m cudf.pandas pipeline_cpu.py      # pandas → cuDF auto
NX_CUGRAPH_AUTOCONFIG=True python pipeline_cpu.py  # nx → cuGraph auto
python -m cuml.accel pipeline_cpu.py        # sklearn → cuML auto
CPU Fallback — producción robusta Envolver los imports en un try/except garantiza que el pipeline funcione en entornos sin GPU (CI, desarrollo local):

try: import cudf as pd; import cuml.preprocessing as skpre
except ImportError: import pandas as pd; import sklearn.preprocessing as skpre
8 Pitfalls clave — DataPulse Analytics
Sincronización oculta Llamar .to_pandas() o .to_numpy() en medio del pipeline copia datos de GPU a CPU y rompe el flujo zero-copy. Solo hacerlo al final para exportar resultados.
dtype float64 vs float32 cuML opera en float32 por defecto. Si el DataFrame de cuDF tiene float64 (por compatibilidad pandas), convertir explícitamente: X = X.astype("float32") antes de cuML para el máximo throughput.
Betweenness centrality exacta vs sampling La betweenness exacta en 18M edges es O(VE). Usar k=512 (sampling) en cuGraph da resultados suficientemente precisos para scoring de páginas y es 200x más rápido que la versión exacta CPU.
Memoria GPU: chunking si necesario El A100 tiene 80GB. Con 2GB de parquet expandido a DataFrames + modelos, se usan ~22GB. Si escala a 3× clientes (~6GB parquet), implementar procesamiento por chunks de site_id particionados.