GCP Data Pipeline: arquitectura de datos en Google Cloud con Cloud Functions, Pub/Sub, Dataproc, Composer y Firestore
Diseñar un GCP Data Pipeline real no consiste simplemente en conectar servicios de Google Cloud. El reto arquitectónico está en decidir dónde ingestar, cómo desacoplar productores y consumidores, cuándo procesar con Apache Spark, cómo orquestar dependencias y dónde servir finalmente los datos. En esta guía construimos una arquitectura end-to-end orientada a producción utilizando Cloud Functions, Pub/Sub, Cloud Storage, Dataproc Spark, Cloud Composer y Firestore.
Lo que aprenderás
Arquitectura de datos en Google Cloud: del evento al dato procesado
Una arquitectura robusta separa las responsabilidades. La función de ingesta no debería asumir el trabajo de procesamiento pesado; el sistema de mensajería no debería convertirse en una base de datos; y el orquestador debería coordinar el workflow sin acoplarse innecesariamente a la lógica de negocio.
Productor
Functions
Storage
Spark
Por encima de este flujo se sitúa Cloud Composer, que utiliza Apache Airflow para coordinar procesos batch, validar disponibilidad de datos, lanzar jobs de Spark, gestionar dependencias y controlar reintentos.
Este diseño introduce una propiedad fundamental: desacoplamiento. Cada componente puede escalar, fallar y evolucionar de forma relativamente independiente.
Matriz de decisión: ¿Cloud Functions, Pub/Sub, Dataproc o Firestore?
Uno de los errores más habituales en Data Engineering consiste en seleccionar el servicio por familiaridad en lugar de hacerlo por patrón de carga. Una función serverless, un sistema de mensajería, un motor distribuido y una base NoSQL resuelven problemas completamente diferentes.
| Servicio / patrón | Rol | Latencia | Escalabilidad | Coste | Uso recomendado |
|---|---|---|---|---|---|
| Cloud Functions | Ejecución serverless orientada a eventos | Baja | Automática | Variable según invocaciones y recursos | Ingesta HTTP, validaciones y transformaciones ligeras |
| Pub/Sub | Mensajería y desacoplamiento | Muy baja | Muy alta | Basado principalmente en volumen de datos | Eventos, colas, integración entre productores y consumidores |
| Cloud Storage | Data Lake / staging | No orientado a serving transaccional | Muy alta | Generalmente eficiente para almacenamiento masivo | Raw data, Avro, Parquet, backups y datasets intermedios |
| Dataproc + Spark | Procesamiento distribuido | Media / batch | Alta | Depende de recursos y duración del cluster/job | ETL pesado, joins, agregaciones y procesamiento de grandes volúmenes |
| Firestore | Serving operacional NoSQL | Baja | Alta | Basado en operaciones, almacenamiento y tráfico | Aplicaciones web/móviles y acceso documental de baja latencia |
| Cloud Composer | Orquestación | No es el motor de procesamiento | Gestionada | Coste asociado al entorno gestionado | DAGs, dependencias, scheduling, retries y workflows empresariales |
Antipatrón crítico: convertir Cloud Functions en un motor ETL
Un pipeline puede comenzar con una Cloud Function que recibe datos por HTTP y los transforma antes de almacenarlos. El problema aparece cuando esa función empieza a asumir tareas que pertenecen a un motor distribuido: grandes joins, procesamiento de ficheros completos, agregaciones pesadas o transformaciones que requieren demasiado tiempo y memoria.
El resultado suele ser una arquitectura difícil de escalar y de observar. Además, mezclar ingesta, procesamiento y persistencia en una única función incrementa el acoplamiento y dificulta los reintentos selectivos.
El patrón recomendado es separar responsabilidades: Cloud Functions para ingesta/eventos → Pub/Sub para desacoplar → Cloud Storage para staging → Dataproc Spark para procesamiento → Firestore para serving operacional.
Regla práctica: si el procesamiento requiere paralelización real, joins complejos o trabaja con datasets grandes, no intentes convertir una función serverless en un cluster Spark improvisado.
Implementación práctica: Google Cloud Functions Python ETL
La primera capa puede recibir un payload HTTP, validarlo y publicarlo en Pub/Sub. La función no necesita conocer qué consumidor procesará posteriormente el evento. Esta separación permite añadir nuevos consumidores sin modificar el productor.
import json
import os
from google.cloud import pubsub_v1
publisher = pubsub_v1.PublisherClient()
PROJECT_ID = os.environ["GOOGLE_CLOUD_PROJECT"]
TOPIC_ID = os.environ["PUBSUB_TOPIC"]
topic_path = publisher.topic_path(PROJECT_ID, TOPIC_ID)
def ingest(request):
"""Recibe datos HTTP y los publica como evento en Pub/Sub."""
payload = request.get_json(silent=True)
if not payload:
return {"error": "Invalid JSON payload"}, 400
required_fields = {"event_id", "timestamp", "data"}
if not required_fields.issubset(payload):
return {"error": "Missing required fields"}, 400
message = json.dumps(payload).encode("utf-8")
future = publisher.publish(
topic_path,
message,
event_id=str(payload["event_id"])
)
message_id = future.result()
return {
"status": "published",
"message_id": message_id
}, 202
En producción, esta capa debería incorporar además autenticación/autorización, validación de esquema, observabilidad, correlación de eventos, políticas de reintento y una estrategia explícita de idempotencia.
Pub/Sub to Cloud Storage: separar eventos de almacenamiento
Pub/Sub funciona como una capa de transporte y desacoplamiento. Cloud Storage, en cambio, puede actuar como zona de aterrizaje persistente para los datos que posteriormente serán procesados por Spark.
Para datasets destinados a procesamiento analítico, formatos columnares como Parquet suelen ser preferibles frente a formatos orientados a intercambio como JSON, especialmente cuando el siguiente paso implica leer subconjuntos de columnas y ejecutar transformaciones distribuidas.
Un patrón habitual consiste en conservar una zona raw inmutable y generar posteriormente datasets procesados. Esto permite reproducir transformaciones, investigar errores y reconstruir resultados sin depender de la disponibilidad del productor original.
Apache Spark en Google Cloud Dataproc
Cuando el volumen o la complejidad supera lo razonable para una función serverless, Dataproc permite ejecutar jobs de Apache Spark utilizando recursos distribuidos.
La principal ventaja arquitectónica no es simplemente disponer de Spark gestionado. Es poder desacoplar el almacenamiento del cómputo: los datos pueden permanecer en Cloud Storage mientras los recursos de procesamiento se utilizan cuando son necesarios.
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, count, sum as spark_sum
def main():
spark = (
SparkSession.builder
.appName("gcp-data-pipeline")
.getOrCreate()
)
source = "gs://data-lake/raw/events/*.parquet"
destination = "gs://data-lake/processed/events"
df = spark.read.parquet(source)
result = (
df.filter(col("status") == "completed")
.groupBy("customer_id")
.agg(
count("*").alias("event_count"),
spark_sum("amount").alias("total_amount")
)
)
(
result.write
.mode("overwrite")
.parquet(destination)
)
spark.stop()
if __name__ == "__main__":
main()
En entornos empresariales conviene evitar que el job dependa de rutas o parámetros codificados. Las rutas, fechas de procesamiento, nombres de datasets y configuraciones deben parametrizarse y quedar bajo control del orquestador.
Google Cloud Composer Airflow tutorial: orquestar el pipeline
Un pipeline empresarial necesita algo más que ejecutar jobs. Necesita controlar dependencias, reintentos, ventanas temporales, sensores, backfills y estados de ejecución. Ahí entra Cloud Composer, el servicio gestionado de Google Cloud basado en Apache Airflow.
from datetime import datetime
from airflow import DAG
from airflow.providers.google.cloud.operators.dataproc import (
DataprocSubmitJobOperator
)
PROJECT_ID = "my-project"
REGION = "europe-west1"
CLUSTER_NAME = "data-processing"
SPARK_JOB = {
"reference": {
"project_id": PROJECT_ID
},
"placement": {
"cluster_name": CLUSTER_NAME
},
"pyspark_job": {
"main_python_file_uri":
"gs://data-pipeline-code/jobs/process.py"
}
}
with DAG(
dag_id="gcp_data_pipeline",
start_date=datetime(2026, 1, 1),
schedule="@daily",
catchup=False,
tags=["gcp", "data-engineering", "spark"]
) as dag:
process_data = DataprocSubmitJobOperator(
task_id="process_data",
project_id=PROJECT_ID,
region=REGION,
job=SPARK_JOB
)
El DAG no debería contener toda la lógica de negocio. Airflow debe actuar como orquestador y delegar el procesamiento a los servicios especializados. Mantener esta frontera clara mejora mantenibilidad y testabilidad.
Patrones de diseño y mejores prácticas
Desacoplamiento
Usa Pub/Sub para separar productores de consumidores. El productor no debería conocer qué procesamiento ocurrirá después.
Idempotencia
Diseña consumidores capaces de procesar un evento más de una vez sin generar efectos duplicados. Los reintentos son parte normal de un sistema distribuido.
Separación de capas
Mantén independientes ingesta, raw data, transformación y serving. Esto facilita reemplazar una tecnología sin reconstruir todo el pipeline.
Observabilidad
Cada etapa debe producir métricas y logs suficientes para responder qué ocurrió, cuándo ocurrió y dónde falló una ejecución.
FinOps
Controla duración de jobs, recursos de Spark, almacenamiento, retención de mensajes y volumen transferido. El rendimiento sin control de costes no es una arquitectura optimizada.
Gobernanza
Centraliza IAM, separación de proyectos/entornos, convenciones de nombres, clasificación de datos, secretos y políticas de acceso.
Dataproc Spark Firestore: cuándo utilizar Firestore como capa de serving
Firestore no debería confundirse con un Data Warehouse. Su valor aparece cuando el dato procesado debe ser consumido con baja latencia por una aplicación web o móvil y el modelo de acceso encaja con una base de datos documental.
Por ejemplo, Spark puede calcular diariamente métricas agregadas por cliente y publicar el resultado final en Firestore. La aplicación puede entonces consultar esos documentos directamente sin ejecutar de nuevo un procesamiento distribuido.
Este patrón crea una frontera clara entre Data Engineering y serving operacional: Spark calcula; Firestore sirve.
FinOps: el error de escalar cómputo antes de optimizar el pipeline
Un pipeline puede consumir recursos innecesariamente por leer repetidamente datasets completos, utilizar formatos ineficientes, mantener clusters activos fuera de la ventana de procesamiento o ejecutar trabajos que no necesitan realmente Spark.
Antes de incrementar recursos, mide volumen de entrada, duración, bytes procesados, utilización de CPU/memoria, frecuencia de ejecución y coste por unidad de dato procesado.
Una optimización arquitectónica puede tener más impacto que aumentar capacidad: convertir JSON a Parquet, particionar datasets, filtrar antes de joins, reducir shuffles, eliminar ejecuciones redundantes o utilizar procesamiento serverless cuando la carga sea pequeña.
Framework de implementación paso a paso
La arquitectura correcta depende del patrón de carga
Un Proyecto Data Engineering GCP paso a paso debería enseñar algo más que cómo desplegar recursos. El verdadero conocimiento está en comprender por qué cada servicio ocupa una posición concreta dentro de la arquitectura.
Cloud Functions resuelve bien la ejecución orientada a eventos. Pub/Sub proporciona desacoplamiento. Cloud Storage aporta una capa persistente y económica para datos. Dataproc Spark permite procesamiento distribuido. Cloud Composer coordina workflows complejos y Firestore proporciona una capa de serving operacional de baja latencia.
La arquitectura se vuelve robusta cuando estas piezas tienen responsabilidades claras y están conectadas mediante contratos, observabilidad, idempotencia y políticas de seguridad. Ese es el salto entre un tutorial que simplemente funciona y una plataforma de datos diseñada para producción.
Preguntas frecuentes sobre GCP Data Pipelines
¿Qué servicios de GCP necesito para construir un Data Pipeline end-to-end?
Una arquitectura habitual puede combinar Cloud Functions para ingesta ligera, Pub/Sub para desacoplar productores y consumidores, Cloud Storage como capa de aterrizaje, Dataproc con Spark para procesamiento distribuido, Cloud Composer con Airflow para orquestación y Firestore para servir datos operacionales a aplicaciones.
¿Cuándo debería utilizar Dataproc Spark en lugar de Cloud Functions para transformar datos?
Cloud Functions encaja mejor en transformaciones pequeñas, orientadas a eventos y con requisitos de ejecución acotados. Dataproc Spark resulta más apropiado cuando existe procesamiento distribuido, grandes volúmenes, joins complejos, transformaciones intensivas o necesidad de aprovechar el ecosistema Spark.
¿Qué papel desempeña Cloud Composer en una arquitectura de datos en GCP?
Cloud Composer proporciona un entorno gestionado basado en Apache Airflow para orquestar dependencias, reintentos, ventanas de ejecución, sensores, tareas y workflows que integran servicios de Google Cloud y sistemas externos.
