Logo GCP con Eduardo

GCP con Eduardo

BigQuery ya ejecuta Spark: Fin del ETL Tradicional y Data Lakes Fragmentados
✦ Guía Técnica & Arquitectura

BigQuery ya ejecuta Spark: lo que cambia de verdad para los Data Engineers

Durante más de una década, la ingeniería de datos ha estado fracturada en dos mundos incompatibles: el Data Warehouse gobernado y estructurado con SQL (BigQuery), y el Data Lake flexible y distribuido en Python/Spark. La llegada de Apache Spark Serverless nativo en BigQuery Studio destruye esta frontera técnica.

Lo que aprenderás en esta guía

Cero Movimiento de Datos: Lectura directa mediante BigQuery Storage API a máxima velocidad gRPC.
Serverless Real: Olvídate de aprovisionar clusters en Dataproc o pelear con infraestructura Terraform.
De SQL a Data Science: Uso conjunto de PySpark, transformaciones matriciales y linaje unificado.
FinOps & Trade-offs: Cuándo usar Spark Serverless y cuándo mantener clusters dedicados.
Google BigQuery Studio Dataproc Serverless Apache Spark 3.x PySpark BigQuery Storage Read API Google Cloud Platform

La pesadilla de los pipelines híbridos clásicos

Imagina un caso habitual: tu equipo de negocio o ciencia de datos necesita cruzar 50 TB de datos para calcular la rotación de stock crítico, valoración financiera de inventarios y cuota de catálogo. SQL estándar se queda corto para analítica avanzada de machine learning o manipulación multidimensional.

Hasta ahora, la solución implicaba un proceso frágil y costoso:

  1. Exportar tablas desde BigQuery hacia buckets de Cloud Storage en ficheros Parquet o AVRO.
  2. Aprovisionar un clúster dedicado de Cloud Dataproc (esperando 5-10 minutos de arranque y configurando redes VPC/IAM).
  3. Leer los ficheros desde Spark, ejecutar los cálculos en memoria distribuida.
  4. Escribir los resultados nuevamente en GCS.
  5. Ejecutar un comando LOAD de vuelta en BigQuery para crear los Data Marts.

¿El resultado? Horas de desfase temporal, costes duplicados de almacenamiento, latencia de red y un equipo de ingeniería perdiendo el 80% de su tiempo en mantenimiento operativo en lugar de extraer valor de los datos.

Matriz de Decisión Arquitectónica: ¿Qué motor utilizar?

La integración de Spark en BigQuery no sustituye a todos los motores, sino que redefine los límites de aplicabilidad técnica:

Patrón de Carga Motor Óptimo Latencia de Arranque Gobernanza / IAM Modelo de Coste
Analítica SQL Directa & BI BigQuery SQL Engine Sub-segundo Nativo BigQuery Slots / Bytes Escaneados
ETL Esporádico / ML Ligero BigQuery Studio + Spark Serverless ~30 - 60 seg (Cold Start) Unificado en el Dataset DCU por segundo (Pay-as-you-go)
Procesamiento Masivo 24/7 Dataproc Dedicado (GCE / GKE) Inmediato (Pre-warm) Service Account por Nodo / VPC Cómputo VM continuo + Reserva
Micro-servicios / Ingesta Low Latency Cloud Run / Dataflow Streaming Milisegundos IAM granular por Endpoint CPU/Memoria por petición o Worker
⚠️ Antipatrón FinOps Crítico: La trampa del Spark Serverless continuo

Spark Serverless cobra por Data Compute Units (DCUs) exactas durante el tiempo de ejecución de la sesión. Si utilizas este patrón para procesar micro-batches continuos cada 30 segundos, el coste de inicialización y la tarifa serverless superará con creces el coste de un clúster Dataproc reservado o una canalización nativa con Dataflow. Reserva Spark sobre BigQuery para cargas por lotes, Data Marts intermedios y prototipado rápido de notebooks.

Implementación Práctica: Pipeline Multi-Mart en BigQuery Studio

A continuación se detalla la implementación modular de un notebook en BigQuery Studio utilizando DataprocSparkSession con el conector optimizado direct para leer y escribir tres Data Marts sin salir de la consola.

inventory_analytics_pipeline.py Python 3.11 / PySpark 3.x
from google.cloud.dataproc_sparksession import DataprocSparkSession
from pyspark.sql import functions as F

# 1. Inicialización de la sesión Serverless conectada al catálogo de BigQuery
spark = DataprocSparkSession.builder.getOrCreate()

PROJECT_ID = "mi-empresa-analytics"
SRC_DATASET = "retail_raw_data"
DEST_DATASET = "retail_datamarts_prod"

# 2. Lectura directa mediante BigQuery Storage API (Zero Egress GCS)
df_inventory = spark.read.format("bigquery") \
    .option("table", f"{PROJECT_ID}:{SRC_DATASET}.fact_inventory") \
    .load()

df_products = spark.read.format("bigquery") \
    .option("table", f"{PROJECT_ID}:{SRC_DATASET}.dim_products") \
    .load() \
    .withColumnRenamed("name", "product_name")

