Lucía tiene una pregunta que no se responde bien con SQL: qué productos se compran juntos. Si un cliente mete unos crampones en el carrito, ¿qué probabilidad hay de que también compre un piolet? Con esa matriz, AlpinaShop podría sugerir complementos en la ficha de producto, agrupar artículos en packs y colocar mejor las categorías. Es un cálculo sobre todos los pares posibles dentro de cada pedido, repetido sobre dos años de historia, y es el tipo de problema que empieza siendo elegante en SQL y termina siendo un JOIN de una tabla consigo misma que nadie quiere mantener.
Es, además, exactamente el tipo de problema para el que se inventó Spark: cálculo iterativo y algorítmico sobre grandes volúmenes, expresado en código, no en consultas.
Aquí conviene ser honesto sobre por qué esta lección existe. Si AlpinaShop empezara hoy desde cero, probablemente resolvería esto con BigQuery y Dataflow y no tocaría Spark nunca. Pero el mundo real no empieza de cero: hay decenas de miles de empresas con clústeres de Hadoop en sus centros de datos, con años de código Spark en producción, con equipos que saben PySpark y no saben Beam, y con bibliotecas —MLlib, GraphX, todo el ecosistema de Python científico— que no tienen equivalente directo. Cloud Dataproc es la respuesta de Google a esa realidad: Hadoop y Spark gestionados, con la particularidad de que le da la vuelta a la idea misma de clúster.
En esta lección verás qué son Hadoop y Spark en lo esencial, crearás el clúster alpinashop-spark, entenderás el patrón que hace de Dataproc algo distinto a un Hadoop alquilado —el clúster efímero sobre almacenamiento en Cloud Storage—, ejecutarás el trabajo PySpark que responde a la pregunta de Lucía, y compararás con criterio Dataproc, Dataproc Serverless y Dataflow.
Contenido
- Hadoop y Spark en una página
- Por qué siguen importando en 2026
- Qué aporta Dataproc
- Anatomía de un clúster de Dataproc
- Crear
alpinashop-spark - El patrón clave: clúster efímero sobre Cloud Storage
- Enviar trabajos:
gcloud dataproc jobs submit - El trabajo PySpark de AlpinaShop: cesta media y productos comprados juntos
- Acciones de inicialización y versiones de imagen
- Dataproc Serverless para Spark
- Dataproc, Serverless y Dataflow: la tabla de decisión
- Notebooks y Spark SQL sobre BigQuery
- Migrar un Hadoop on-premise a Google Cloud
- Coste y VM Spot
- Hadoop y Spark en una página
Hadoop nació en 2006 para resolver un problema concreto: procesar más datos de los que cabían en una máquina, usando muchos ordenadores baratos que fallan a menudo. Tiene tres piezas:
- HDFS, un sistema de ficheros distribuido que parte los ficheros en bloques y los replica (por defecto tres veces) por los discos de las máquinas del clúster.
- YARN, el gestor de recursos que decide qué proceso corre en qué máquina.
- MapReduce, el modelo de programación original: partes el trabajo en una fase map (transformar cada registro) y otra reduce (agregar por clave), y el framework se encarga del reparto y de los fallos.
MapReduce funcionaba y era lentísimo, por una razón de diseño: escribía en disco entre cada fase. Un algoritmo iterativo que necesita veinte pasadas sobre los mismos datos hacía veinte rondas de escritura y lectura en disco.
Spark apareció en 2014 con la corrección obvia: mantener los datos intermedios en memoria. Para el mismo algoritmo iterativo, la mejora era de uno o dos órdenes de magnitud. Y añadió una API mucho más agradable.
Spark se organiza hoy en torno a:
| Componente | Qué es | Uso |
|---|---|---|
| Spark Core / RDD | La abstracción original: colección distribuida y resiliente | Control fino, operaciones no expresables en SQL |
| DataFrame / Spark SQL | Tablas con esquema y un optimizador (Catalyst) | El 90 % del uso actual; SQL sobre datos distribuidos |
| MLlib | Biblioteca de machine learning distribuido | Clustering, recomendación, clasificación a escala |
| Structured Streaming | Procesamiento continuo con la API de DataFrame | Alternativa a Beam dentro del mundo Spark |
| GraphX / GraphFrames | Algoritmos sobre grafos | Redes, caminos, comunidades |
La arquitectura de ejecución, que hay que tener en la cabeza para entender lo que viene:
flowchart TD
D["Driver<br/>tu programa PySpark<br/>construye el plan"]
CM["Gestor de recursos (YARN)"]
E1["Executor 1<br/>tareas + cache en memoria"]
E2["Executor 2"]
E3["Executor N"]
S["Almacenamiento<br/>HDFS o Cloud Storage"]
D -->|pide recursos| CM
CM -->|asigna| E1 & E2 & E3
D -->|envia tareas| E1 & E2 & E3
E1 & E2 & E3 <--> S
El driver ejecuta tu código, construye un grafo de operaciones y lo trocea en tareas. Los executors ejecutan esas tareas sobre particiones de los datos. Igual que en Beam, las transformaciones son perezosas: no ocurre nada hasta que llamas a una acción (count(), collect(), write()). Esa pereza es lo que permite al optimizador reordenar y fusionar operaciones.
- Por qué siguen importando en 2026
Con BigQuery y Dataflow disponibles, ¿por qué aprender esto? Cuatro razones concretas y una consecuencia.
Código existente. Una empresa que migra a la nube con 200.000 líneas de PySpark en producción no las va a reescribir. Reescribir no aporta valor de negocio, introduce errores y consume meses. Dataproc permite mover esas cargas sin tocar el código, y ya se optimizará después.
Personas. El mercado tiene muchos más ingenieros que saben Spark que ingenieros que saben Beam. Si el equipo de datos de AlpinaShop contrata mañana, es más probable que el candidato traiga Spark. Elegir la tecnología que tu equipo sabe usar es una decisión de arquitectura legítima, no una concesión.
MLlib y el ecosistema de Python. Algoritmos como ALS para recomendación, k-means a escala o FP-Growth para reglas de asociación están implementados, probados y distribuidos. Escribirlos en Beam sería absurdo. Y dentro de un job de Spark puedes usar pandas, NumPy o scikit-learn sobre particiones concretas.
Formatos y ecosistema abiertos. Hive, Presto/Trino, HBase, Kafka, Iceberg, Delta Lake, Hudi: un mundo entero de herramientas abiertas que habla el idioma de Hadoop. Si AlpinaShop quisiera irse de Google Cloud, un lago de datos en Parquet sobre almacenamiento de objetos y procesado con Spark viaja a cualquier proveedor sin cambios.
La consecuencia: Dataproc no es un servicio de segunda ni una reliquia. Es la vía de migración menos traumática y la herramienta correcta cuando el problema es algorítmico y el equipo sabe Spark.
- Qué aporta Dataproc
Montar un clúster de Hadoop a mano —instalar, configurar YARN, dimensionar HDFS, ajustar la memoria de los executors, integrar la autenticación— es un trabajo de semanas y una fuente permanente de mantenimiento. Dataproc lo reduce a un comando.
| Aspecto | Hadoop autogestionado | Dataproc |
|---|---|---|
| Tiempo de creación | Días o semanas | Menos de 2 minutos |
| Configuración | Manual, por componente | Preconfigurada y coherente |
| Escalado | Comprar y montar hardware | Cambiar un número; autoescalado disponible |
| Actualizaciones | Proyecto en sí mismo | Versiones de imagen gestionadas |
| Coste en reposo | El hardware, siempre | Cero si el clúster no existe |
| Integración con la nube | Hay que construirla | Cloud Storage, BigQuery, Logging, IAM de serie |
Los dos minutos de creación no son un dato de marketing: son lo que cambia el modelo mental. Si crear un clúster cuesta semanas, el clúster es una instalación permanente que hay que cuidar, compartir entre equipos y mantener siempre encendida por si acaso. Si cuesta noventa segundos, el clúster pasa a ser desechable: se crea para un trabajo, se destruye al terminar, y cada trabajo puede tener el suyo con la versión y las bibliotecas que necesita.
Eso es lo que veremos en el apartado 6, y es la idea más importante de la lección.
- Anatomía de un clúster de Dataproc
Un clúster de Dataproc tiene tres tipos de nodo:
| Nodo | Función | Cantidad | Notas |
|---|---|---|---|
| Maestro | Ejecuta el driver, YARN ResourceManager, HDFS NameNode | 1 (o 3 en alta disponibilidad) | Si cae con 1 nodo, el trabajo muere |
| Workers primarios | Ejecutan tareas; aportan disco a HDFS | Mínimo 2 (o 0 en modo nodo único) | VM estándar, estables |
| Workers secundarios | Solo cómputo; no almacenan HDFS | 0 a N | Pueden ser Spot: hasta ~80 % más baratos |
Los workers secundarios merecen atención. Como no participan en HDFS, pueden desaparecer sin que se pierda ningún dato. Por eso pueden ser VM Spot (las mismas de 02-01, que Google puede reclamar con 30 segundos de aviso). Si una desaparece a mitad de tarea, YARN reasigna esa tarea a otro nodo y el trabajo continúa, más lento pero correcto.
Esa es la combinación ganadora para AlpinaShop: pocos workers primarios estándar para dar estabilidad, y muchos secundarios Spot para dar potencia barata.
flowchart TD
subgraph C["Cluster alpinashop-spark"]
M["Maestro<br/>n2-standard-4<br/>driver + YARN RM"]
W1["Worker primario 1<br/>n2-standard-4"]
W2["Worker primario 2<br/>n2-standard-4"]
S1["Secundario Spot 1"]
S2["Secundario Spot 2"]
S3["Secundario Spot N<br/>autoescalado"]
end
GCS["Cloud Storage<br/>gs://alpinashop-catalogo<br/>gs://alpinashop-datalake"]
BQ["BigQuery<br/>alpinashop_analitica"]
M --- W1 & W2 & S1 & S2 & S3
W1 & W2 & S1 & S2 & S3 <--> GCS
W1 & W2 <--> BQ
Fíjate en que el almacenamiento está fuera del clúster. Ese es el siguiente apartado.
- Crear
alpinashop-spark
alpinashop-sparkgcloud config set project alpinashop-datos
# Bucket propio para el lago de datos y los artefactos de Spark
gcloud storage buckets create gs://alpinashop-datalake \
--project=alpinashop-datos --location=europe-west1 \
--uniform-bucket-level-access
# Cuenta de servicio dedicada, con minimo privilegio
gcloud iam service-accounts create sa-dataproc-analitica \
--display-name="Clusteres de Dataproc de analitica"
SA="[email protected]"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/dataproc.worker"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/bigquery.dataEditor"
gcloud projects add-iam-policy-binding alpinashop-datos \
--member="serviceAccount:$SA" --role="roles/bigquery.jobUser"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-datalake \
--member="serviceAccount:$SA" --role="roles/storage.objectAdmin"
gcloud storage buckets add-iam-policy-binding gs://alpinashop-catalogo \
--member="serviceAccount:$SA" --role="roles/storage.objectViewer"Y el clúster:
gcloud dataproc clusters create alpinashop-spark \
--region=europe-west1 \
--zone=europe-west1-b \
--service-account="$SA" \
--subnet=sn-datos-euw1 \
--no-address \
--master-machine-type=n2-standard-4 \
--master-boot-disk-size=100GB \
--master-boot-disk-type=pd-balanced \
--num-workers=2 \
--worker-machine-type=n2-standard-4 \
--worker-boot-disk-size=200GB \
--num-secondary-workers=2 \
--secondary-worker-type=spot \
--image-version=2.2-debian12 \
--optional-components=JUPYTER \
--enable-component-gateway \
--bucket=alpinashop-datalake \
--max-idle=30m \
--properties="spark:spark.sql.adaptive.enabled=true,\
spark:spark.dynamicAllocation.enabled=true,\
spark:spark.sql.sources.partitionOverwriteMode=dynamic" \
--labels=entorno=produccion,equipo=datos,centro-coste=analitica,aplicacion=analiticaRepaso de lo que importa:
| Opción | Por qué está |
|---|---|
--subnet=sn-datos-euw1 + --no-address |
El clúster vive en alpinashop-vpc, sin IP pública. Sale a internet por el Cloud NAT de 03-01 y accede a las API por Private Google Access |
--service-account |
Identidad propia con permisos mínimos, no la cuenta por defecto de Compute |
--num-secondary-workers=2 --secondary-worker-type=spot |
Potencia barata que puede desaparecer sin romper nada |
--image-version=2.2-debian12 |
Versión fijada. Sin esto, Google elegiría la más reciente y un trabajo que funcionaba ayer podría fallar mañana |
--optional-components=JUPYTER + --enable-component-gateway |
Notebooks accesibles desde la consola con autenticación de IAM, sin abrir puertos |
--bucket=alpinashop-datalake |
Bucket de trabajo para logs y ficheros temporales del clúster |
--max-idle=30m |
El clúster se autodestruye tras 30 minutos sin trabajos. La opción más rentable de todo el comando |
spark.sql.adaptive.enabled |
Ejecución adaptativa: Spark reajusta particiones y estrategias de JOIN en tiempo real. Mitiga bastante el sesgo de datos |
spark.dynamicAllocation.enabled |
Spark pide y libera executors según necesita |
Añadir autoescalado al clúster requiere una política aparte:
# politica-autoescalado.yaml
workerConfig:
minInstances: 2
maxInstances: 2 # los primarios NO escalan: aportan HDFS
secondaryWorkerConfig:
minInstances: 0
maxInstances: 20 # los Spot si escalan, hasta 20
basicAlgorithm:
cooldownPeriod: 2m
yarnConfig:
scaleUpFactor: 1.0 # anade el 100 % de lo que YARN pide
scaleDownFactor: 0.5 # retira la mitad de lo sobrante: prudente
gracefulDecommissionTimeout: 10m # espera a que terminen las tareas en cursogcloud dataproc autoscaling-policies import pol-autoescalado-analitica \
--region=europe-west1 --source=politica-autoescalado.yaml
gcloud dataproc clusters update alpinashop-spark --region=europe-west1 \
--autoscaling-policy=pol-autoescalado-analiticaEl gracefulDecommissionTimeout es importante: sin él, al reducir el clúster se matan nodos con tareas en curso y hay que rehacerlas. Con él, se espera a que terminen. Y scaleDownFactor: 0.5 evita el efecto acordeón de subir y bajar constantemente.
Verificación:
gcloud dataproc clusters list --region=europe-west1 \
--format="table(clusterName, status.state, config.workerConfig.numInstances)"
gcloud dataproc clusters describe alpinashop-spark --region=europe-west1
- El patrón clave: clúster efímero sobre Cloud Storage
Aquí está la idea que cambia todo.
En un Hadoop tradicional, el almacenamiento y el cómputo están en las mismas máquinas. HDFS vive en los discos de los workers. Eso tenía una razón excelente en 2006: mover datos por la red era carísimo comparado con leerlos del disco local, así que el principio era "lleva el cómputo al dato".
Pero tiene una consecuencia devastadora: si apagas el clúster, pierdes los datos. Por eso los clústeres Hadoop están siempre encendidos, aunque solo trabajen tres horas al día. Se paga el hardware las veinticuatro.
En Google Cloud, la red interna es tan rápida que la premisa ya no se sostiene: leer de Cloud Storage no es significativamente más lento que leer de un disco local. Eso permite invertir el diseño:
flowchart LR
subgraph Antes["Hadoop tradicional"]
H["Cluster permanente<br/>computo + HDFS<br/>encendido 24x7"]
end
subgraph Ahora["Patron Dataproc"]
GCS["Cloud Storage<br/>gs://alpinashop-datalake<br/>PERMANENTE, barato"]
C1["Cluster efimero A<br/>Spark 3.5<br/>vive 20 min"]
C2["Cluster efimero B<br/>Spark 3.3 + libreria X<br/>vive 5 min"]
end
GCS <--> C1
GCS <--> C2
El conector de Cloud Storage viene preinstalado en Dataproc y hace que las rutas gs:// se comporten como rutas de HDFS para cualquier código de Spark o Hadoop:
# El mismo codigo que leia de HDFS...
df = spark.read.parquet("hdfs:///datos/pedidos/")
# ...lee de Cloud Storage cambiando el prefijo. Nada mas.
df = spark.read.parquet("gs://alpinashop-datalake/pedidos/")Las ventajas de separar almacenamiento y cómputo son concretas:
- Pagas cómputo solo cuando calculas. Un clúster que vive 20 minutos al día cuesta el 1,4 % de uno permanente.
- Los datos sobreviven al clúster. Puedes destruirlo con total tranquilidad.
- Varios clústeres sobre los mismos datos. El de Lucía con Spark 3.5 y el de un proveedor externo con una versión antigua, simultáneamente, sin interferirse.
- Durabilidad muy superior. Cloud Storage replica con garantías de 11 nueves; HDFS con tres copias en tres discos del mismo rack no se acerca.
- Los datos son accesibles desde fuera de Spark. BigQuery los lee como tabla externa, Dataflow los procesa, la aplicación los descarga.
- Sin ceremonia de actualización. Para pasar a una versión nueva de Spark, creas un clúster nuevo. No migras nada.
Los inconvenientes, para ser justos: latencia algo mayor en operaciones de muchos ficheros pequeños, y que Cloud Storage no tiene renombrado atómico de directorios, lo que afecta a ciertos patrones de escritura. Se mitiga con formatos columnares y ficheros de tamaño razonable (128-512 MB), que es lo que hay que hacer de todas formas.
La regla de AlpinaShop: HDFS solo como espacio de trabajo temporal dentro de un trabajo. Nada que deba sobrevivir al clúster se escribe en HDFS. Y --max-idle en todos los clústeres, sin excepción. Un clúster olvidado un fin de semana cuesta más que todo el resto de la analítica del mes.
Para trabajos programados, ni siquiera se mantiene un clúster: se crea, se ejecuta y se destruye en un solo paso con los workflow templates:
# 1) Plantilla de flujo de trabajo
gcloud dataproc workflow-templates create wf-cesta-media --region=europe-west1
# 2) Cluster gestionado: nace y muere con el flujo
gcloud dataproc workflow-templates set-managed-cluster wf-cesta-media \
--region=europe-west1 \
--cluster-name=cluster-efimero-cesta \
--service-account="$SA" --subnet=sn-datos-euw1 --no-address \
--master-machine-type=n2-standard-4 \
--worker-machine-type=n2-standard-4 --num-workers=2 \
--num-secondary-workers=4 --secondary-worker-type=spot \
--image-version=2.2-debian12
# 3) El trabajo que se ejecutara
gcloud dataproc workflow-templates add-job pyspark \
gs://alpinashop-datalake/jobs/cesta_media.py \
--step-id=cesta-media --workflow-template=wf-cesta-media \
--region=europe-west1 \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31
# 4) Ejecutar: crea el cluster, lanza el job, destruye el cluster
gcloud dataproc workflow-templates instantiate wf-cesta-media --region=europe-west1Ese comando final es el que invocará el orquestador de 04-06. Coste total: los minutos que dure el trabajo. Cero el resto del mes.
- Enviar trabajos:
gcloud dataproc jobs submit
gcloud dataproc jobs submitDataproc admite varios tipos de trabajo:
# PySpark
gcloud dataproc jobs submit pyspark gs://alpinashop-datalake/jobs/mi_job.py \
--cluster=alpinashop-spark --region=europe-west1 \
-- arg1 arg2
# Spark (JAR de Scala o Java)
gcloud dataproc jobs submit spark --cluster=alpinashop-spark --region=europe-west1 \
--class=com.alpinashop.Informe --jars=gs://alpinashop-datalake/jars/informes.jar
# Spark SQL desde un fichero
gcloud dataproc jobs submit spark-sql --cluster=alpinashop-spark \
--region=europe-west1 --file=gs://alpinashop-datalake/sql/ventas.sql
# Hive, Pig, Presto/Trino tambien estan disponiblesTodo lo que va después de -- son argumentos para tu programa, no para gcloud. Es una confusión muy frecuente.
Opciones útiles al enviar:
gcloud dataproc jobs submit pyspark gs://alpinashop-datalake/jobs/cesta_media.py \
--cluster=alpinashop-spark \
--region=europe-west1 \
--py-files=gs://alpinashop-datalake/jobs/utilidades.zip \
--files=gs://alpinashop-datalake/config/categorias.json \
--jars=gs://spark-lib/bigquery/spark-3.5-bigquery-0.42.0.jar \
--properties="spark.executor.memory=6g,spark.executor.cores=2,spark.sql.shuffle.partitions=200" \
--labels=proceso=cesta-media \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31 \
--salida=gs://alpinashop-datalake/resultados/cesta/--py-files: módulos Python propios, empaquetados en.zipo.egg, distribuidos a todos los executors.--files: ficheros de datos que el trabajo necesita, accesibles por nombre en el directorio de trabajo.--jars: dependencias Java, como el conector de BigQuery.spark.sql.shuffle.partitions: número de particiones tras un shuffle. El valor por defecto (200) es un mal ajuste para casi todo el mundo: demasiado para datos pequeños, insuficiente para grandes. Una regla razonable es 2-3 veces el número total de cores del clúster.
Seguimiento:
gcloud dataproc jobs list --region=europe-west1 --cluster=alpinashop-spark \
--format="table(reference.jobId, status.state, statusHistory[0].stateStartTime)"
gcloud dataproc jobs wait JOB_ID --region=europe-west1 # sigue los logs en vivoLos logs van automáticamente a Cloud Logging (06-06) y la interfaz de Spark History Server queda accesible por el component gateway, incluso después de destruir el clúster si se configura un servidor de historial persistente.
- El trabajo PySpark de AlpinaShop: cesta media y productos comprados juntos
Ahora el trabajo real. Lee las líneas de pedido del lago, calcula la cesta media por mes y país, y construye la matriz de coocurrencia de productos.
"""
cesta_media.py -- Analisis de cesta de AlpinaShop con PySpark.
Entrada : gs://alpinashop-datalake/pedidos/ (Parquet, particionado por fecha)
Salidas : gs://alpinashop-datalake/resultados/cesta/ (Parquet)
alpinashop-datos.alpinashop_analitica.productos_juntos (BigQuery)
Envio:
gcloud dataproc jobs submit pyspark gs://.../cesta_media.py \
--cluster=alpinashop-spark --region=europe-west1 \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31
"""
import argparse
from pyspark.sql import SparkSession, functions as F, Window
PROYECTO = "alpinashop-datos"
DATASET = "alpinashop_analitica"
LAGO = "gs://alpinashop-datalake"
def crear_sesion():
"""La SparkSession es el punto de entrada. En Dataproc, la configuracion
de recursos y el maestro los aporta YARN: no hay que indicarlos."""
return (
SparkSession.builder
.appName("alpinashop-cesta-media")
# Bucket temporal que necesita el conector de BigQuery para escribir
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate()
)def cargar_lineas(spark, desde, hasta):
"""Lee las lineas de pedido del lago y las filtra por fecha.
Parquet guarda el esquema y las estadisticas por bloque, asi que Spark
puede saltarse ficheros enteros: es el 'predicate pushdown'.
"""
lineas = (
spark.read.parquet(f"{LAGO}/pedidos/lineas/")
.filter((F.col("fecha_pedido") >= desde) & (F.col("fecha_pedido") <= hasta))
# Seleccionar columnas pronto reduce memoria y shuffle
.select("pedido_id", "fecha_pedido", "sku", "cantidad",
"precio_unitario", "importe_linea")
)
cabeceras = (
spark.read.parquet(f"{LAGO}/pedidos/cabeceras/")
.filter((F.col("fecha_pedido") >= desde) & (F.col("fecha_pedido") <= hasta))
.filter(~F.col("estado").isin("cancelado", "devuelto"))
.select("pedido_id", "pais", "canal", "cliente_id")
)
# broadcast(): las cabeceras del periodo caben en memoria de cada executor.
# Evita el shuffle del JOIN, igual que la entrada lateral de Beam en 04-02.
return lineas.join(F.broadcast(cabeceras), on="pedido_id", how="inner")
def calcular_cesta_media(df):
"""Cesta media por mes y pais: importe total y numero de articulos."""
por_pedido = (
df.groupBy("pedido_id", "pais", F.trunc("fecha_pedido", "month").alias("mes"))
.agg(
F.sum("importe_linea").alias("importe_pedido"),
F.sum("cantidad").alias("articulos_pedido"),
F.countDistinct("sku").alias("skus_distintos"),
)
)
return (
por_pedido.groupBy("mes", "pais")
.agg(
F.count("*").alias("num_pedidos"),
F.round(F.avg("importe_pedido"), 2).alias("cesta_media_eur"),
F.round(F.expr("percentile_approx(importe_pedido, 0.5)"), 2)
.alias("cesta_mediana_eur"),
F.round(F.avg("articulos_pedido"), 2).alias("articulos_medios"),
F.round(F.avg("skus_distintos"), 2).alias("skus_medios"),
)
.orderBy("mes", "pais")
)F.broadcast() merece atención: le dice a Spark que replique el DataFrame pequeño en todos los executors en lugar de repartir ambos lados por la red. Es la misma optimización que la entrada lateral de Beam en 04-02, y es la diferencia entre un JOIN de segundos y uno de minutos. Solo funciona si el lado pequeño cabe en memoria (unos cientos de MB como mucho).
def calcular_productos_juntos(df, soporte_minimo=20):
"""Matriz de coocurrencia: que pares de SKU aparecen en el mismo pedido.
El algoritmo es un self-join del conjunto de pedidos consigo mismo,
con dos precauciones fundamentales de rendimiento.
"""
# 1) Un pedido puede tener el mismo SKU en varias lineas: nos quedamos
# con pares (pedido, sku) unicos para no contar dos veces.
pedido_sku = df.select("pedido_id", "sku").distinct()
# 2) PRECAUCION CRITICA: descartar pedidos con demasiadas lineas.
# Un pedido de 200 SKU genera 200*199/2 = 19.900 pares el solo,
# y esos pedidos corporativos raros dominarian el calculo y la memoria.
tam = pedido_sku.groupBy("pedido_id").agg(F.count("*").alias("n_skus"))
pedidos_validos = tam.filter((F.col("n_skus") >= 2) & (F.col("n_skus") <= 30))
base = pedido_sku.join(F.broadcast(pedidos_validos.select("pedido_id")),
on="pedido_id", how="inner")
izq = base.withColumnRenamed("sku", "sku_a")
der = base.withColumnRenamed("sku", "sku_b")
pares = (
izq.join(der, on="pedido_id")
# 3) sku_a < sku_b elimina los pares consigo mismo Y los duplicados
# invertidos: (A,B) se cuenta una vez, no dos como (A,B) y (B,A).
.filter(F.col("sku_a") < F.col("sku_b"))
.groupBy("sku_a", "sku_b")
.agg(F.count("*").alias("veces_juntos"))
.filter(F.col("veces_juntos") >= soporte_minimo)
)
# 4) Metricas de reglas de asociacion: soporte, confianza y lift.
conteo_sku = (base.groupBy("sku").agg(F.count("*").alias("veces_total")))
total_pedidos = base.select("pedido_id").distinct().count()
resultado = (
pares
.join(F.broadcast(conteo_sku.withColumnRenamed("sku", "sku_a")
.withColumnRenamed("veces_total", "total_a")),
on="sku_a")
.join(F.broadcast(conteo_sku.withColumnRenamed("sku", "sku_b")
.withColumnRenamed("veces_total", "total_b")),
on="sku_b")
.withColumn("soporte", F.col("veces_juntos") / F.lit(total_pedidos))
.withColumn("confianza_a_b", F.col("veces_juntos") / F.col("total_a"))
.withColumn("confianza_b_a", F.col("veces_juntos") / F.col("total_b"))
.withColumn(
"lift",
(F.col("veces_juntos") * F.lit(total_pedidos))
/ (F.col("total_a") * F.col("total_b")),
)
)
# 5) Top 5 acompanantes de cada producto, con funcion de ventana
ventana = Window.partitionBy("sku_a").orderBy(F.desc("lift"))
return (
resultado
.withColumn("puesto", F.row_number().over(ventana))
.filter(F.col("puesto") <= 5)
.select("sku_a", "sku_b", "veces_juntos",
F.round("soporte", 5).alias("soporte"),
F.round("confianza_a_b", 4).alias("confianza"),
F.round("lift", 3).alias("lift"),
"puesto")
)Cómo se interpretan las tres métricas, porque son las que Lucía llevará a la reunión:
| Métrica | Qué significa | Ejemplo AlpinaShop |
|---|---|---|
| Soporte | Proporción de pedidos que contienen ambos productos | 0,012 → el 1,2 % de los pedidos llevan crampones y piolet |
| Confianza | Si compra A, probabilidad de que compre B | 0,34 → un tercio de quien compra crampones compra piolet |
| Lift | Cuánto más probable es que vayan juntos frente al azar | 8,5 → ocho veces y media más de lo esperado: asociación fortísima |
El lift es el que hay que mirar. La confianza engaña con los productos superventas: si el 60 % de los pedidos llevan calcetines técnicos, cualquier producto tendrá alta confianza hacia ellos sin que exista relación real. El lift corrige por la popularidad de cada producto. Lift mayor que 1 indica asociación real; lift cercano a 1, independencia.
def main():
parser = argparse.ArgumentParser()
parser.add_argument("--fecha-desde", required=True)
parser.add_argument("--fecha-hasta", required=True)
parser.add_argument("--salida", default=f"{LAGO}/resultados/cesta/")
args = parser.parse_args()
spark = crear_sesion()
spark.sparkContext.setLogLevel("WARN")
df = cargar_lineas(spark, args.fecha_desde, args.fecha_hasta)
# cache(): el DataFrame se usa en DOS calculos distintos. Sin cache,
# Spark releeria y refiltraria todo el lago dos veces.
df.cache()
cesta = calcular_cesta_media(df)
(cesta.coalesce(1) # un solo fichero: el resultado es pequeno
.write.mode("overwrite")
.parquet(f"{args.salida}/cesta_media"))
juntos = calcular_productos_juntos(df)
(juntos.write.format("bigquery")
.option("table", f"{PROYECTO}.{DATASET}.productos_juntos")
.option("writeMethod", "direct")
.mode("overwrite")
.save())
df.unpersist()
print(f"Cesta media: {cesta.count()} filas")
print(f"Pares de productos: {juntos.count()} filas")
spark.stop()
if __name__ == "__main__":
main()Tres detalles de rendimiento que hay que interiorizar:
cache()materializa el DataFrame en memoria de los executors. Sin él, como Spark es perezoso, cada acción posterior recalcularía toda la cadena desde la lectura del lago. Con dos consumidores, ahorra la mitad del trabajo. Yunpersist()al terminar, para liberar memoria.coalesce(1)reduce a una sola partición antes de escribir. Es correcto solo para resultados pequeños; con datos grandes, concentrarlo todo en un executor lo tumbaría. Para volúmenes grandes se usarepartition(n).writeMethod=directusa la Storage Write API de BigQuery en lugar de pasar por ficheros temporales en el bucket. Es más rápido y evita gestionar limpieza.
Y la advertencia que corresponde: este análisis usa cliente_id y datos de pedido. Si en algún momento se incorporan datos personales identificables, el tratamiento debe ajustarse al RGPD y ser revisado por un profesional de compliance. Para el análisis de cesta no hace falta saber quién es el cliente, solo qué había en el pedido; mantenerlo así es minimización de datos por diseño. Todos los datos de este curso son ficticios.
- Acciones de inicialización y versiones de imagen
Una acción de inicialización es un script que se ejecuta en cada nodo al crearse el clúster. Sirve para instalar dependencias que no vienen en la imagen.
# Script propio en el bucket
cat > init-alpinashop.sh <<'EOF'
#!/bin/bash
set -euxo pipefail
# Bibliotecas de Python para el analisis de cesta
pip install --no-cache-dir mlxtend==0.23.1 pyarrow==16.1.0
# Solo en el maestro: utilidades de diagnostico
ROL=$(/usr/share/google/get_metadata_value attributes/dataproc-role)
if [[ "$ROL" == "Master" ]]; then
pip install --no-cache-dir jupyterlab-git
fi
EOF
gcloud storage cp init-alpinashop.sh gs://alpinashop-datalake/init/
gcloud dataproc clusters create alpinashop-spark-ml \
--region=europe-west1 --subnet=sn-datos-euw1 --no-address \
--service-account="$SA" \
--initialization-actions=gs://alpinashop-datalake/init/init-alpinashop.sh \
--initialization-action-timeout=10m \
--image-version=2.2-debian12 \
--max-idle=30mConsejos sobre las acciones de inicialización:
- El script se ejecuta en todos los nodos, incluidos los que añada el autoescalado. Debe ser idempotente y rápido: un script de cinco minutos multiplica por cinco el tiempo de arranque de cada nodo nuevo.
- Usa
get_metadata_value attributes/dataproc-rolepara distinguir maestro de worker. set -euxo pipefailhace que el script falle ruidosamente en lugar de dejar el nodo a medias.- Si las dependencias son muchas, es mejor construir una imagen personalizada que instalarlas en cada arranque.
Sobre las versiones de imagen: cada versión de Dataproc empaqueta un conjunto concreto de Spark, Hadoop, Python y el sistema operativo. La serie 2.2-debian12, por ejemplo, trae Spark 3.5 y Python 3.11.
| Práctica | Consecuencia |
|---|---|
| No indicar versión | Google elige la más reciente: tu trabajo puede romperse solo |
Indicar la serie (2.2-debian12) |
Actualizaciones menores automáticas dentro de la serie. Equilibrio razonable |
Indicar la versión exacta (2.2.28-debian12) |
Reproducibilidad total. Recomendado en producción crítica |
Para AlpinaShop: serie fijada en desarrollo, versión exacta en los flujos de trabajo programados. Y en el fichero del flujo, versionado en Git.
- Dataproc Serverless para Spark
Aunque un clúster efímero es mucho mejor que uno permanente, sigue habiendo que dimensionarlo: cuántos workers, qué máquina, cuánta memoria por executor. Dataproc Serverless elimina esa decisión: envías el trabajo de Spark y Google se encarga de todo.
gcloud dataproc batches submit pyspark \
gs://alpinashop-datalake/jobs/cesta_media.py \
--batch=cesta-media-$(date +%Y%m%d-%H%M%S) \
--region=europe-west1 \
--version=2.2 \
--service-account="$SA" \
--subnet=sn-datos-euw1 \
--deps-bucket=gs://alpinashop-datalake \
--properties="spark.executor.instances=4,\
spark.dynamicAllocation.enabled=true,\
spark.dynamicAllocation.maxExecutors=20" \
--labels=entorno=produccion,equipo=datos,centro-coste=analitica \
-- --fecha-desde=2025-01-01 --fecha-hasta=2026-03-31No hay clusters create. No hay --max-idle porque no hay nada que apagar. No hay --num-workers.
| Aspecto | Dataproc con clúster | Dataproc Serverless |
|---|---|---|
| Gestión | Creas y destruyes clústeres | Ninguna |
| Arranque | 90-120 s | 30-60 s |
| Dimensionado | Lo eliges tú | Automático |
| Facturación | Por VM y hora | Por DCU (unidades de cómputo de datos) mientras dura |
| Componentes | Todo el ecosistema: Hive, HBase, Presto, Jupyter | Solo Spark |
| Acciones de inicialización | Sí | No; se usan imágenes de contenedor personalizadas |
| VM Spot | Sí, muy barato | No aplicable |
| Sesión interactiva | Notebook en el clúster | Sesiones interactivas Serverless |
El criterio: si tu trabajo es Spark puro y no necesitas Hive ni HBase ni un clúster de larga vida, empieza por Serverless. Es menos que administrar y menos que olvidarse encendido. Usa clúster cuando necesites componentes del ecosistema, sesiones interactivas largas, o cuando el descuento de las VM Spot en cargas muy grandes compense la gestión.
Para AlpinaShop, la decisión razonada es: el análisis de cesta va a Serverless, porque es Spark puro, se ejecuta mensualmente y nadie quiere acordarse de apagar nada. El clúster alpinashop-spark se mantiene únicamente como entorno de exploración con Jupyter, con --max-idle=30m, y se destruye cuando no se use durante un mes.
- Dataproc, Serverless y Dataflow: la tabla de decisión
| Criterio | Dataproc (clúster) | Dataproc Serverless | Dataflow |
|---|---|---|---|
| Modelo de programación | Spark / Hadoop / Hive | Spark | Apache Beam |
| Lote | Sí | Sí | Sí |
| Streaming | Structured Streaming | Limitado | Su punto fuerte |
| Arranque | ~2 min | ~40 s | ~2 min (lote) |
| Infraestructura | La defines tú | Ninguna | Ninguna |
| Autoescalado | Con política | Automático | Automático |
| Coste en reposo | El clúster, si lo dejas | Cero | Cero (salvo streaming activo) |
| Semántica exactamente una vez | Hay que construirla | Igual | De serie |
| Ventanas y tiempo del evento | Manual | Manual | Modelo completo |
| Ecosistema de bibliotecas | Enorme (MLlib, pandas…) | Grande | Limitado a Beam |
| Portabilidad | Alta (Spark corre en todas partes) | Media | Alta (Beam tiene varios runners) |
| Elegir si… | Tienes código Spark, necesitas Hive/HBase, el equipo sabe Spark | Spark puro sin querer gestionar nada | Streaming, tiempo del evento, pipelines nuevos |
La decisión final para AlpinaShop queda así:
- Streaming de pedidos y visitas → Dataflow. El modelo de ventanas y marcas de agua de Beam no tiene equivalente cómodo en Spark, y las plantillas de Pub/Sub a BigQuery resuelven el caso base sin código.
- Cargas y transformaciones por lotes nuevas → Dataflow, por coherencia con lo anterior y porque el equipo ya lo tiene montado.
- Análisis algorítmico: cesta, coocurrencia, futuros modelos con MLlib → Dataproc Serverless.
- Exploración interactiva sobre el lago → Notebook en
alpinashop-spark, con--max-idle. - Transformación expresable en SQL sobre datos ya en BigQuery → BigQuery, sin mover nada.
- Notebooks y Spark SQL sobre BigQuery
Con --optional-components=JUPYTER y --enable-component-gateway, la consola de Dataproc muestra un enlace a JupyterLab protegido por IAM. Nada de abrir puertos ni túneles SSH: quien tenga el rol adecuado entra, y quien no, no.
El conector de BigQuery para Spark permite leer y escribir tablas directamente:
from pyspark.sql import SparkSession, functions as F
spark = (SparkSession.builder
.appName("exploracion-lucia")
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate())
# Leer una tabla completa: el conector usa la Storage Read API,
# que lee en paralelo y en formato columnar. No exporta a ficheros.
productos = (spark.read.format("bigquery")
.option("table", "alpinashop-datos.alpinashop_analitica.productos")
.load())
# MEJOR: delegar el filtro a BigQuery para traer menos datos
lineas = (spark.read.format("bigquery")
.option("table", "alpinashop-datos.alpinashop_analitica.lineas_pedido")
.option("filter", "fecha_pedido >= '2026-01-01'") # se ejecuta en BigQuery
.load())
# A partir de aqui, Spark SQL normal
lineas.createOrReplaceTempView("lineas")
productos.createOrReplaceTempView("productos")
resumen = spark.sql("""
SELECT p.categoria,
COUNT(DISTINCT l.pedido_id) AS pedidos,
ROUND(SUM(l.importe_linea), 2) AS ventas_eur
FROM lineas l
JOIN productos p ON p.sku = l.sku
GROUP BY p.categoria
ORDER BY ventas_eur DESC
""")
resumen.show(truncate=False)La opción filter es la que marca la diferencia: se ejecuta en BigQuery, no en Spark. Sin ella, el conector traería la tabla entera por la red para que Spark la filtrase, pagando la lectura completa en BigQuery y perdiendo tiempo. Es el mismo principio que vimos con EXTERNAL_QUERY en 04-01: filtra lo más cerca posible del origen.
Y la pregunta obligada: si puedes hacer esto en Spark SQL, ¿por qué no hacerlo en BigQuery directamente? Casi siempre deberías. El conector tiene sentido cuando el resultado del SQL alimenta un algoritmo de MLlib, cuando cruzas tablas de BigQuery con ficheros del lago que no están cargados, o cuando el código Spark ya existe. Para una agregación pura, BigQuery es más rápido y más barato.
- Migrar un Hadoop on-premise a Google Cloud
Este es el escenario para el que Dataproc está más justificado. Supón que AlpinaShop absorbe a un competidor con un clúster Hadoop de 30 nodos.
Qué se conserva:
- El código Spark, Hive y PySpark: en su inmensa mayoría funciona sin cambios.
- Las consultas de Hive: Dataproc incluye Hive, y el metastore puede migrarse.
- Los formatos de datos: Parquet, ORC y Avro son idénticos.
- Los flujos de Oozie, aunque conviene reemplazarlos por Composer (04-06).
Qué se replantea, obligatoriamente:
| Elemento on-premise | En Google Cloud | Por qué cambia |
|---|---|---|
| HDFS permanente | Cloud Storage | Es el cambio fundamental: sin él no hay clústeres efímeros ni ahorro |
| Un clúster gigante compartido | Varios clústeres pequeños por carga | Cada equipo con su versión y su presupuesto; sin colas ni vecinos ruidosos |
| Dimensionado para el pico | Autoescalado + Spot | Se pagaba el pico las 24 horas |
| Kerberos | IAM y cuentas de servicio | Modelo de identidad de la nube (03-04) |
| Metastore de Hive local | Dataproc Metastore gestionado | Sobrevive a los clústeres efímeros: es la pieza que lo hace posible |
| Oozie / cron | Cloud Composer o Workflows | 04-06 |
| Impala / Presto para consultas | BigQuery | Suele ser el mayor salto de rendimiento y de simplicidad |
| Flume / Kafka de ingesta | Pub/Sub (04-04) o Managed Kafka | Gestionado |
La estrategia recomendada, y la única que suele salir bien, es por fases:
- Copiar los datos a Cloud Storage con el Storage Transfer Service o
hadoop distcp, sin tocar nada más. El clúster on-premise sigue funcionando. - Levantar Dataproc Metastore y registrar las tablas apuntando a
gs://en lugar dehdfs://. - Ejecutar los trabajos existentes en Dataproc contra los datos ya en Cloud Storage, comparando resultados con los del clúster antiguo. Esta fase de doble ejecución es innegociable: es la única forma de demostrar que los números coinciden.
- Apagar el clúster on-premise cuando la comparación cuadre durante varias semanas.
- Solo entonces, optimizar: pasar las consultas de Hive a BigQuery, los trabajos de streaming a Dataflow, adoptar Serverless.
El error clásico es intentar el paso 5 a la vez que el 3, es decir, migrar y modernizar simultáneamente. Cuando los números no cuadran, no se sabe si es por la migración o por la reescritura, y el proyecto se atasca durante meses.
- Coste y VM Spot
Dataproc factura dos cosas:
- Las VM subyacentes (Compute Engine, discos y red), a tarifa normal.
- Una tarifa de gestión de Dataproc, del orden de 0,01 $ por vCPU y hora, verificable en la documentación oficial.
Es decir, el sobrecoste de Dataproc frente a montar Hadoop tú mismo en VM es pequeño: pagas poco por no administrar nada.
Ejemplo con alpinashop-spark (1 maestro + 2 workers + 2 secundarios, todos n2-standard-4, 20 vCPU en total), como orden de magnitud:
| Escenario | Horas al mes | Coste aproximado |
|---|---|---|
| Clúster permanente 24×7 | 720 | Del orden de 1.400 € |
Clúster con --max-idle=30m, 2 h de uso al día |
~75 | Del orden de 150 € |
| Clúster efímero por flujo, 20 min al día | ~10 | Del orden de 20 € |
| Con secundarios Spot en lugar de estándar | ~10 | Del orden de 12 € |
| Dataproc Serverless, mismo trabajo mensual | ~1 | Céntimos |
Verifica los precios vigentes en la documentación oficial; lo que importa aquí es el factor 100 entre la primera fila y la tercera. Ese factor es la lección entera.
Palancas de ahorro, en orden de impacto:
--max-idlesiempre. Un clúster olvidado un puente de cuatro días cuesta más que un año de Serverless.- Clústeres efímeros por flujo de trabajo. Elimina el problema de raíz.
- Workers secundarios Spot. Descuentos de hasta el 80 %, con el matiz de que no aportan HDFS y pueden desaparecer.
- Serverless para trabajos ocasionales.
- Discos ajustados. Con los datos en Cloud Storage, HDFS es solo espacio de shuffle: 200 GB por worker sobran para casi todo.
- Región coherente. Clúster y buckets en
europe-west1: leer datos de otra región cuesta salida y latencia. - Descuentos por uso comprometido solo si acabas teniendo un clúster permanente, cosa que este apartado sugiere evitar.
- Formatos columnares. Parquet con ficheros de 128-512 MB lee mucho menos y evita el problema de los ficheros pequeños, que es el mayor asesino de rendimiento en Spark sobre almacenamiento de objetos.
Errores Comunes y Consejos
Crear un clúster permanente por costumbre. Es la herencia mental del Hadoop on-premise y es el error más caro. Si el clúster no tiene --max-idle, no debería existir.
Escribir datos importantes en HDFS. Desaparecen al destruir el clúster. HDFS en Dataproc es memoria de trabajo, no almacenamiento.
No fijar la versión de imagen. Google actualiza la versión por defecto y un trabajo que funcionaba deja de funcionar sin que nadie haya tocado el código.
Usar collect() sobre un DataFrame grande. Trae todos los datos al driver, que es una sola máquina. Es la causa número uno de OutOfMemoryError en Spark. Usa show(), take(n) o escribe a un fichero.
Olvidar cache() cuando un DataFrame se usa varias veces. Spark recalcula toda la cadena en cada acción. Y el error inverso: cachear todo lo que se mueve, hasta llenar la memoria y provocar volcado a disco.
Dejar spark.sql.shuffle.partitions en 200. Con datos pequeños genera 200 tareas minúsculas con más sobrecarga que trabajo; con datos grandes, particiones enormes que no caben en memoria.
Muchos ficheros pequeños en el lago. Diez mil ficheros de 1 MB son mucho más lentos de leer que veinte de 500 MB, porque cada apertura tiene latencia. Compacta.
Poner los datos en una región y el clúster en otra. Salida de datos entre regiones facturada y latencia añadida en cada lectura.
Consejo: usa Serverless por defecto. Empieza por ahí y crea clúster solo cuando descubras que necesitas algo que Serverless no da. Es el camino con menos deuda operativa.
Consejo: --dry-run no existe, pero el subconjunto sí. Antes de lanzar un trabajo sobre dos años, ejecútalo sobre una semana. Los errores de lógica aparecen igual y cuestan cien veces menos.
Consejo: mira siempre la interfaz de Spark. El component gateway da acceso a la UI de Spark, donde se ven las etapas, las tareas y —lo más útil— la distribución de tiempos entre tareas. Si una tarea tarda cien veces más que la mediana, tienes sesgo de datos, exactamente igual que en Dataflow.
Ejercicios
Ejercicio 1: clúster efímero con autodestrucción
Crea un clúster llamado alpinashop-spark-pruebas en europe-west1 que: viva en sn-datos-euw1 sin IP pública, use la cuenta de servicio sa-dataproc-analitica, tenga 1 maestro n2-standard-2 y 2 workers n2-standard-2, añada 2 workers secundarios Spot, fije la imagen 2.2-debian12, se autodestruya tras 15 minutos de inactividad y en cualquier caso a las 2 horas de vida, y lleve las etiquetas estándar de AlpinaShop. Después comprueba su estado y bórralo explícitamente.
Ejercicio 2: trabajo PySpark de devoluciones
Escribe un trabajo PySpark que lea gs://alpinashop-datalake/pedidos/cabeceras/ y gs://alpinashop-datalake/pedidos/lineas/ en Parquet, y calcule por categoría de producto: número de pedidos con al menos una devolución, tasa de devolución sobre el total de pedidos de esa categoría, e importe medio devuelto. El resultado debe escribirse en alpinashop-datos.alpinashop_analitica.devoluciones_categoria. Aplica al menos dos optimizaciones vistas en la lección y explica por qué las aplicas.
Ejercicio 3: decidir la herramienta
Para cada uno de estos cinco encargos de AlpinaShop, elige entre BigQuery, Dataflow, Dataproc con clúster, Dataproc Serverless o bq load, y justifica la elección en dos o tres frases:
- Cargar cada noche un fichero Parquet de 4 GB del ERP en una tabla de BigQuery, sin transformación alguna.
- Calcular el ranking mensual de ventas por categoría a partir de tablas que ya están en
alpinashop_analitica. - Procesar los eventos de
pedidos-nuevossegún llegan, agrupándolos por hora del evento y tolerando 90 minutos de retraso. - Entrenar cada semana un modelo de recomendación con ALS de MLlib sobre dos años de historial de compras.
- Ejecutar el código PySpark heredado del competidor absorbido, 15.000 líneas que usan Hive y funciones definidas por el usuario, mientras se decide qué hacer con él.
Soluciones
Solución 1
SA="[email protected]"
gcloud dataproc clusters create alpinashop-spark-pruebas \
--region=europe-west1 \
--zone=europe-west1-b \
--subnet=sn-datos-euw1 \
--no-address \
--service-account="$SA" \
--master-machine-type=n2-standard-2 \
--master-boot-disk-size=100GB \
--num-workers=2 \
--worker-machine-type=n2-standard-2 \
--worker-boot-disk-size=100GB \
--num-secondary-workers=2 \
--secondary-worker-type=spot \
--image-version=2.2-debian12 \
--max-idle=15m \
--max-age=2h \
--labels=entorno=desarrollo,equipo=datos,centro-coste=analitica,aplicacion=analitica# Comprobar el estado
gcloud dataproc clusters describe alpinashop-spark-pruebas --region=europe-west1 \
--format="yaml(status.state, config.lifecycleConfig, config.gceClusterConfig.internalIpOnly)"
gcloud dataproc clusters list --region=europe-west1 \
--format="table(clusterName, status.state, config.softwareConfig.imageVersion)"
# Borrado explicito, sin esperar a max-idle
gcloud dataproc clusters delete alpinashop-spark-pruebas --region=europe-west1 --quietLas dos opciones de ciclo de vida son complementarias y conviene poner las dos:
--max-idle=15m: se destruye si nadie lo usa durante 15 minutos. Cubre el caso "me he ido a comer".--max-age=2h: se destruye a las 2 horas pase lo que pase. Cubre el caso "un trabajo en bucle infinito mantiene el clúster ocupado y--max-idlenunca se dispara". Es el fallo que produce las facturas de fin de semana.
--no-address con --subnet=sn-datos-euw1 exige que Private Google Access esté activado en esa subred (03-01) o los nodos no podrán descargar paquetes ni hablar con las API. Si el clúster se queda en CREATING y luego falla, esa es la primera causa a revisar.
Solución 2
"""devoluciones.py -- Tasa de devolucion por categoria."""
from pyspark.sql import SparkSession, functions as F
PROYECTO, DATASET = "alpinashop-datos", "alpinashop_analitica"
LAGO = "gs://alpinashop-datalake"
spark = (SparkSession.builder
.appName("alpinashop-devoluciones")
.config("temporaryGcsBucket", "alpinashop-datalake")
.getOrCreate())
# OPTIMIZACION 1: seleccionar solo las columnas necesarias en la lectura.
# Parquet es columnar: las columnas no pedidas no se leen del bucket.
cabeceras = (spark.read.parquet(f"{LAGO}/pedidos/cabeceras/")
.select("pedido_id", "estado", "fecha_pedido", "total_pedido"))
lineas = (spark.read.parquet(f"{LAGO}/pedidos/lineas/")
.select("pedido_id", "sku", "importe_linea"))
# El catalogo se lee de BigQuery, no del lago
productos = (spark.read.format("bigquery")
.option("table", f"{PROYECTO}.{DATASET}.productos")
.option("filter", "activo = true")
.load()
.select("sku", "categoria"))
# OPTIMIZACION 2: broadcast del catalogo (miles de filas, cabe en memoria).
# Evita por completo el shuffle del JOIN con la tabla de lineas, que es
# la grande. Sin esto, ambos lados se reparticionarian por 'sku'.
lineas_cat = lineas.join(F.broadcast(productos), on="sku", how="inner")
# Un pedido devuelto lo esta entero: marcamos a nivel de cabecera
cab = cabeceras.withColumn(
"es_devuelto", F.when(F.col("estado") == "devuelto", 1).otherwise(0)
)
# OPTIMIZACION 3: cache, porque el DataFrame unido se usa dos veces
detalle = lineas_cat.join(F.broadcast(cab), on="pedido_id", how="inner")
detalle.cache()
# Pedidos distintos por categoria (un pedido puede tocar varias categorias)
por_categoria = (
detalle.groupBy("categoria")
.agg(
F.countDistinct("pedido_id").alias("pedidos_totales"),
F.countDistinct(
F.when(F.col("es_devuelto") == 1, F.col("pedido_id"))
).alias("pedidos_devueltos"),
F.round(
F.avg(F.when(F.col("es_devuelto") == 1, F.col("importe_linea"))), 2
).alias("importe_medio_devuelto"),
)
.withColumn(
"tasa_devolucion_pct",
F.round(100 * F.col("pedidos_devueltos") / F.col("pedidos_totales"), 2),
)
.orderBy(F.desc("tasa_devolucion_pct"))
)
(por_categoria.write.format("bigquery")
.option("table", f"{PROYECTO}.{DATASET}.devoluciones_categoria")
.option("writeMethod", "direct")
.mode("overwrite")
.save())
por_categoria.show(truncate=False)
detalle.unpersist()
spark.stop()Las optimizaciones y su justificación:
- Proyección temprana de columnas. Parquet es columnar; pedir cuatro columnas en lugar de veinte reduce proporcionalmente los bytes leídos del bucket y la memoria de los executors. Es el equivalente exacto de no hacer
SELECT *en BigQuery. broadcast()en los dosJOIN. El catálogo de productos son miles de filas y las cabeceras del periodo también son pequeñas comparadas con las líneas. Difundirlas evita repartir por la red la tabla grande, que es el coste dominante.cache()sobredetalle. Aunque en esta versión final solo hay una agregación, en cuanto se añada un segundo cálculo (por país, por mes) Spark recalcularía toda la cadena. Es la preparación correcta; conunpersist()al terminar para no retener memoria.- Filtro delegado a BigQuery (
option("filter", "activo = true")). Se ejecuta allí y llegan menos filas.
Nota metodológica que hay que explicitar al presentar el resultado: un pedido con productos de tres categorías cuenta como pedido en las tres, así que las cifras por categoría no suman el total de pedidos. Es correcto para medir tasa por categoría, pero hay que decirlo en el informe o alguien restará y no le cuadrará.
Solución 3
1. Cargar un Parquet de 4 GB cada noche sin transformación → bq load.
No hay transformación, luego no hace falta motor de procesamiento. La carga por lotes en BigQuery es gratuita, Parquet lleva el esquema incorporado y no hay que dimensionar nada. Usar Dataflow o Spark aquí sería pagar cómputo por hacer una copia. Un solo comando, orquestado en 04-06.
2. Ranking mensual sobre tablas ya en alpinashop_analitica → BigQuery.
Los datos ya están ahí y la operación es una agregación con función de ventana, exactamente lo que hicimos en 04-01 con RANK() OVER y QUALIFY. Sacar los datos de BigQuery para procesarlos fuera y volver a meterlos es el antipatrón clásico: coste de lectura, coste de cómputo, coste de escritura y latencia, para obtener un resultado peor. Si hay que refrescarlo a menudo, vista materializada.
3. Eventos de pedidos-nuevos por hora del evento con 90 minutos de tolerancia → Dataflow.
Es literalmente el caso de uso para el que existe el modelo de Beam: streaming no acotado, agrupación por tiempo del evento, ventanas fijas, marca de agua y allowed_lateness. Structured Streaming de Spark podría, pero con un modelo de tiempo menos expresivo y sin la integración nativa con Pub/Sub. Además Dataflow da semántica de exactamente una vez hacia BigQuery sin trabajo adicional.
4. Entrenar ALS de MLlib semanalmente → Dataproc Serverless. ALS es un algoritmo iterativo distribuido implementado en MLlib, sin equivalente en Beam ni en SQL puro. Es Spark del principio al fin. Serverless en lugar de clúster porque se ejecuta una vez por semana: nadie quiere mantener ni acordarse de apagar un clúster que trabaja una hora cada siete días. (El paso siguiente, servir ese modelo en producción, es territorio de Vertex AI en el módulo 5.)
5. 15.000 líneas de PySpark heredadas con Hive y UDF → Dataproc con clúster. Aquí manda la restricción práctica: el código existe, funciona y usa Hive, que Serverless no incluye. La prioridad es que siga funcionando con el mínimo cambio, así que clúster de Dataproc con Dataproc Metastore, ejecutando el código tal cual contra los datos ya copiados a Cloud Storage. La modernización —pasar consultas a BigQuery, streaming a Dataflow— es una fase posterior y separada, nunca simultánea a la migración: si los números no cuadran, hay que poder saber si es por el traslado o por la reescritura.
Conclusión
Has recorrido el mundo de Hadoop y Spark con la perspectiva justa: qué son, qué problema resolvieron —procesar más datos de los que caben en una máquina, con máquinas que fallan—, por qué Spark desplazó a MapReduce manteniendo los datos intermedios en memoria, y por qué en 2026 siguen importando aunque existan BigQuery y Dataflow: hay código escrito, hay personas que lo saben usar, hay bibliotecas como MLlib sin equivalente, y hay un ecosistema abierto que da portabilidad real.
Has visto qué aporta Dataproc: clústeres en menos de dos minutos, configurados y coherentes, integrados con Cloud Storage, BigQuery, IAM y Logging. Conoces la anatomía —maestro, workers primarios que sostienen HDFS, workers secundarios que solo aportan cómputo y por eso pueden ser Spot— y has creado alpinashop-spark dentro de sn-datos-euw1, sin IP pública, con la cuenta sa-dataproc-analitica, con la imagen fijada, con Jupyter accesible por el component gateway y con una política de autoescalado que solo escala los secundarios y los retira con elegancia.
Pero lo importante de esta lección no es un comando, es una inversión conceptual: el clúster efímero sobre almacenamiento en Cloud Storage. Como la red interna hace que leer de gs:// sea comparable a leer de disco local, ya no hace falta que el dato viva en el clúster. Y si el dato no vive en el clúster, el clúster puede morir. De ahí salen --max-idle, --max-age, los flujos de trabajo con clúster gestionado que nace y muere con el job, y el factor cien de diferencia en la factura entre un clúster permanente y uno efímero.
Has escrito el trabajo PySpark que responde a la pregunta de Lucía: la cesta media por mes y país con su mediana al lado, y la matriz de productos comprados juntos con soporte, confianza y lift —la métrica que corrige por la popularidad y evita concluir que todo el mundo compra calcetines con todo—, aplicando broadcast para evitar shuffle, cache para no recalcular, un tope de líneas por pedido para que los pedidos corporativos raros no dominen el cálculo, y el truco de sku_a < sku_b para contar cada par una sola vez. Conoces las acciones de inicialización, sus riesgos y por qué fijar la versión de imagen no es una manía.
Has conocido Dataproc Serverless, que elimina incluso la decisión de dimensionar, y has fijado el criterio de AlpinaShop: el análisis de cesta a Serverless, el clúster solo como entorno de exploración con Jupyter y autodestrucción. Y tienes la tabla de decisión completa entre Dataproc, Serverless y Dataflow, con streaming y tiempo del evento del lado de Beam, algoritmos y ecosistema del lado de Spark, y SQL sobre datos ya cargados del lado de BigQuery, sin mover nada. Sabes cómo se migra un Hadoop on-premise por fases, con la regla de oro de no modernizar y migrar a la vez. Y sabes dónde está el dinero: --max-idle, clústeres efímeros, Spot, Serverless, discos ajustados y ficheros grandes en formato columnar.
Queda una promesa pendiente desde hace dos lecciones. Todo lo que has construido —el pipeline de Dataflow, el trabajo de Spark, las tablas de BigQuery— funciona sobre datos que ya están en algún sitio. Pero la tienda sigue siendo una isla: cuando un cliente confirma un pedido, la aplicación Flask tiene que avisar al almacén para que lo prepare, a facturación para que emita la factura, al servicio de correo para la confirmación, y ahora también a la analítica. Si lo hace llamando a los cuatro uno detrás de otro, la venta se queda colgada esperando al más lento, y si uno falla, no está claro qué ha ocurrido con los otros tres. Es un diseño frágil que se rompe justo el día de más ventas del año.
En 04-04, Cloud Pub/Sub, romperemos ese acoplamiento. Crearemos por fin el topic pedidos-nuevos y las suscripciones sub-almacen, sub-facturacion y sub-analitica; entenderás las garantías reales de la mensajería —entrega al menos una vez, orden no garantizado— y por qué eso obliga a que tus consumidores sean idempotentes; verás los temas de mensajes fallidos, los reintentos con retroceso exponencial, la reproducción de mensajes con seek, los filtros por atributo y las suscripciones directas a BigQuery que ingieren sin escribir una línea de código. Y por fin conectaremos de verdad las notificaciones del bucket alpinashop-catalogo que quedaron prometidas en 02-02.
Curso de Google Cloud Platform (GCP)
Módulo 1: Introducción a Google Cloud Platform
- ¿Qué es Google Cloud Platform?
- Configuración de tu cuenta de GCP
- Descripción general de la consola de GCP
- Proyectos, jerarquía de recursos y facturación
- Regiones, zonas y modelo de responsabilidad compartida
- Cloud Shell y la CLI de gcloud
Módulo 2: Servicios principales de GCP
- Compute Engine: máquinas virtuales en Google Cloud
- Cloud Storage: almacenamiento de objetos
- Cloud SQL: bases de datos relacionales gestionadas
- App Engine: plataforma como servicio
- Google Kubernetes Engine (GKE)
- Bases de datos NoSQL: Firestore, Bigtable y Spanner
- Cómo elegir el servicio de cómputo adecuado
Módulo 3: Redes y seguridad
- Redes VPC
- Balanceo de carga en la nube
- Cloud CDN
- Gestión de identidad y acceso (IAM)
- Cloud Armor
- Secretos y cifrado: Secret Manager y Cloud KMS
- Cloud DNS, certificados TLS y publicación segura de servicios
Módulo 4: Datos y análisis
- BigQuery: el almacén de datos analítico
- Cloud Dataflow: procesamiento de datos por lotes y en streaming
- Cloud Dataproc: Spark y Hadoop gestionados
- Cloud Pub/Sub: mensajería asíncrona
- Cloud Data Fusion: integración de datos sin código
- Orquestación de pipelines con Cloud Composer y Workflows
- Gobierno del dato y cuadros de mando con Dataplex y Looker Studio
Módulo 5: Aprendizaje automático e IA
- Vertex AI: la plataforma de machine learning de GCP
- AutoML: modelos a medida sin escribir código
- TensorFlow en GCP: entrenamiento y servicio de modelos
- API de lenguaje natural
- API de visión
- IA generativa en Vertex AI: modelos Gemini y embeddings
- MLOps: del modelo al producto con Vertex AI Pipelines
Módulo 6: DevOps y monitoreo
- Cloud Build: integración continua en GCP
- Cloud Source Repositories y gestión del código fuente
- Cloud Functions: funciones sin servidor
- Cloud Monitoring (antes Stackdriver): métricas, paneles y alertas
- Cloud Deployment Manager e infraestructura como código nativa
- Cloud Logging y Cloud Trace: logs, trazas y diagnóstico
- Terraform en GCP: infraestructura como código en la práctica
Módulo 7: Temas avanzados de GCP
- Híbrido y multinube con Anthos
- Computación sin servidor con Cloud Run
- Redes avanzadas: VPC compartida, peering y conectividad híbrida
- Mejores prácticas de seguridad
- Gestión y optimización de costos
- Fiabilidad: SLO, alta disponibilidad y recuperación ante desastres
- Gobierno a escala: organización, políticas y auditoría
