Streaming en Tiempo Real con Pub/Sub, Dataflow y BigQuery usando Python
Diseñar un pipeline de streaming no consiste simplemente en sustituir un batch de cinco minutos por un proceso continuo. En producción hay que diseñar para backpressure, duplicados, eventos fuera de orden, escalabilidad, observabilidad, tolerancia a fallos y coste. Esta guía muestra cómo evolucionar desde un modelo batch hacia una arquitectura de procesamiento continuo en Google Cloud utilizando Pub/Sub → Dataflow → BigQuery.
Del batch periódico al procesamiento continuo
El cambio arquitectónico fundamental es pasar de preguntar periódicamente «¿qué datos nuevos existen?» a mantener un flujo de eventos que se procesa a medida que llega.
El stack de Google Cloud
Cada componente tiene una responsabilidad clara. La clave arquitectónica consiste en no convertir BigQuery en un sistema de transporte de eventos ni Dataflow en un simple script de transformación.
Flujo de datos
Los dispositivos o aplicaciones producen eventos. Pub/Sub desacopla productores y consumidores. Dataflow procesa el stream y BigQuery proporciona la capa analítica.
Responsabilidad de cada capa
Pub/Sub absorbe variaciones de producción; Dataflow ejecuta el procesamiento distribuido; BigQuery almacena y consulta los datos analíticos.
Batch frente a Streaming
Streaming no es automáticamente mejor. El patrón correcto depende de la latencia de negocio, volumen, coste operativo y complejidad requerida.
| Criterio | Batch tradicional | Streaming con Pub/Sub + Dataflow | Implicación arquitectónica |
|---|---|---|---|
| Latencia | Minutos u horas. | Segundos o subminutos, dependiendo del pipeline. | Streaming es apropiado cuando la decisión pierde valor esperando al siguiente batch. |
| Ingesta | Archivos, dumps o cargas periódicas. | Eventos independientes publicados continuamente. | Pub/Sub desacopla productores y consumidores. |
| Procesamiento | Job que comienza y termina. | Pipeline persistente. | Hay que diseñar autoscaling, backpressure y recuperación. |
| Orden de eventos | Normalmente controlado por el lote. | Puede existir retraso o desorden. | La lógica debe considerar event time, no solamente processing time. |
| Coste | Puede ser más sencillo de controlar. | Existe consumo continuo de infraestructura. | Autoscaling, throughput y diseño de transformaciones afectan directamente al FinOps. |
| Casos de uso | ETL nocturno, reporting periódico, reconciliaciones. | Telemetría, fraude, IoT, alertas, monitorización operacional. | Elegir streaming solamente cuando la latencia tenga valor económico o operacional. |
El error más peligroso: confundir streaming con “insertar todo inmediatamente”
Un diseño ingenuo puede intentar publicar cada evento y realizar una escritura individual inmediatamente sobre el destino analítico. Esto aumenta el overhead de red, operaciones y serialización y puede provocar que el coste y la latencia crezcan cuando aumenta el throughput.
La solución es desacoplar la llegada de eventos de la persistencia. Pub/Sub funciona como buffer durable y Dataflow controla el procesamiento distribuido. La escritura hacia BigQuery debe utilizar un mecanismo de streaming apropiado y una estrategia coherente con el volumen y los requisitos de latencia.
En arquitecturas modernas conviene evaluar la
BigQuery Storage Write API frente al mecanismo clásico
de Streaming Inserts (insertAll). La decisión debe basarse
en throughput, semántica de entrega, control de streams y requisitos
operacionales, no simplemente en cuál API resulta más sencilla de llamar.
Principio Staff/Architect: el pipeline debe absorber picos de entrada sin obligar al consumidor analítico a procesar cada evento como una transacción OLTP independiente.
Pipeline Dataflow Streaming con Python
El siguiente patrón utiliza Apache Beam para leer continuamente de Pub/Sub, interpretar eventos JSON, normalizarlos y escribirlos en BigQuery. La lógica de negocio queda separada de la infraestructura del pipeline.
import json
import apache_beam as beam
from apache_beam.options.pipeline_options import (
PipelineOptions,
StandardOptions,
)
class ParseTelemetry(beam.DoFn):
"""Convierte el mensaje Pub/Sub en un registro analítico."""
def process(self, message):
try:
payload = json.loads(message.decode("utf-8"))
yield {
"device_id": str(payload["device_id"]),
"event_ts": str(payload["event_ts"]),
"temperature": float(payload["temperature"]),
"battery": float(payload.get("battery", 0)),
"region": str(payload.get("region", "unknown")),
}
except (ValueError, KeyError, TypeError) as exc:
# En producción, enviaríamos el evento a una DLQ
# con suficiente contexto para poder reprocesarlo.
yield beam.pvalue.TaggedOutput(
"dead_letter",
{
"raw_message": message.decode("utf-8", errors="replace"),
"error": str(exc),
},
)
class AddDerivedMetrics(beam.DoFn):
"""Añade métricas derivadas sin acoplarlas al origen."""
def process(self, row):
temperature = row["temperature"]
row["temperature_alert"] = (
temperature >= 80.0 or temperature <= -10.0
)
yield row
def run():
options = PipelineOptions(
streaming=True,
save_main_session=True,
)
options.view_as(StandardOptions).streaming = True
input_subscription = (
"projects/PROJECT_ID/subscriptions/telemetry-sub"
)
output_table = (
"PROJECT_ID:analytics.telemetry"
)
dead_letter_table = (
"PROJECT_ID:analytics.telemetry_dead_letter"
)
with beam.Pipeline(options=options) as pipeline:
parsed = (
pipeline
| "ReadFromPubSub"
>> beam.io.ReadFromPubSub(
subscription=input_subscription
)
| "ParseTelemetry"
>> beam.ParDo(ParseTelemetry()).with_outputs(
"dead_letter",
main="valid",
)
)
valid_rows = (
parsed.valid
| "AddDerivedMetrics"
>> beam.ParDo(AddDerivedMetrics())
)
(
valid_rows
| "WriteToBigQuery"
>> beam.io.WriteToBigQuery(
output_table,
schema={
"fields": [
{"name": "device_id", "type": "STRING"},
{"name": "event_ts", "type": "TIMESTAMP"},
{"name": "temperature", "type": "FLOAT64"},
{"name": "battery", "type": "FLOAT64"},
{"name": "region", "type": "STRING"},
{"name": "temperature_alert", "type": "BOOL"},
]
},
write_disposition=(
beam.io.BigQueryDisposition.WRITE_APPEND
),
create_disposition=(
beam.io.BigQueryDisposition.CREATE_IF_NEEDED
),
)
)
(
parsed.dead_letter
| "WriteDeadLetter"
>> beam.io.WriteToBigQuery(
dead_letter_table,
write_disposition=(
beam.io.BigQueryDisposition.WRITE_APPEND
),
create_disposition=(
beam.io.BigQueryDisposition.CREATE_IF_NEEDED
),
)
)
if __name__ == "__main__":
run()
Nota de arquitectura:
el ejemplo utiliza WriteToBigQuery para mantener el código
didáctico. En un entorno de alto throughput conviene evaluar explícitamente
la configuración de escritura de Beam y la BigQuery Storage Write API,
además de definir una estrategia de deduplicación y una DLQ operacional.
El problema real no es recibir eventos: es saber cuándo ocurrieron
Un evento puede generarse a las 10:00:00, publicarse a las 10:00:02, llegar al worker a las 10:00:04 y persistirse a las 10:00:05. Si agregamos utilizando solamente el instante de procesamiento, podemos atribuir el evento a una ventana incorrecta.
Processing Time
Es el instante en el que el sistema procesa el elemento. Es sencillo de implementar, pero no representa necesariamente el momento real en el que sucedió el evento.
Event Time
Utiliza el timestamp generado por el productor. Es el patrón adecuado para métricas temporales de negocio y requiere considerar timestamps tardíos, watermarks y triggers.
Patrones que convierten un demo en una plataforma de producción
Desacoplamiento mediante Pub/Sub
Los productores no deben conocer la velocidad de procesamiento de Dataflow. El broker absorbe variaciones de throughput y permite escalar consumidores de forma independiente.
Idempotencia y deduplicación
Un sistema distribuido debe asumir que un evento puede volver a procesarse. Diseña una clave de evento estable y determina dónde se aplicará la deduplicación.
Dead-Letter Queue
Un evento corrupto no debería bloquear el pipeline completo. Los mensajes inválidos deben aislarse, conservar contexto y quedar disponibles para reprocesamiento.
Event Time + Windows
Para agregaciones temporales utiliza event time y define explícitamente ventanas, watermarks y comportamiento ante eventos tardíos.
Schema Governance
El contrato de eventos debe evolucionar de forma compatible. Evita que productores independientes cambien tipos o semántica sin coordinación.
Observabilidad end-to-end
No basta con saber que Dataflow está “Running”. Hay que correlacionar backlog de Pub/Sub, throughput, latencia, errores y comportamiento del destino.
Cómo leer el grafo de Dataflow en producción
El grafo visual de Dataflow es una herramienta de diagnóstico. Cada etapa representa una parte del flujo y sus métricas permiten identificar dónde se está acumulando trabajo.
Señales que debes observar
Backlog: si crece continuamente, la capacidad de procesamiento no está siguiendo la velocidad de entrada.
Throughput: compara elementos recibidos y procesados para detectar un cuello de botella.
Latencia: mide cuánto tarda un evento desde su llegada hasta estar disponible para análisis.
Interpretación del grafo
Una etapa con mayor tiempo de procesamiento, menor throughput o crecimiento de backlog puede convertirse en el cuello de botella.
El autoscaling puede aumentar workers, pero no soluciona una transformación mal diseñada, una dependencia externa lenta o una escritura al destino que limite el pipeline.
Regla operacional: monitoriza el pipeline desde el productor hasta BigQuery. Un Dataflow saludable no implica necesariamente un sistema saludable si Pub/Sub acumula backlog o la disponibilidad de los datos analíticos se degrada.
Checklist para migrar de Batch a Streaming
La migración debe comenzar por el SLA de negocio y no por la tecnología. Si nadie necesita los datos en tiempo real, mantener un batch optimizado puede ser una decisión arquitectónicamente superior.
Define el SLA de latencia
Determina si el negocio necesita segundos, minutos u horas. Evita pagar por streaming cuando el SLA real sea compatible con batch.
Define el contrato de eventos
Establece identificador único, timestamp de evento, versión de esquema, productor, región y campos obligatorios.
Introduce Pub/Sub como frontera de desacoplamiento
Configura topic, suscripción, retención, dead-letter topic y permisos siguiendo el principio de mínimo privilegio.
Implementa Dataflow con Apache Beam
Separa lectura, validación, enriquecimiento, transformación y escritura. Esto facilita testing, observabilidad y evolución del pipeline.
Diseña la semántica de entrega
Decide cómo tratarás duplicados, reintentos, eventos tardíos y errores. No confundas “exactly-once” a nivel de una etapa con una garantía end-to-end de negocio.
Optimiza BigQuery
Diseña tablas, particiones y clustering pensando en los patrones de consulta. El streaming no sustituye al diseño físico del modelo analítico.
Define observabilidad y alertas
Establece umbrales para backlog, latencia, errores, throughput, disponibilidad y crecimiento de costes.
Ejecuta pruebas de carga
Simula picos superiores al tráfico esperado y valida que el sistema pueda recuperarse después de una acumulación de backlog.
Streaming no significa eliminar Batch
Una arquitectura madura combina ambos paradigmas. Los eventos pueden procesarse en streaming para obtener una visión operacional inmediata, mientras que procesos batch posteriores realizan reconciliaciones, compactaciones, backfills o controles de calidad.
Después, un proceso batch puede verificar integridad, detectar anomalías, recalcular periodos afectados o reconstruir particiones históricas. Esta combinación reduce el riesgo de exigir al camino crítico del streaming responsabilidades que no necesitan latencia baja.
El objetivo arquitectónico no es conseguir el pipeline más sofisticado, sino la menor complejidad que satisfaga el SLA de negocio.
Preguntas frecuentes sobre Pub/Sub, Dataflow y BigQuery
¿Por qué utilizar Dataflow entre Pub/Sub y BigQuery?
Porque Dataflow permite ejecutar procesamiento distribuido sobre el stream: validación, enriquecimiento, transformación, ventanas, manejo de eventos tardíos y rutas de errores. Pub/Sub se ocupa principalmente de desacoplar la ingesta y BigQuery de la analítica.
¿Es lo mismo Streaming Inserts que BigQuery Storage Write API?
No. Streaming Inserts suele hacer referencia al mecanismo tradicional de inserción de filas mediante la API de BigQuery, mientras que la BigQuery Storage Write API ofrece una interfaz moderna para escrituras de datos, con diferentes capacidades de streams y semánticas. Para un pipeline nuevo conviene evaluar la opción adecuada según volumen, latencia, coste y garantías requeridas.
¿Qué métricas debo vigilar en un pipeline Dataflow Streaming?
Como mínimo: backlog de Pub/Sub, throughput de entrada y salida, latencia, errores, elementos enviados a DLQ, número de workers, autoscaling, uso de CPU, tiempos por etapa y comportamiento de las escrituras en BigQuery. En producción también es recomendable disponer de métricas de negocio, no solamente métricas técnicas.