df_stores = spark.read.format("bigquery") \
    .option("table", f"{PROJECT_ID}:{SRC_DATASET}.dim_stores") \
    .load() \
    .withColumnRenamed("name", "store_name")

# 3. Master Join en memoria distribuida
df_joined = df_inventory \
    .join(df_products, on="product_id", how="inner") \
    .join(df_stores, on="store_id", how="inner")

# --- DATA MART 1: Valoración de Inventario por Tienda ---
df_store_valuation = df_joined.groupBy("store_id", "store_name") \
    .agg(
        F.sum(F.col("quantity") * F.col("price")).alias("total_inventory_value"),
        F.sum("quantity").alias("total_items_count")
    ) \
    .orderBy(F.col("total_inventory_value").desc())

df_store_valuation.write.format("bigquery") \
    .option("table", f"{PROJECT_ID}:{DEST_DATASET}.dm_store_valuation") \
    .mode("overwrite") \
    .save()

# --- DATA MART 2: Alertas de Quiebre de Stock Crítico (< 5 unidades) ---
df_low_stock = df_joined.filter(F.col("quantity") < 5) \
    .select("store_name", "product_name", "quantity", "price") \
    .orderBy("quantity", "store_name")

df_low_stock.write.format("bigquery") \
    .option("table", f"{PROJECT_ID}:{DEST_DATASET}.dm_low_stock_alerts") \
    .mode("overwrite") \
    .save()

# --- DATA MART 3: Cuota de Valor y Catálogo por Marca ---
df_brand_share = df_joined.groupBy("brand") \
    .agg(
        F.sum(F.col("quantity") * F.col("price")).alias("brand_total_value"),
        F.countDistinct("product_id").alias("unique_products")
    ) \
    .orderBy(F.col("brand_total_value").desc())

df_brand_share.write.format("bigquery") \
    .option("table", f"{PROJECT_ID}:{DEST_DATASET}.dm_brand_share") \
    .mode("overwrite") \
    .save()

Patrones de Diseño y Rendimiento en BigQuery Studio

1. Mecánica del conector directo (gRPC + Apache Arrow)

El parámetro spark.read.format("bigquery") interactúa con la API de almacenamiento de BigQuery. Los flujos de datos son decodificados directamente en vectores de memoria Arrow compatibles con Spark, eliminando la serialización de registros fila por fila.

2. Predicate Pushdown y Column Pruning

Asegúrate de aplicar los métodos .select() y .filter() lo antes posible en tu grafo de ejecución de Spark. El conector traslada estas condiciones al motor interno de BigQuery, de modo que solo se transfieren por red las columnas y particiones estrictamente requeridas.

3. Gobernanza Unificada con Dataplex

Al no generar copias dispersas en buckets de Cloud Storage, mantienes las políticas de acceso a nivel de columna, enmascaramiento de datos sensibles (PII) y linaje corporativo totalmente centralizados dentro del catálogo de BigQuery y Dataplex.

Framework de Implementación Empresarial Paso a Paso

  1. Habilitación de APIs: Activa bigquerystudio.googleapis.com y dataproc.googleapis.com en tu proyecto de GCP.
  2. Configuración de Permisos IAM: Otorga a los Data Engineers los roles de BigQuery Studio User, Dataproc Serverless Worker y BigQuery Data Editor en los datasets destino.
  3. Creación del Dataset de Destino: Asegúrate de que el Dataset destino resida en la misma región multirregional (ej. US o EU) que las tablas de origen para evitar latencias de cruce regional.
  4. Despliegue del Notebook: Abre BigQuery Studio, selecciona Nuevo Notebook con Spark, vincula el entorno de ejecución Serverless y ejecuta las celdas de forma interactiva.

Preguntas Frecuentes (FAQ)

¿Cuándo conviene usar Spark Serverless en BigQuery en lugar de Dataproc tradicional?

Spark Serverless en BigQuery Studio es óptimo para cargas de trabajo analíticas ad-hoc, pipelines de machine learning ligero y procesos batch esporádicos. Si tienes un cluster ejecutándose 24/7 con alta densidad de jobs continuos, Dataproc con instancias reservadas sigue ofreciendo un coste unitario más predecible.

¿Cómo accede Spark a las tablas de BigQuery sin mover datos por Cloud Storage?

Utiliza la BigQuery Storage Read/Write API basada en gRPC y formato columnar Apache Arrow. Los ejecutores de Spark acceden a los streams de almacenamiento directamente en paralelo con pushdown de predicados y proyecciones, evitando el paso intermedio de exportar archivos Parquet/AVRO a Cloud Storage.

¿Qué latencia de arranque tiene una sesión de Spark Serverless en BigQuery Studio?

El tiempo de arranque en frío para inicializar la sesión serverless de Dataproc Spark ronda entre los 30 y 90 segundos. Una vez instanciada la sesión, la ejecución de las celdas y queries Spark se realiza en segundos.

✦ Sobre el Autor

Eduardo Martínez Agrelo

AI & Data Architect

Especialista en modernización de plataformas de datos y arquitecturas analíticas en Google Cloud Platform. Ayudo a organizaciones a diseñar ecosistemas escalables de Big Data, IA generativa y optimización FinOps en la nube.