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.