Streaming en Tiempo Real con Pub/Sub, Dataflow y BigQuery (Python)
Cómo diseñar una arquitectura de procesamiento de eventos en tiempo real en Google Cloud, evolucionando desde cargas batch hacia un pipeline continuo de telemetría con Pub/Sub, Dataflow, Apache Beam y BigQuery, incluyendo decisiones de arquitectura, FinOps, tolerancia a fallos y monitorización del grafo de ejecución.
Lo que aprenderás
El objetivo no es simplemente conectar Pub/Sub con BigQuery, sino entender cómo construir un sistema de streaming operable cuando el volumen, la latencia y el coste dejan de ser problemas teóricos.
Stack tecnológico
Una arquitectura cloud-native orientada a eventos, diseñada para separar la ingesta, el procesamiento distribuido y la capa analítica.
De Batch a Streaming: el cambio de paradigma
En un sistema batch, la pregunta habitual es «¿cuándo ejecuto el siguiente job?». En streaming, la unidad de diseño pasa a ser el evento: cada mensaje debe poder atravesar el sistema de forma continua, observable y tolerante a fallos.
La responsabilidad de Pub/Sub es proporcionar un backbone de mensajería desacoplado. Dataflow ejecuta el procesamiento distribuido mediante Apache Beam. BigQuery actúa como capa analítica para consultar la telemetría prácticamente desde el momento en que se ingesta.
Esta separación permite escalar productores y consumidores de forma independiente. Además, evita introducir lógica de transformación compleja directamente en los productores de eventos.
Matriz de decisión arquitectónica
No todos los mecanismos de ingesta resuelven el mismo problema. La decisión debe basarse en latencia, volumen, transformación requerida, consistencia, coste operativo y modelo de consumo.
| Patrón | Latencia | Transformación | Escalabilidad | Coste / Operación | Uso recomendado |
|---|---|---|---|---|---|
| Batch → BigQuery | Minutos / horas | Alta, pero offline | Alta | Simple de operar | ETL periódico, reporting, cargas históricas. |
| Pub/Sub → BigQuery | Baja | Limitada antes de persistir | Alta | Baja complejidad | Ingesta directa cuando no necesitas un pipeline de transformación complejo. |
| Pub/Sub → Dataflow → BigQuery | Baja | Muy alta | Muy alta | Mayor complejidad operativa | Telemetría, enriquecimiento, validación, ventanas y procesamiento empresarial. |
| Pub/Sub → Dataflow → Storage Write API | Muy baja | Muy alta | Muy alta | Orientado a grandes volúmenes | Ingesta moderna de alto throughput y arquitecturas de streaming exigentes. |
| OLTP → CDC → Data Warehouse | Segundos / minutos | Alta | Alta | Media / alta | Replicación de cambios operacionales y analítica casi en tiempo real. |
El antipatrón más caro: confundir throughput con latencia
En streaming, aumentar la frecuencia de escritura no significa necesariamente conseguir una arquitectura mejor. El diseño debe equilibrar latencia, tamaño de los lotes internos, throughput y coste.
Un error habitual consiste en intentar que cada evento individual termine en BigQuery como una operación aislada. Con millones de eventos por minuto, este enfoque puede incrementar la presión sobre el sistema de ingestión y producir una arquitectura difícil de gobernar económicamente.
La solución es diseñar el pipeline alrededor del throughput real. Dataflow debe procesar los eventos de forma distribuida y utilizar el mecanismo de escritura de BigQuery apropiado para el caso de uso.
Importante: Streaming Inserts y BigQuery Storage Write API no son sinónimos. Streaming Inserts es el mecanismo clásico de streaming de filas, mientras que Storage Write API es la interfaz moderna diseñada para ingestión de alto rendimiento. La elección debe evaluarse contra volumen, semántica de entrega, requisitos de latencia y capacidades actuales de Apache Beam.
También es fundamental recordar que la entrega de mensajes y la escritura analítica deben diseñarse pensando en reintentos. En sistemas distribuidos, «procesado una vez» y «efecto final exactamente una vez» no son conceptos intercambiables.
Implementación práctica con Python y Apache Beam
El siguiente ejemplo representa un pipeline de producción simplificado: consume mensajes JSON desde Pub/Sub, valida la estructura de telemetría, normaliza los tipos y escribe los registros en BigQuery mediante STREAMING_INSERTS.
import json
import logging
from datetime import datetime, timezone
import apache_beam as beam
from apache_beam.io.gcp.bigquery import WriteToBigQuery
from apache_beam.options.pipeline_options import (
PipelineOptions,
StandardOptions,
)
PROJECT_ID = "my-gcp-project"
REGION = "europe-west1"
PUBSUB_SUBSCRIPTION = (
f"projects/{PROJECT_ID}/subscriptions/telemetry-sub"
)
BQ_TABLE = f"{PROJECT_ID}:telemetry.events"
class ParseTelemetry(beam.DoFn):
"""Parsea y normaliza eventos JSON procedentes de Pub/Sub."""
def process(self, message):
try:
payload = json.loads(message.decode("utf-8"))
device_id = str(payload["device_id"])
event_timestamp = payload["event_timestamp"]
temperature = float(payload["temperature"])
humidity = float(payload["humidity"])
yield {
"device_id": device_id,
"event_timestamp": event_timestamp,
"temperature": temperature,
"humidity": humidity,
"ingested_at": datetime.now(
timezone.utc
).isoformat(),
}
except (ValueError, TypeError, KeyError, json.JSONDecodeError):
logging.exception("Invalid telemetry event")
def run():
options = PipelineOptions(
project=PROJECT_ID,
region=REGION,
streaming=True,
save_main_session=True,
)
options.view_as(StandardOptions).streaming = True
table_schema = {
"fields": [
{
"name": "device_id",
"type": "STRING",
"mode": "REQUIRED",
},
{
"name": "event_timestamp",
"type": "TIMESTAMP",
"mode": "REQUIRED",
},
{
"name": "temperature",
"type": "FLOAT64",
"mode": "NULLABLE",
},
{
"name": "humidity",
"type": "FLOAT64",
"mode": "NULLABLE",
},
{
"name": "ingested_at",
"type": "TIMESTAMP",
"mode": "NULLABLE",
},
]
}
with beam.Pipeline(options=options) as pipeline:
(
pipeline
| "ReadFromPubSub"
>> beam.io.ReadFromPubSub(
subscription=PUBSUB_SUBSCRIPTION
)
| "ParseTelemetry"
>> beam.ParDo(ParseTelemetry())
| "WriteToBigQuery"
>> WriteToBigQuery(
table=BQ_TABLE,
schema=table_schema,
method=WriteToBigQuery.Method.STREAMING_INSERTS,
write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND,
create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED,
)
)
if __name__ == "__main__":
run()
Qué ocurre realmente dentro del grafo de Dataflow
Una de las capacidades más importantes para operar streaming en producción es aprender a leer el grafo de ejecución como una representación del flujo de datos y no simplemente como una lista de pasos del código.
① ReadFromPubSub
Los workers reciben eventos desde la suscripción. El comportamiento del backlog permite detectar si los productores están generando datos más rápidamente de lo que el pipeline puede procesar.
- Backlog de mensajes
- Throughput de entrada
- Errores de recepción
- Distribución de carga
② ParseTelemetry
Es la frontera entre datos externos y el modelo interno del pipeline. Aquí deben controlarse JSON inválido, campos ausentes, tipos incorrectos y reglas básicas de calidad.
- Validación de esquema
- Normalización
- Dead-letter strategy
- Observabilidad de errores
③ Transformaciones
En una arquitectura real esta capa puede incluir enriquecimiento con catálogos, deduplicación, ventanas temporales, agregaciones o clasificación de eventos.
- Windowing
- Triggers
- Stateful processing
- Side inputs
④ BigQuery Sink
La última fase persiste los eventos para consumo analítico. El diseño de la tabla es tan importante como el código: particionado, clustering, tipos, retención y estrategia de consulta afectan directamente al coste.
- WRITE_APPEND
- Particionado temporal
- Clustering
- Gobernanza y retención
Patrones de diseño para producción
La versión empresarial del pipeline requiere resolver mucho más que la conectividad entre servicios.
Backpressure y elasticidad
El sistema debe soportar picos de producción sin asumir que la tasa de eventos es constante. Pub/Sub desacopla productores y consumidores, mientras Dataflow permite ajustar la capacidad de procesamiento.
Dead-letter architecture
Un mensaje corrupto no debería bloquear indefinidamente el pipeline. Los eventos que no cumplen el contrato deben aislarse para análisis, reparación y replay controlado.
Idempotencia
Los reintentos forman parte del modelo distribuido. Cuando un evento puede procesarse más de una vez, la arquitectura debe definir cómo evitar efectos duplicados o cómo reconciliarlos posteriormente.
Event time frente a processing time
La telemetría puede llegar tarde. Para análisis temporales correctos, el timestamp del evento debe distinguirse del instante en que Dataflow lo procesa o de cuándo llega a BigQuery.
Observabilidad como feature
Un pipeline sin métricas de throughput, latencia, backlog y errores no está preparado para producción. Las métricas deben formar parte del diseño desde el primer día.
FinOps por diseño
El coste debe modelarse desde la arquitectura: volumen de mensajes, capacidad de Dataflow, frecuencia de escritura, almacenamiento, consultas posteriores y retención.
BigQuery: diseña la tabla para el patrón de acceso
Ingerir rápido no sirve de mucho si las consultas analíticas posteriores obligan a recorrer grandes cantidades de datos innecesariamente.
Particionado temporal
Para telemetría, normalmente existe una dimensión temporal natural. Aprovecharla para particionar permite reducir el volumen de datos que deben escanearse en consultas acotadas por fecha o intervalo temporal.
Clustering
Si las consultas filtran habitualmente por device_id, región, tipo de evento u otra dimensión de alta utilidad analítica, el clustering puede complementar el particionado.
Separar ingestión de consumo
No conviene asumir que la tabla de ingestión debe ser idéntica al modelo dimensional final. En arquitecturas maduras es habitual separar una capa raw o de eventos de las estructuras optimizadas para analítica.
Checklist de implementación empresarial
Utiliza este framework para pasar de un prototipo de streaming a un servicio operable en producción.
Define el contrato del evento
Establece schema, tipos, campos obligatorios, timestamp del evento, identificador único y estrategia para eventos inválidos.
Diseña Pub/Sub
Define topic, subscriptions, retención, política de reintentos, dead-letter topic y estrategia de consumidores.
Construye el pipeline Beam
Mantén separadas ingesta, parsing, validación, transformación, enriquecimiento y escritura.
Selecciona el sink de BigQuery
Evalúa Streaming Inserts frente a Storage Write API según volumen, latencia, semántica requerida y capacidades de la versión de Beam utilizada.
Optimiza el modelo analítico
Configura particionado, clustering, retención y políticas de acceso antes de que la tabla se convierta en un almacén de eventos masivo.
Instrumenta observabilidad
Monitoriza backlog, throughput, latencia, errores, utilización de workers y comportamiento del sink de BigQuery.
Prueba los fallos
Simula mensajes corruptos, picos de tráfico, caída temporal de dependencias, reintentos y recuperación del pipeline.
Establece FinOps
Define presupuestos, alertas y métricas de coste. El objetivo es conocer el coste por volumen procesado y detectar anomalías.
Monitorización del grafo en Dataflow
En la consola de Dataflow, el grafo de ejecución proporciona una visión operacional del pipeline. Para diagnosticar un incidente, no basta con comprobar que el job está en estado Running.
Un pipeline puede estar técnicamente ejecutándose mientras acumula backlog porque su capacidad de procesamiento es inferior a la tasa de entrada. Por ello, la primera pregunta debe ser: ¿el sistema está manteniendo el ritmo de producción?
Señales de salud
- Backlog estable o decreciente.
- Throughput de entrada y salida equilibrado.
- Latencia dentro del SLO.
- Ausencia de errores sostenidos.
Señales de degradación
- Backlog creciendo continuamente.
- Workers saturados.
- Transformaciones con latencia anormal.
- Errores de escritura hacia BigQuery.
Desde una perspectiva de arquitectura, el grafo permite localizar dónde se produce la pérdida de capacidad: entrada, transformación, estado, enriquecimiento o sink. Esa información es crítica para diferenciar un problema de infraestructura de un problema de diseño.
Preguntas frecuentes sobre Streaming en GCP
¿Cuándo tiene sentido usar Dataflow entre Pub/Sub y BigQuery?
Dataflow es especialmente adecuado cuando necesitas transformar, validar, enriquecer, deduplicar, agrupar o enrutar eventos antes de almacenarlos en BigQuery. También permite implementar ventanas, triggers, estado y una estrategia operacional basada en Apache Beam. Si simplemente necesitas mover eventos sin procesamiento significativo, una integración más directa puede ser suficiente.
¿Streaming Inserts y Storage Write API son lo mismo?
No. Son mecanismos diferentes para escribir datos en BigQuery. Streaming Inserts representa el mecanismo clásico de inserción de filas mediante streaming. Storage Write API proporciona una interfaz moderna orientada a streams y a escenarios de ingestión de alto rendimiento. Para una arquitectura nueva conviene evaluar ambos mecanismos en función de volumen, latencia, semántica de entrega, compatibilidad con Beam y requisitos operativos.
¿Cómo sé si mi pipeline de Dataflow tiene suficiente capacidad?
No debes mirar únicamente el número de workers. Compara continuamente la tasa de entrada con la tasa de procesamiento, observa el backlog de Pub/Sub y mide la latencia de extremo a extremo. Si el backlog crece de forma sostenida mientras la entrada permanece por encima de la capacidad de procesamiento, necesitas revisar autoscaling, paralelismo, transformaciones costosas o el sink.
