Logo GCP con Eduardo

GCP con Eduardo

Descarga el código de la lección

Streaming en Tiempo Real con Pub/Sub, Dataflow y BigQuery (Python)
✦ Guía Técnica & Arquitectura

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.

01 · Lo que aprenderás

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.

Diseñar una arquitectura Pub/Sub → Dataflow → BigQuery para telemetría en tiempo real.
Implementar transformaciones streaming con Apache Beam y Python.
Gestionar ventanas, datos tardíos, duplicados, errores y backpressure.
Monitorizar el grafo, backlog, latencia, throughput y autoscaling en producción.
02 · Stack

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.

Google Cloud Pub/Sub Apache Beam Google Cloud Dataflow BigQuery Python Cloud Monitoring IAM Dead-Letter Queue

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.

IoT / Apps Pub/Sub Dataflow BigQuery

Responsabilidad de cada capa

Pub/Sub absorbe variaciones de producción; Dataflow ejecuta el procesamiento distribuido; BigQuery almacena y consulta los datos analíticos.

Buffer + Compute + Analytics
03 · Matriz de decisión

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.
Pub/Sub
Ingesta y desacoplamiento de productores.
Dataflow
Procesamiento distribuido y escalable.
BigQuery
Persistencia analítica y consultas SQL.
04 · Antipatrones & FinOps

El error más peligroso: confundir streaming con “insertar todo inmediatamente”

⚠ Error de producción

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.

05 · Implementación

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.

streaming_pipeline.py Python · Apache Beam
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.

06 · Event Time

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.

Receive Process Aggregate

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.

event_ts Window Watermark
07 · Patrones de diseño

Patrones que convierten un demo en una plataforma de producción

01

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.

02

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.

03

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.

04

Event Time + Windows

Para agregaciones temporales utiliza event time y define explícitamente ventanas, watermarks y comportamiento ante eventos tardíos.

05

Schema Governance

El contrato de eventos debe evolucionar de forma compatible. Evita que productores independientes cambien tipos o semántica sin coordinación.

06

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.

08 · Monitorización

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.

09 · Framework empresarial

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.

1

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.

2

Define el contrato de eventos

Establece identificador único, timestamp de evento, versión de esquema, productor, región y campos obligatorios.

3

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.

4

Implementa Dataflow con Apache Beam

Separa lectura, validación, enriquecimiento, transformación y escritura. Esto facilita testing, observabilidad y evolución del pipeline.

5

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.

6

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.

7

Define observabilidad y alertas

Establece umbrales para backlog, latencia, errores, throughput, disponibilidad y crecimiento de costes.

8

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.

10 · Decisiones de arquitectura

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.

Eventos Pub/Sub Dataflow Streaming BigQuery

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.

11 · FAQ

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.

Sobre el Autor: Eduardo Martínez Agrelo

AI & Data Architect

Eduardo Martínez Agrelo es AI & Data Architect y desarrolla contenidos técnicos centrados en arquitectura de datos, ingeniería de datos, inteligencia artificial y sistemas distribuidos. Su enfoque combina fundamentos técnicos con decisiones de arquitectura orientadas a rendimiento, escalabilidad, observabilidad, gobierno y eficiencia económica en entornos empresariales.

© 2026 Eduardo Martínez Agrelo · AI & Data Architect

Streaming · Pub/Sub · Dataflow · BigQuery · Python · Google Cloud