Volvamos un momento a la aplicación de AlpinaShop. Un cliente pulsa "Confirmar pedido". La función de Flask que atiende esa petición tiene que hacer varias cosas: guardar el pedido en alpinashop-pedidos, cobrar en la pasarela de pago, avisar al almacén para que lo prepare, avisar a facturación para que emita la factura, enviar el correo de confirmación, descontar stock, y —desde el módulo 4— informar a la analítica.

Hoy, ese código es una lista de llamadas, una detrás de otra. Y esa lista tiene tres problemas que no se ven hasta el peor día del año.

El cliente espera a todos. El tiempo de respuesta de la compra es la suma de los siete pasos. Si el servicio de correo tarda cuatro segundos, el cliente ve una rueda girando cuatro segundos de más por algo que no le importa: él ya ha comprado.

Un fallo lo tira todo. Si el sistema del almacén está reiniciándose, la llamada falla. ¿Y ahora qué? ¿Se deshace el pedido, que ya está cobrado? ¿Se continúa y se pierde el aviso? ¿Se reintenta y se corre el riesgo de duplicar la factura, que sí se envió? No hay ninguna respuesta buena.

Cada novedad toca el código de la compra. Añadir la analítica significa modificar, probar y desplegar la función más crítica de la tienda. Y luego llegará el programa de fidelización, y el antifraude, y el aviso al proveedor. La función de compra crece sin parar y cada cambio arriesga la caja registradora.

Cloud Pub/Sub resuelve los tres de una vez con una idea muy simple: en vez de llamar a siete sistemas, la tienda publica un hecho —"se ha confirmado el pedido PED-2026-0042"— y sigue con lo suyo. Quien esté interesado en ese hecho, se suscribe. La tienda no sabe ni le importa cuántos son.

En esta lección crearás el topic pedidos-nuevos y sus suscripciones, publicarás desde Flask, consumirás de las cuatro formas posibles, entenderás las garantías reales del servicio —que no son las que la gente supone— y aprenderás por qué la idempotencia no es un adorno de arquitectura sino un requisito ineludible.

Contenido

  1. El problema del acoplamiento, antes y después
  2. El modelo publicación-suscripción
  3. Las garantías reales de Pub/Sub
  4. Crear el topic pedidos-nuevos y sus suscripciones
  5. Publicar desde la aplicación Flask
  6. Consumo pull con el cliente asíncrono
  7. Consumo push a un endpoint HTTPS con OIDC
  8. Confirmaciones, plazos y reintentos
  9. Temas de mensajes fallidos (dead letter)
  10. Idempotencia: el consumidor debe tolerar duplicados
  11. Retención, seek e instantáneas
  12. Filtros de suscripción por atributo
  13. Suscripciones de BigQuery y de Cloud Storage
  14. Ordenación por clave y sus implicaciones
  15. Pub/Sub Lite y otras alternativas
  16. Monitorización y coste
  17. Las notificaciones del bucket alpinashop-catalogo

  1. El problema del acoplamiento, antes y después

flowchart TD
    subgraph Antes["ANTES: llamadas sincronas encadenadas"]
        W1["Flask: confirmar_pedido()"]
        A1["Almacen"]
        F1["Facturacion"]
        M1["Correo"]
        AN1["Analitica"]
        W1 -->|"espera 200 ms"| A1
        W1 -->|"espera 350 ms"| F1
        W1 -->|"espera 4 s"| M1
        W1 -->|"espera 800 ms"| AN1
    end
flowchart TD
    subgraph Despues["DESPUES: un hecho publicado, N interesados"]
        W2["Flask: confirmar_pedido()"]
        T["Topic pedidos-nuevos"]
        S1["sub-almacen"]
        S2["sub-facturacion"]
        S3["sub-analitica"]
        S4["sub-email"]
        A2["Servicio de almacen"]
        F2["Servicio de facturacion"]
        AN2["Dataflow -> BigQuery"]
        M2["Servicio de correo"]

        W2 -->|"publica: 15 ms"| T
        T --> S1 --> A2
        T --> S2 --> F2
        T --> S3 --> AN2
        T --> S4 --> M2
    end

Lo que cambia en concreto:

Aspecto Llamadas encadenadas Con Pub/Sub
Latencia de la compra Suma de todos los pasos (~5,4 s) Solo la publicación (~15 ms)
Un consumidor caído Rompe la compra El mensaje espera en su suscripción
Añadir un consumidor Modificar y desplegar la tienda Crear una suscripción. Cero cambios
Pico de tráfico Cada sistema debe aguantar el pico Pub/Sub absorbe; cada uno consume a su ritmo
Reprocesar un día Imposible sin scripts ad hoc seek a un instante anterior
Trazabilidad Logs repartidos Métricas por suscripción

La expresión técnica de esto es desacoplamiento: el productor no conoce a los consumidores, no sabe cuántos hay, no espera su respuesta y no falla si fallan. El precio es que el sistema pasa a ser eventualmente consistente: cuando la tienda responde "pedido confirmado", el almacén todavía no lo sabe. Lo sabrá en unos milisegundos, o en unos segundos si estaba ocupado. Ese matiz hay que aceptarlo conscientemente, porque cambia cómo se diseñan las pantallas: no puedes mostrar "preparando envío" inmediatamente después de comprar si el almacén aún no se ha enterado.

  1. El modelo publicación-suscripción

Cuatro conceptos, y conviene precisarlos porque el vocabulario se usa mal a menudo.

Topic (tema). El canal donde se publica. Representa un tipo de hecho: pedidos-nuevos, imagenes-subidas, stock-agotado. Un topic no guarda nada por sí mismo; es un punto de entrada.

Suscripción. La cola de un consumidor concreto sobre un topic. Aquí es donde viven realmente los mensajes. Cada suscripción recibe una copia independiente de cada mensaje publicado, y lleva su propia contabilidad de qué ha confirmado y qué no.

Esto último es la clave que más cuesta interiorizar:

Se publica 1 mensaje en pedidos-nuevos
   └── sub-almacen      recibe su copia  → la confirma a los 0,2 s
   └── sub-facturacion  recibe su copia  → la confirma a los 0,5 s
   └── sub-analitica    recibe su copia  → su consumidor está caído: espera 6 horas

Que el almacén confirme no afecta en absoluto a la copia de analítica. Son colas separadas alimentadas por el mismo topic.

Y el corolario práctico: si conectas dos instancias de tu servicio de almacén a la misma suscripción sub-almacen, cada mensaje irá a una de las dos. Eso es reparto de carga, y es lo que quieres para escalar. Si en cambio creas dos suscripciones distintas para el mismo servicio, cada instancia procesará todos los mensajes y el trabajo se hará dos veces. Una suscripción por consumidor lógico, tantas instancias como quieras dentro.

Mensaje. Tiene tres partes:

  • data: el cuerpo, bytes (típicamente JSON codificado en UTF-8). Máximo 10 MB.
  • attributes: pares clave-valor de texto, hasta 100. Son metadatos, y su gran virtud es que se pueden filtrar sin abrir el cuerpo.
  • Campos del sistema: messageId (único, asignado por Pub/Sub), publishTime, orderingKey opcional.

Confirmación (ack). El consumidor dice "procesado, no me lo vuelvas a mandar". Hasta que llega esa confirmación, Pub/Sub considera el mensaje pendiente y lo reenviará.

  1. Las garantías reales de Pub/Sub

Este apartado es el más importante de la lección, porque la mayoría de los errores de producción con mensajería vienen de suponer garantías que no existen.

Garantía ¿La da Pub/Sub? Consecuencia práctica
Entrega al menos una vez Sí, por defecto Tu consumidor recibirá duplicados. Hay que asumirlo
Entrega como máximo una vez No Nunca se pierde un mensaje por diseño
Exactamente una vez Sí, si se activa en la suscripción (con condiciones) Reduce, no elimina, la necesidad de idempotencia
Orden de llegada No, salvo con clave de ordenación Los mensajes pueden llegar desordenados
Durabilidad Sí, replicado en varias zonas No se pierden aunque caiga una zona
Retención 7 días por defecto, hasta 31 Se pueden reproducir
Entrega en orden entre topics distintos No, en ningún caso No asumas relación temporal entre topics

Las dos primeras filas merecen desarrollo, porque son contraintuitivas.

Por qué habrá duplicados. El escenario es sencillo: tu consumidor recibe el mensaje, lo procesa correctamente, envía la confirmación, y la confirmación se pierde por la red. Pub/Sub no la recibe, considera que el mensaje sigue pendiente y lo reenvía. Tu consumidor lo procesa por segunda vez. No hay ningún fallo en tu código y aun así ha ocurrido.

También ocurre si el consumidor tarda más que el ack deadline, si se reinicia a mitad de procesamiento, o si el autoescalado reasigna el mensaje. Es normal, es esperable y no es un error. Por eso el apartado 10 existe.

Exactamente una vez. Pub/Sub ofrece suscripciones con esta semántica (--enable-exactly-once-delivery), con dos matices importantes: solo aplica dentro de una región, y garantiza que no habrá una segunda entrega confirmada, no que tu procesamiento sea atómico. Si tu consumidor escribe en la base de datos y luego se cae antes de confirmar, el trabajo ya está hecho y el mensaje volverá. La idempotencia sigue siendo necesaria; exactamente una vez solo reduce la frecuencia con la que la necesitas.

Por qué no hay orden. Pub/Sub reparte los mensajes entre muchos servidores para escalar. Si el mensaje A se publica un milisegundo antes que el B, pueden acabar en servidores distintos y llegar en cualquier orden. Para AlpinaShop esto importa poco en pedidos-nuevos —cada pedido es independiente—, pero importaría muchísimo si publicáramos cambios de estado del mismo pedido: "confirmado", "enviado", "entregado" fuera de orden dejaría el pedido marcado como confirmado después de entregado. Ese caso se resuelve con clave de ordenación (apartado 14).

  1. Crear el topic pedidos-nuevos y sus suscripciones

gcloud config set project alpinashop-datos
gcloud services enable pubsub.googleapis.com

# El topic: un tipo de hecho de negocio
gcloud pubsub topics create pedidos-nuevos \
  --message-retention-duration=7d \
  --labels=entorno=produccion,equipo=desarrollo,centro-coste=plataforma,aplicacion=tienda

--message-retention-duration en el topic habilita la reproducción: permite hacer seek a un instante del pasado incluso para suscripciones creadas después. Sin esto, solo se puede volver atrás dentro de lo que retenga cada suscripción.

Ahora el topic de mensajes fallidos, que crearemos antes de las suscripciones porque lo van a necesitar:

gcloud pubsub topics create pedidos-nuevos-fallidos \
  --labels=entorno=produccion,equipo=desarrollo,aplicacion=tienda

# Y una suscripcion sobre el, para poder inspeccionar lo que caiga ahi
gcloud pubsub subscriptions create sub-pedidos-fallidos \
  --topic=pedidos-nuevos-fallidos \
  --message-retention-duration=31d \
  --ack-deadline=600

Las tres suscripciones de consumo:

# 1) ALMACEN: pull, procesamiento rapido, tolera duplicados por idempotencia
gcloud pubsub subscriptions create sub-almacen \
  --topic=pedidos-nuevos \
  --ack-deadline=30 \
  --message-retention-duration=7d \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5 \
  --min-retry-delay=10s \
  --max-retry-delay=600s \
  --labels=consumidor=almacen

# 2) FACTURACION: exactamente una vez, porque duplicar una factura es grave
gcloud pubsub subscriptions create sub-facturacion \
  --topic=pedidos-nuevos \
  --ack-deadline=60 \
  --message-retention-duration=7d \
  --enable-exactly-once-delivery \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5 \
  --labels=consumidor=facturacion

# 3) ANALITICA: la consumira el pipeline de Dataflow de 04-02
gcloud pubsub subscriptions create sub-analitica \
  --topic=pedidos-nuevos \
  --ack-deadline=120 \
  --message-retention-duration=7d \
  --labels=consumidor=analitica

Comentarios sobre las diferencias, que son deliberadas:

  • sub-facturacion con exactamente una vez. Emitir dos facturas del mismo pedido es un problema contable real. Aunque el consumidor debe ser idempotente igualmente, esta garantía reduce la exposición. Tiene un coste: menor rendimiento máximo por suscripción.
  • sub-analitica sin dead letter. El pipeline de Dataflow ya tiene su propio mecanismo de cuarentena (04-02): los mensajes malformados van a pedidos_streaming_errores. Duplicar el mecanismo complicaría el diagnóstico.
  • ack deadline distintos. 30 s para el almacén (operación rápida), 120 s para analítica (Dataflow procesa por lotes internos).

El dead letter necesita permisos explícitos, y esto se olvida siempre:

PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROY_NUM}@gcp-sa-pubsub.iam.gserviceaccount.com"

# Pub/Sub necesita poder PUBLICAR en el topic de fallidos...
gcloud pubsub topics add-iam-policy-binding pedidos-nuevos-fallidos \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.publisher"

# ...y CONFIRMAR en la suscripcion de origen para retirar el mensaje
gcloud pubsub subscriptions add-iam-policy-binding sub-almacen \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.subscriber"
gcloud pubsub subscriptions add-iam-policy-binding sub-facturacion \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.subscriber"

Sin estos dos permisos, la configuración de dead letter se acepta sin protestar y no funciona: los mensajes se reintentan eternamente en lugar de desviarse. Es un fallo silencioso clásico.

Verificación:

gcloud pubsub topics list --format="table(name)"
gcloud pubsub subscriptions list \
  --format="table(name, topic, ackDeadlineSeconds, deadLetterPolicy.maxDeliveryAttempts)"

  1. Publicar desde la aplicación Flask

Ahora la parte de Dani, el desarrollador backend. La publicación desde la aplicación de catálogo:

"""publicador.py -- Publicacion de eventos de pedido en Pub/Sub."""
import json
import logging
from concurrent import futures

from google.api_core import retry
from google.cloud import pubsub_v1

PROYECTO = "alpinashop-datos"
TOPIC = "pedidos-nuevos"

# 1) CONFIGURACION DE LOTES.
#    El cliente acumula mensajes y los envia juntos: menos llamadas de red,
#    menos coste. Publica cuando se cumple el PRIMERO de los tres limites.
config_lote = pubsub_v1.types.BatchSettings(
    max_messages=100,        # o 100 mensajes...
    max_bytes=1024 * 1024,   # ...o 1 MB acumulado...
    max_latency=0.05,        # ...o 50 ms de espera. Lo que ocurra antes.
)

# 2) CONFIGURACION DE REINTENTOS con retroceso exponencial.
config_reintentos = retry.Retry(
    initial=0.1,     # primer reintento a los 100 ms
    maximum=60.0,    # nunca esperar mas de 60 s entre intentos
    multiplier=2.0,  # 0,1 -> 0,2 -> 0,4 -> 0,8 ...
    deadline=600.0,  # se rinde a los 10 minutos
)

# 3) El cliente es CARO de crear: uno por proceso, reutilizado.
#    Crearlo dentro de la funcion de compra seria un error de rendimiento grave.
publicador = pubsub_v1.PublisherClient(batch_settings=config_lote)
ruta_topic = publicador.topic_path(PROYECTO, TOPIC)

logger = logging.getLogger(__name__)


def publicar_pedido(pedido: dict) -> futures.Future:
    """Publica un evento de pedido confirmado. NO bloquea.

    Devuelve un Future. La aplicacion puede seguir respondiendo al cliente
    sin esperar la confirmacion de Pub/Sub.
    """
    cuerpo = json.dumps(pedido, ensure_ascii=False).encode("utf-8")

    futuro = publicador.publish(
        ruta_topic,
        data=cuerpo,
        # ATRIBUTOS: metadatos filtrables sin abrir el cuerpo (apartado 12).
        # Todos los valores deben ser CADENAS.
        tipo_evento="pedido_confirmado",
        origen=pedido.get("canal", "web"),
        pais=pedido["envio"]["pais"],
        version_esquema="1",
        momento_evento=pedido["momento"],   # lo usara Dataflow (04-02)
        pedido_id=pedido["pedido_id"],
    )

    def _al_terminar(fut):
        try:
            id_mensaje = fut.result()
            logger.info("Publicado %s como %s", pedido["pedido_id"], id_mensaje)
        except Exception:
            # CRITICO: si la publicacion falla definitivamente, hay que
            # registrarlo para poder recuperarlo. Un log de error aqui
            # debe disparar una alerta (06-04).
            logger.exception("FALLO AL PUBLICAR el pedido %s",
                             pedido["pedido_id"])

    futuro.add_done_callback(_al_terminar)
    return futuro

Y su uso en la vista de Flask:

from flask import Flask, jsonify, request

app = Flask(__name__)


@app.post("/api/pedidos")
def confirmar_pedido():
    datos = request.get_json()

    # 1) Lo IMPRESCINDIBLE y sincrono: persistir y cobrar.
    #    Si esto falla, el cliente debe enterarse.
    pedido = guardar_en_cloud_sql(datos)
    cobrar_en_pasarela(pedido)

    # 2) Lo DERIVADO: se publica el hecho y se sigue.
    #    Almacen, facturacion, correo y analitica se enteraran solos.
    publicar_pedido({
        "pedido_id": pedido.id,
        "momento": pedido.creado_en.isoformat(),
        "cliente_id": pedido.cliente_id,
        "canal": pedido.canal,
        "total": str(pedido.total),
        "envio": {"pais": pedido.pais, "ciudad": pedido.ciudad},
        "lineas": [
            {"sku": l.sku, "cantidad": l.cantidad, "precio": str(l.precio)}
            for l in pedido.lineas
        ],
    })

    # 3) Respuesta inmediata: no esperamos a Pub/Sub ni a los consumidores.
    return jsonify({"pedido_id": pedido.id, "estado": "confirmado"}), 201

Cuatro cuestiones de diseño que merecen ser explícitas:

Qué va dentro del mensaje. Aquí publicamos un evento gordo, con las líneas incluidas. La alternativa es un evento fino con solo el pedido_id, obligando a cada consumidor a consultar la base de datos. El gordo evita esa carga de lecturas pero acopla el esquema del mensaje; el fino es más flexible pero multiplica las consultas a Cloud SQL. Para AlpinaShop, con ~1.200 pedidos al mes, el evento gordo es claramente mejor: menos llamadas a la base de datos operativa y los consumidores son autosuficientes.

version_esquema como atributo. El día que cambie la estructura del mensaje, los consumidores antiguos podrán reconocer y rechazar lo que no entienden en vez de romperse. Cuesta un atributo y ahorra un incidente.

El momento_evento. Es el atributo que el pipeline de Dataflow de 04-02 usa como timestamp_attribute. Sin él, las ventanas de tiempo del evento no funcionan.

Qué pasa si la publicación falla. Es el riesgo real de este diseño: el pedido está cobrado y el aviso no salió. Para AlpinaShop, el registro con alerta es suficiente dado el volumen. En sistemas de mayor exigencia se usa el patrón outbox: se escribe el evento en una tabla de la misma base de datos dentro de la misma transacción del pedido, y un proceso aparte lo publica. Así el evento y el pedido son atómicos.

Prueba rápida desde la línea de comandos:

gcloud pubsub topics publish pedidos-nuevos \
  --message='{"pedido_id":"PED-2026-0042","momento":"2026-03-14T10:22:31Z","cliente_id":"CLI-8821","canal":"web","total":"192.27","envio":{"pais":"ES","ciudad":"Barcelona"}}' \
  --attribute=tipo_evento=pedido_confirmado,origen=web,pais=ES,version_esquema=1,momento_evento=2026-03-14T10:22:31Z

  1. Consumo pull con el cliente asíncrono

En el modo pull, el consumidor pide mensajes. La biblioteca de Python usa streaming pull: mantiene una conexión abierta y recibe mensajes según llegan, con muy poca latencia.

"""consumidor_almacen.py -- Servicio de almacen de AlpinaShop."""
import json
import logging
import signal
import sys

from google.cloud import pubsub_v1

PROYECTO = "alpinashop-datos"
SUSCRIPCION = "sub-almacen"

logger = logging.getLogger(__name__)

# CONTROL DE FLUJO: limita cuantos mensajes se tienen a la vez sin confirmar.
# Sin esto, el cliente puede aceptar miles de mensajes, no darles salida
# a tiempo, dejar vencer el ack deadline y provocar reentregas masivas.
control_flujo = pubsub_v1.types.FlowControl(
    max_messages=50,
    max_bytes=10 * 1024 * 1024,
)


def procesar(mensaje):
    """Callback ejecutado por la biblioteca en un hilo del pool."""
    try:
        pedido = json.loads(mensaje.data.decode("utf-8"))
    except json.JSONDecodeError:
        # Mensaje corrupto: NO tiene arreglo con reintentos.
        # Se confirma para retirarlo y se registra para investigar.
        logger.error("Mensaje ilegible, descartado: %s", mensaje.message_id)
        mensaje.ack()
        return

    tipo = mensaje.attributes.get("tipo_evento")
    if tipo != "pedido_confirmado":
        # No es para nosotros; lo confirmamos sin hacer nada.
        mensaje.ack()
        return

    try:
        # La idempotencia se implementa aqui dentro (apartado 10)
        crear_orden_de_preparacion(pedido)
        mensaje.ack()
        logger.info("Preparacion creada para %s", pedido["pedido_id"])

    except ErrorTemporal as exc:
        # Fallo transitorio (BD saturada, red): NACK para reintentar ya.
        logger.warning("Fallo temporal en %s: %s", pedido["pedido_id"], exc)
        mensaje.nack()

    except Exception:
        # Fallo desconocido: NACK. Tras 5 intentos ira al dead letter.
        logger.exception("Fallo procesando %s", pedido.get("pedido_id"))
        mensaje.nack()


def main():
    suscriptor = pubsub_v1.SubscriberClient()
    ruta = suscriptor.subscription_path(PROYECTO, SUSCRIPCION)

    futuro = suscriptor.subscribe(ruta, callback=procesar,
                                  flow_control=control_flujo)
    logger.info("Escuchando en %s...", ruta)

    # Apagado limpio: al recibir SIGTERM (Kubernetes, Cloud Run),
    # se dejan de aceptar mensajes nuevos y se terminan los que hay.
    def apagar(signum, frame):
        logger.info("Senal %s recibida, cerrando...", signum)
        futuro.cancel()
        futuro.result(timeout=30)
        sys.exit(0)

    signal.signal(signal.SIGTERM, apagar)
    signal.signal(signal.SIGINT, apagar)

    with suscriptor:
        try:
            futuro.result()
        except Exception:
            logger.exception("El suscriptor ha fallado")
            futuro.cancel()


if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    main()

Tres puntos que hacen la diferencia entre un consumidor de juguete y uno de producción:

El control de flujo. Sin FlowControl, la biblioteca acepta tantos mensajes como le manden. Si tu procesamiento es lento, muchos vencerán su plazo, se reenviarán, tu consumidor los procesará otra vez, y entrarás en una espiral en la que cada vez hay más trabajo duplicado. Limitar los mensajes en vuelo es lo que evita esa espiral.

ack() frente a nack(). ack() retira el mensaje definitivamente. nack() lo devuelve para reentrega inmediata. Si no haces ninguna de las dos, el mensaje se reenvía al vencer el plazo, lo que retrasa el reintento pero funciona igual.

Distinguir errores recuperables de irrecuperables. Un JSON corrupto no mejorará por reintentarlo cinco veces: se descarta con registro. Una base de datos saturada sí: nack(). Confundirlos hace que los mensajes envenenados consuman recursos indefinidamente o que se pierdan datos recuperables.

También existe el pull síncrono, útil para procesos por lotes:

# Leer hasta 10 mensajes sin confirmarlos (para inspeccionar)
gcloud pubsub subscriptions pull sub-almacen --limit=10 --format=json

# Leerlos y confirmarlos
gcloud pubsub subscriptions pull sub-almacen --limit=10 --auto-ack

  1. Consumo push a un endpoint HTTPS con OIDC

En el modo push, Pub/Sub hace un POST HTTPS a una URL tuya. Tú no mantienes ningún proceso escuchando: el servicio te llama.

Aspecto Pull Push
Quién inicia El consumidor Pub/Sub
Infraestructura Proceso siempre en marcha Endpoint HTTPS; puede escalar a cero
Control del ritmo Total, con FlowControl Limitado (ventana deslizante automática)
Confirmación Explícita (ack()) Implícita: HTTP 2xx confirma, otro código es nack
Ideal para Alto volumen, procesamiento continuo Cloud Run, Cloud Functions, volumen moderado

Push encaja perfectamente con lo que viene en el curso: Cloud Run (07-02, donde acabará el catálogo según la decisión DA-001) y Cloud Functions (06-03) escalan a cero, así que no tiene sentido tener un proceso esperando mensajes.

La configuración segura usa autenticación OIDC: Pub/Sub adjunta un token firmado por Google que identifica una cuenta de servicio, y tu endpoint lo verifica. Sin esto, cualquiera que descubra tu URL podría inyectar mensajes falsos.

# 1) Cuenta de servicio que representara a Pub/Sub ante tu endpoint
gcloud iam service-accounts create sa-pubsub-invocador \
  --display-name="Identidad de Pub/Sub para invocar servicios push"

SA_INV="[email protected]"

# 2) Permitir que invoque el servicio de Cloud Run del almacen
gcloud run services add-iam-policy-binding svc-almacen \
  --region=europe-west1 \
  --member="serviceAccount:${SA_INV}" \
  --role="roles/run.invoker"

# 3) Permitir que Pub/Sub genere tokens en nombre de esa cuenta
PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
gcloud iam service-accounts add-iam-policy-binding "$SA_INV" \
  --member="serviceAccount:service-${PROY_NUM}@gcp-sa-pubsub.iam.gserviceaccount.com" \
  --role="roles/iam.serviceAccountTokenCreator"

# 4) Suscripcion push con OIDC
gcloud pubsub subscriptions create sub-almacen-push \
  --topic=pedidos-nuevos \
  --push-endpoint="https://svc-almacen-xxxxx.europe-west1.run.app/eventos/pedidos" \
  --push-auth-service-account="$SA_INV" \
  --push-auth-token-audience="https://svc-almacen-xxxxx.europe-west1.run.app" \
  --ack-deadline=60 \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5

El endpoint receptor:

import base64
import json

from flask import Flask, request

app = Flask(__name__)


@app.post("/eventos/pedidos")
def recibir_evento():
    """Endpoint push de Pub/Sub.

    Cloud Run ya ha validado el token OIDC antes de que llegue aqui,
    porque el servicio requiere autenticacion y la cuenta invocadora
    tiene run.invoker. Si el servicio fuera publico, habria que
    verificar la cabecera Authorization manualmente.
    """
    sobre = request.get_json(silent=True)
    if not sobre or "message" not in sobre:
        # 400: no se reintentara. Es un error de formato, no transitorio.
        return "peticion mal formada", 400

    mensaje = sobre["message"]

    # El cuerpo viene en base64 dentro del sobre
    cuerpo = base64.b64decode(mensaje.get("data", "")).decode("utf-8")
    atributos = mensaje.get("attributes", {})
    id_mensaje = mensaje["messageId"]
    intento = int(sobre.get("deliveryAttempt", 1))

    try:
        pedido = json.loads(cuerpo)
    except json.JSONDecodeError:
        # 200 sin procesar: confirmamos para que NO se reintente
        # un mensaje que nunca sera valido.
        app.logger.error("Mensaje %s ilegible, descartado", id_mensaje)
        return "", 204

    try:
        crear_orden_de_preparacion(pedido)
    except ErrorTemporal:
        app.logger.warning("Fallo temporal en %s, intento %s",
                           id_mensaje, intento)
        # 500: Pub/Sub reintentara con retroceso exponencial
        return "reintentar", 500

    # 204: confirmado
    return "", 204

La regla de oro del push: el código de respuesta HTTP es la confirmación. 2xx retira el mensaje; cualquier otra cosa (o un tiempo de espera agotado) lo devuelve a la cola. Devolver 200 en un except genérico "para que no moleste" es la forma más rápida de perder datos en silencio.

  1. Confirmaciones, plazos y reintentos

El ack deadline es el tiempo que Pub/Sub espera la confirmación antes de reenviar. Por defecto 10 segundos; configurable entre 10 y 600.

sequenceDiagram
    participant PS as Pub/Sub
    participant C as Consumidor
    PS->>C: entrega mensaje M (deadline 30 s)
    Note over C: procesando... 25 s
    C->>PS: modifyAckDeadline(+60 s)
    Note over C: sigue procesando... 40 s
    C->>PS: ack(M)
    Note over PS: M retirado de esta suscripcion

Si el consumidor no confirma ni amplía el plazo, el mensaje vuelve a la cola. Elegir bien el plazo es importante:

  • Demasiado corto: mensajes que se están procesando correctamente se reenvían y se duplica el trabajo.
  • Demasiado largo: si un consumidor muere, sus mensajes tardan mucho en reasignarse a otro.

Regla práctica: el percentil 99 de tu tiempo de procesamiento, con margen. Si el 99 % de las órdenes de almacén se crean en menos de 8 segundos, 30 segundos es razonable.

La buena noticia es que la biblioteca de Python amplía el plazo automáticamente mientras tu callback sigue ejecutándose (lease management), hasta un máximo configurable. Si necesitas ampliar a mano:

# Ampliar el plazo desde dentro del procesamiento
mensaje.modify_ack_deadline(120)

Los reintentos siguen la política que configuraste:

gcloud pubsub subscriptions update sub-almacen \
  --min-retry-delay=10s \
  --max-retry-delay=600s

Con retroceso exponencial, los reintentos se espacian: 10 s, 20 s, 40 s, 80 s… hasta 600 s. Esto es lo correcto cuando el fallo es porque un sistema está saturado: reintentar cada segundo lo hundiría más. Sin retroceso, un consumidor caído genera una tormenta de reintentos que impide que se recupere.

  1. Temas de mensajes fallidos (dead letter)

Imagina un mensaje con un sku que no existe en el catálogo. El consumidor del almacén falla. Reintenta. Falla. Reintenta. Para siempre. Ese mensaje se llama envenenado, y sin dead letter tiene tres consecuencias: consume recursos indefinidamente, ensucia los logs, y —si hay ordenación activada— bloquea todos los mensajes posteriores de su clave.

El tema de mensajes fallidos resuelve esto: tras N intentos, el mensaje se desvía a otro topic y se retira de la suscripción original.

gcloud pubsub subscriptions update sub-almacen \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5

Cuando un mensaje llega al dead letter, conserva su cuerpo y sus atributos y añade metadatos sobre el origen y el número de intentos. Un proceso de revisión, típicamente semanal:

"""revisar_fallidos.py -- Inspeccion de la cola de mensajes fallidos."""
import json

from google.cloud import pubsub_v1

suscriptor = pubsub_v1.SubscriberClient()
ruta = suscriptor.subscription_path("alpinashop-datos", "sub-pedidos-fallidos")

respuesta = suscriptor.pull(
    request={"subscription": ruta, "max_messages": 100},
    timeout=30,
)

for recibido in respuesta.received_messages:
    m = recibido.message
    print("---")
    print("ID original :", m.attributes.get(
        "CloudPubSubDeadLetterSourceMessageId", "?"))
    print("Suscripcion :", m.attributes.get(
        "CloudPubSubDeadLetterSourceSubscription", "?"))
    print("Intentos    :", m.attributes.get(
        "CloudPubSubDeadLetterSourceDeliveryCount", "?"))
    print("Publicado   :", m.publish_time)
    print("Cuerpo      :", m.data.decode("utf-8")[:300])

    # Decision manual: corregir y republicar, o descartar con registro.

Reglas de dead letter que conviene fijar como política de equipo:

  1. Toda suscripción de producción tiene uno. Sin excepción. Es la diferencia entre "algo falló y lo tenemos guardado" y "algo falló y no sabemos qué era".
  2. max-delivery-attempts entre 5 y 10. Menos, y un fallo transitorio manda mensajes buenos a la cola de fallidos. Más, y se tarda demasiado en detectar el problema.
  3. Una alerta sobre el número de mensajes en el dead letter. Un dead letter que nadie mira es un agujero negro con pasos extra.
  4. No dirijas un dead letter al topic original. Es un bucle infinito, y es un error que se comete más de lo que parece.

  1. Idempotencia: el consumidor debe tolerar duplicados

Ya sabemos que habrá duplicados. La solución no es evitarlos —no se puede— sino que procesar dos veces el mismo mensaje tenga el mismo efecto que procesarlo una vez. Eso es la idempotencia.

Tres estrategias, de menos a más robusta.

A. Operaciones naturalmente idempotentes

La mejor, cuando es posible: diseñar la operación para que repetirla no cambie nada.

# NO idempotente: sumar. Dos veces suma el doble.
db.execute("UPDATE stock SET unidades = unidades - %s WHERE sku = %s",
           (cantidad, sku))

# Idempotente: fijar un valor absoluto calculado desde la fuente de verdad.
db.execute("UPDATE stock SET unidades = %s WHERE sku = %s",
           (unidades_calculadas, sku))

B. Inserción condicional por clave de negocio

Usar una restricción de unicidad de la base de datos:

def crear_orden_de_preparacion(pedido):
    """El indice UNIQUE sobre pedido_id hace el trabajo por nosotros."""
    with db.transaction() as tx:
        tx.execute(
            """
            INSERT INTO ordenes_preparacion (pedido_id, estado, creada_en)
            VALUES (%s, 'pendiente', NOW())
            ON CONFLICT (pedido_id) DO NOTHING
            """,
            (pedido["pedido_id"],),
        )
        # Si ya existia, ON CONFLICT no hace nada y no hay error.
        # El mensaje duplicado se procesa sin efecto: idempotente.

Es la opción preferible cuando existe una clave de negocio natural, como aquí pedido_id. Cero infraestructura añadida y la garantía la da la base de datos.

C. Registro de mensajes procesados

Cuando la operación no es idempotente ni hay clave natural (enviar un correo, llamar a una API externa), hay que llevar un registro. Firestore, que AlpinaShop ya usa para el carrito (02-06), es ideal por su latencia baja y sus transacciones:

"""idempotencia.py -- Registro de mensajes ya procesados."""
import datetime

from google.cloud import firestore

db = firestore.Client(project="alpinashop-datos")
COLECCION = "eventos_procesados"
TTL_DIAS = 14   # mayor que la retencion maxima de Pub/Sub (7 dias)


def procesar_una_sola_vez(clave_idempotencia: str, consumidor: str, accion):
    """Ejecuta `accion` solo si esta clave no se ha procesado antes.

    La clave incluye el consumidor: el mismo pedido debe procesarse
    una vez en el almacen Y una vez en facturacion. Son marcas distintas.
    """
    doc_id = f"{consumidor}__{clave_idempotencia}"
    ref = db.collection(COLECCION).document(doc_id)

    @firestore.transactional
    def _reservar(tx):
        instantanea = ref.get(transaction=tx)
        if instantanea.exists:
            return False          # ya procesado: no hacemos nada
        tx.set(ref, {
            "consumidor": consumidor,
            "clave": clave_idempotencia,
            "procesado_en": firestore.SERVER_TIMESTAMP,
            # Campo de TTL: Firestore borra el documento automaticamente
            "caduca_en": datetime.datetime.utcnow()
                         + datetime.timedelta(days=TTL_DIAS),
        })
        return True

    primera_vez = _reservar(db.transaction())
    if not primera_vez:
        return False

    accion()
    return True

Y su uso:

def procesar_facturacion(mensaje):
    pedido = json.loads(mensaje.data.decode("utf-8"))

    # La clave de idempotencia es el identificador de NEGOCIO,
    # no el message_id de Pub/Sub. Motivo: si el pedido se republica
    # por un reproceso con seek, el message_id sera distinto pero
    # el pedido es el mismo y NO debe facturarse dos veces.
    ejecutado = procesar_una_sola_vez(
        clave_idempotencia=pedido["pedido_id"],
        consumidor="facturacion",
        accion=lambda: emitir_factura(pedido),
    )

    if not ejecutado:
        logger.info("Pedido %s ya facturado, se ignora el duplicado",
                    pedido["pedido_id"])

    mensaje.ack()

El comentario sobre la clave es la parte más importante de todo el apartado. Usar message_id protege contra las reentregas de Pub/Sub, pero no contra un reproceso con seek ni contra una republicación desde la aplicación. Usar el identificador de negocio (pedido_id) protege contra ambos. Es la diferencia entre un sistema que aguanta un reproceso y uno que factura dos veces el día que alguien reproduce una hora de mensajes.

Un detalle práctico: el TTL del registro debe ser mayor que la retención máxima de la suscripción. Si Pub/Sub puede reentregar durante 7 días y tu registro caduca a los 3, un mensaje reentregado el día 5 se procesaría de nuevo.

  1. Retención, seek e instantáneas

Los mensajes se retienen en la suscripción entre 10 minutos y 31 días (7 por defecto). Con seek se mueve el punto de lectura.

# 1) Reprocesar las ultimas 3 horas: util tras corregir un error del consumidor
gcloud pubsub subscriptions seek sub-analitica \
  --time="$(date -u -d '3 hours ago' '+%Y-%m-%dT%H:%M:%SZ')"

# 2) Descartar TODO lo pendiente: util al desatascar una cola inservible
gcloud pubsub subscriptions seek sub-almacen --time="$(date -u '+%Y-%m-%dT%H:%M:%SZ')"

Casos reales en los que seek salva el día:

  • Un error en el consumidor de analítica calculó mal el IVA durante seis horas. Se corrige el código, se despliega, se hace seek a hace seis horas y se reprocesa todo. Con consumidores idempotentes, esto es seguro. Sin idempotencia, es un desastre.
  • Una cola acumuló 400.000 mensajes obsoletos por un consumidor caído durante el fin de semana. seek al presente los descarta de golpe.

Las instantáneas (snapshots) capturan el estado de confirmación de una suscripción para poder volver a él:

# ANTES de desplegar una version arriesgada del consumidor
gcloud pubsub snapshots create snap-analitica-pre-v2 --subscription=sub-analitica

# ... se despliega, y si sale mal ...
gcloud pubsub subscriptions seek sub-analitica --snapshot=snap-analitica-pre-v2

# Limpieza (las instantaneas caducan a los 7 dias, pero mejor ser explicito)
gcloud pubsub snapshots delete snap-analitica-pre-v2

Es el equivalente a un punto de restauración antes de un despliegue de riesgo, y debería formar parte del procedimiento de despliegue de cualquier consumidor que escriba en sistemas de negocio.

  1. Filtros de suscripción por atributo

Una suscripción puede filtrar por los atributos del mensaje, de modo que solo recibe lo que le interesa. El filtrado ocurre en el lado de Pub/Sub: los mensajes descartados ni se entregan ni se facturan como entrega.

# Suscripcion que solo recibe pedidos de Espana
gcloud pubsub subscriptions create sub-almacen-es \
  --topic=pedidos-nuevos \
  --message-filter='attributes.pais = "ES"'

# Solo pedidos grandes, para revision manual antifraude
gcloud pubsub subscriptions create sub-revision-fraude \
  --topic=pedidos-nuevos \
  --message-filter='attributes.tipo_evento = "pedido_confirmado" AND attributes.importe_alto = "true"'

# Todo menos el canal telefonico
gcloud pubsub subscriptions create sub-analitica-digital \
  --topic=pedidos-nuevos \
  --message-filter='NOT (attributes.origen = "telefono")'

# Prefijo: cualquier evento cuyo tipo empiece por "pedido_"
gcloud pubsub subscriptions create sub-todo-pedidos \
  --topic=eventos-tienda \
  --message-filter='hasPrefix(attributes.tipo_evento, "pedido_")'

Operadores disponibles: =, !=, AND, OR, NOT, hasPrefix(), y attributes:clave para comprobar existencia.

Dos limitaciones que hay que conocer:

  • Solo se filtra por atributos, nunca por el cuerpo. Si quieres filtrar por el importe, tienes que publicarlo como atributo. De ahí que el publicador del apartado 5 incluya pais y origen como atributos aunque también estén dentro del JSON.
  • El filtro es inmutable. No se puede modificar una vez creada la suscripción; hay que crear otra.

Los filtros permiten un patrón muy limpio: un topic por dominio, con muchos tipos de evento, y cada consumidor filtra lo suyo. Para AlpinaShop, un topic eventos-tienda con pedido_confirmado, pedido_cancelado, carrito_abandonado y stock_bajo es más manejable que cuatro topics, porque conserva el orden relativo de los eventos del mismo dominio y simplifica el publicador.

  1. Suscripciones de BigQuery y de Cloud Storage

Aquí llega una de las mejores funciones del servicio, y la que más código ahorra.

Una suscripción de BigQuery escribe los mensajes directamente en una tabla. Sin Dataflow, sin código, sin nada que desplegar.

# 1) Tabla destino con el esquema del mensaje
bq mk --table \
  --time_partitioning_field=fecha_pedido \
  --clustering_fields=canal \
  alpinashop-datos:alpinashop_analitica.pedidos_evento \
  pedido_id:STRING,momento:TIMESTAMP,fecha_pedido:DATE,cliente_id:STRING,canal:STRING,total:NUMERIC

# 2) Permiso para que Pub/Sub escriba
PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROY_NUM}@gcp-sa-pubsub.iam.gserviceaccount.com"
bq add-iam-policy-binding --member="serviceAccount:${SA_PUBSUB}" \
  --role="roles/bigquery.dataEditor" \
  alpinashop-datos:alpinashop_analitica.pedidos_evento

# 3) La suscripcion
gcloud pubsub subscriptions create sub-bq-pedidos \
  --topic=pedidos-nuevos \
  --bigquery-table=alpinashop-datos:alpinashop_analitica.pedidos_evento \
  --use-table-schema \
  --drop-unknown-fields \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5
  • --use-table-schema: se espera JSON cuyos campos coincidan con las columnas. La alternativa, --use-topic-schema, valida contra un esquema Avro o Protobuf registrado en el topic.
  • --drop-unknown-fields: los campos del JSON que no existan como columna se ignoran en vez de provocar un fallo. Imprescindible si el esquema del mensaje puede evolucionar.
  • El dead letter aquí es fundamental: sin él, un mensaje que no encaje con el esquema se reintenta indefinidamente.

La comparación honesta:

Criterio Suscripción de BigQuery Dataflow (04-02)
Código Ninguno Un pipeline que mantener
Coste Solo el de Pub/Sub y la escritura Trabajadores 24×7: decenas de €/mes
Transformación Ninguna: el JSON va tal cual Cualquiera
Ventanas y tiempo del evento No Sí
Enriquecimiento con otras fuentes No Sí
Agregación en vuelo No Sí
Cuándo usarla El mensaje ya tiene la forma de la tabla Hay que transformar, agregar o ventanear

La decisión de AlpinaShop: para volcar los eventos crudos a una tabla de aterrizaje, suscripción de BigQuery, porque es gratis en esfuerzo y no hay nada que mantener. El pipeline de Dataflow se reserva para lo que sí necesita transformación: el agregado por hora y canal con ventanas de tiempo del evento y tolerancia a datos tardíos. Es el patrón habitual y correcto: datos crudos por la vía barata, agregados por la vía potente.

La suscripción de Cloud Storage hace lo equivalente con ficheros, agrupando mensajes por tiempo o tamaño:

gcloud pubsub subscriptions create sub-gcs-pedidos \
  --topic=pedidos-nuevos \
  --cloud-storage-bucket=alpinashop-datalake \
  --cloud-storage-file-prefix=eventos/pedidos/ \
  --cloud-storage-file-suffix=.json \
  --cloud-storage-max-duration=300s \
  --cloud-storage-max-bytes=10MB \
  --cloud-storage-output-format=json

Es una forma excelente de tener un archivo inmutable de todos los eventos en el lago, que luego procesa Spark (04-03) o carga BigQuery. Y es la red de seguridad definitiva: si un día hay que reconstruir todo el almacén analítico, los eventos crudos están ahí.

  1. Ordenación por clave y sus implicaciones

Por defecto no hay orden. Para garantizarlo, se publica con clave de ordenación:

publicador = pubsub_v1.PublisherClient(
    publisher_options=pubsub_v1.types.PublisherOptions(enable_message_ordering=True)
)

# Todos los eventos del MISMO pedido comparten clave: llegan en orden
publicador.publish(
    ruta_topic,
    data=cuerpo,
    ordering_key=pedido["pedido_id"],     # <-- la clave
    tipo_evento="pedido_enviado",
)
gcloud pubsub subscriptions create sub-estados-pedido \
  --topic=eventos-tienda \
  --enable-message-ordering

La garantía es: los mensajes con la misma clave de ordenación, publicados en la misma región, se entregan en orden de publicación. Mensajes con claves distintas no tienen relación entre sí, lo que permite seguir paralelizando.

Los tres costes de activar la ordenación, que hay que sopesar:

  1. Menor rendimiento. Los mensajes de una clave se procesan secuencialmente. Con pedido_id como clave no hay problema: hay miles de pedidos distintos y el paralelismo se mantiene. Con pais como clave, todos los pedidos de España irían en serie.
  2. Bloqueo por mensaje envenenado. Si un mensaje de la clave PED-2026-0042 falla repetidamente, todos los mensajes posteriores de ese pedido quedan bloqueados hasta que se resuelva o se desvíe al dead letter. Ordenación sin dead letter es una bomba de relojería.
  3. Publicación más lenta. El cliente debe serializar los envíos de cada clave, lo que reduce el efecto del envío por lotes.

Regla de oro: activa la ordenación solo cuando el orden sea semánticamente necesario, y elige la clave con la máxima cardinalidad posible que preserve ese orden. Para AlpinaShop: pedido_id sí (los estados de un pedido deben ir en orden); pais no.

Y la alternativa que suele ser mejor: hacer que el orden no importe. Si cada evento incluye un número de versión o una marca de tiempo, el consumidor puede descartar los que sean anteriores a lo que ya procesó, y el orden de llegada deja de ser un problema.

# Consumidor tolerante al desorden, sin necesidad de ordering_key
def aplicar_estado(pedido_id, nuevo_estado, version_evento):
    db.execute(
        """
        UPDATE pedidos SET estado = %s, version = %s
        WHERE pedido_id = %s AND version < %s
        """,
        (nuevo_estado, version_evento, pedido_id, version_evento),
    )
    # Si llega un evento antiguo, la condicion version < %s no se cumple
    # y el UPDATE no afecta a ninguna fila. Desorden absorbido.

  1. Pub/Sub Lite y otras alternativas

Pub/Sub Lite fue una variante de menor coste, con capacidad aprovisionada por el usuario en lugar de servida bajo demanda, pensada para volúmenes muy altos y predecibles. Requería dimensionar particiones y almacenamiento a mano, a cambio de un precio por GB considerablemente inferior.

Google anunció su retirada y el servicio dejó de estar disponible en marzo de 2026. Si encuentras referencias a Pub/Sub Lite en documentación o cursos antiguos, están obsoletas; verifica siempre el estado en la documentación oficial. Las alternativas actuales son Pub/Sub estándar o, si necesitas la API de Kafka, Google Cloud Managed Service for Apache Kafka.

La comparación que sí sigue siendo útil:

Servicio Modelo Cuándo tiene sentido
Pub/Sub Global, sin aprovisionar, pago por uso El valor por defecto: la inmensa mayoría de los casos
Managed Kafka Kafka gestionado, con particiones y grupos de consumidores Ya tienes código Kafka, o necesitas su semántica de log
Cloud Tasks Cola de tareas con programación y control de ritmo Encolar trabajo con un destinatario único y conocido
Eventarc Enrutado de eventos de la plataforma Reaccionar a eventos de servicios de Google (usa Pub/Sub por debajo)
Memorystore (Redis Pub/Sub) Mensajería en memoria, sin persistencia Notificaciones efímeras donde perder mensajes es aceptable

La confusión más habitual es Pub/Sub frente a Cloud Tasks. Pub/Sub es para hechos que interesan a N consumidores desconocidos. Cloud Tasks es para encargos dirigidos a un destinatario concreto, con control fino del ritmo de entrega y programación diferida. "Se ha confirmado un pedido" es Pub/Sub. "Envía este correo dentro de 30 minutos, con un máximo de 10 por segundo" es Cloud Tasks.

  1. Monitorización y coste

Las métricas que hay que vigilar, con sus alertas:

Métrica Qué indica Alerta razonable
subscription/oldest_unacked_message_age Antigüedad del mensaje pendiente más antiguo La métrica reina. > 600 s: el consumidor no da abasto o está caído
subscription/num_undelivered_messages Mensajes acumulados sin entregar Crecimiento sostenido: los consumidores van por detrás
subscription/dead_letter_message_count Mensajes desviados a fallidos > 0 en una hora: hay que mirarlo
topic/send_request_count Publicaciones Caída brusca: la aplicación ha dejado de publicar
subscription/push_request_count por código Salud del endpoint push Muchos 5xx: el endpoint falla
subscription/ack_message_count Ritmo de confirmación Comparar con publicaciones

La alerta imprescindible:

gcloud alpha monitoring policies create \
  --notification-channels="$CANAL_INFRA" \
  --display-name="Pub/Sub: mensajes sin confirmar en pedidos-nuevos" \
  --condition-display-name="oldest_unacked_message_age > 10 min" \
  --condition-threshold-value=600 \
  --condition-threshold-duration=300s \
  --condition-filter='metric.type="pubsub.googleapis.com/subscription/oldest_unacked_message_age" AND resource.type="pubsub_subscription"'

oldest_unacked_message_age es la mejor métrica porque detecta a la vez los tres fallos posibles: consumidor caído (crece indefinidamente), consumidor lento (crece despacio) y mensaje envenenado bloqueando una clave ordenada (se queda clavado en un valor alto).

El coste se factura principalmente por volumen de datos:

Concepto Orden de magnitud (verificar en la documentación oficial)
Publicación + entrega ~40 $ por TiB, con una franja mensual gratuita
Almacenamiento de mensajes retenidos ~0,27 $ por GiB y mes
Egreso entre regiones Tarifas de red

Un cálculo realista para AlpinaShop: 1.200 pedidos al mes, mensajes de ~2 KB, con 4 suscripciones. Eso son 1.200 publicaciones y 4.800 entregas: unos 12 MB mensuales. Céntimos, o directamente dentro de la franja gratuita. Los eventos de navegación del catálogo, mucho más numerosos, seguirían siendo baratos.

Un detalle de facturación que sorprende: cada suscripción cuenta como una entrega. Diez suscripciones sobre el mismo topic multiplican por diez el volumen facturado de entrega. No es motivo para no usar suscripciones, pero sí para no dejar suscripciones huérfanas: una suscripción olvidada sin consumidor acumula mensajes, se factura su almacenamiento y no sirve para nada.

# Buscar suscripciones sin consumo: candidatas a borrar
gcloud pubsub subscriptions list --format="table(name, topic)" | while read -r s _; do
  echo "$s"
done

  1. Las notificaciones del bucket alpinashop-catalogo

En 02-02 quedó pendiente una promesa: que Cloud Storage puede avisar cuando aparece un objeto nuevo. Ahora tenemos las piezas para cerrarla de verdad.

El caso de AlpinaShop: cuando alguien sube la imagen original de un producto a productos/<sku>/original/, hay que generar automáticamente las versiones web/ y thumb/. Hoy eso lo hace Marta a mano con un script.

# 1) Topic para los eventos del bucket
gcloud pubsub topics create imagenes-subidas \
  --labels=entorno=produccion,equipo=infra,aplicacion=catalogo

# 2) Permiso para que el agente de Cloud Storage publique
SA_GCS=$(gcloud storage service-agent --project=alpinashop-prod)
gcloud pubsub topics add-iam-policy-binding imagenes-subidas \
  --member="serviceAccount:${SA_GCS}" --role="roles/pubsub.publisher"

# 3) La notificacion, acotada al prefijo y al evento que interesan
gcloud storage buckets notifications create gs://alpinashop-catalogo \
  --topic=projects/alpinashop-datos/topics/imagenes-subidas \
  --event-types=OBJECT_FINALIZE \
  --object-prefix=productos/ \
  --payload-format=json

# 4) Suscripcion para el procesador de imagenes
gcloud pubsub subscriptions create sub-procesar-imagenes \
  --topic=imagenes-subidas \
  --ack-deadline=300 \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5

Tipos de evento disponibles:

Evento Cuándo se dispara
OBJECT_FINALIZE Se crea un objeto nuevo o se sobrescribe
OBJECT_DELETE Se borra (o se sobrescribe la versión anterior)
OBJECT_ARCHIVE Una versión pasa a archivada (con versionado activo)
OBJECT_METADATA_UPDATE Cambian los metadatos

El consumidor:

def procesar_imagen(mensaje):
    """Genera las versiones web y thumb de una imagen recien subida."""
    # Los datos del objeto vienen en los ATRIBUTOS, no hace falta el cuerpo
    bucket = mensaje.attributes["bucketId"]
    nombre = mensaje.attributes["objectId"]
    generacion = mensaje.attributes["objectGeneration"]

    # 1) FILTRO DE BUCLE INFINITO: si no comprobamos esto, al escribir
    #    web/ y thumb/ se dispararian nuevas notificaciones que generarian
    #    mas imagenes, indefinidamente. Es el error clasico.
    if "/original/" not in nombre:
        mensaje.ack()
        return

    if not nombre.lower().endswith((".jpg", ".jpeg", ".png", ".webp")):
        mensaje.ack()
        return

    # 2) IDEMPOTENCIA: la generacion identifica la version exacta del objeto.
    #    Si el mensaje se reentrega, la clave es la misma y no se repite.
    clave = f"{bucket}/{nombre}#{generacion}"

    procesar_una_sola_vez(
        clave_idempotencia=clave,
        consumidor="procesador-imagenes",
        accion=lambda: generar_versiones(bucket, nombre),
    )
    mensaje.ack()

El filtro del bucle infinito merece subrayarse: es el fallo más caro de este patrón. Escribir en el mismo bucket que dispara la notificación genera una recursión que solo se detecta cuando llega la factura. Las tres defensas son: filtrar por prefijo en la notificación (--object-prefix=productos/), comprobar la ruta en el consumidor, y —mejor todavía— escribir la salida en un bucket distinto del de entrada.

Este mismo patrón se aplica al fichero mensual del transportista y a las exportaciones nocturnas: en cuanto el fichero aterriza en exportaciones/<yyyy>/<mm>/<dd>/, una notificación dispara la carga en BigQuery. Sin cron, sin comprobar carpetas cada cinco minutos, sin ventanas de espera. La lógica del procesamiento la pondremos en Cloud Functions en 06-03, y quién lo orquesta todo es la lección siguiente.

Errores Comunes y Consejos

Suponer que no habrá duplicados. Los habrá, con o sin exactamente una vez. Todo consumidor de producción debe ser idempotente. No es opcional.

Usar message_id como clave de idempotencia. Protege contra reentregas, pero no contra seek ni contra republicaciones. Usa el identificador de negocio.

Configurar dead letter y olvidar los permisos. El comando se acepta y la función no opera: los mensajes se reintentan para siempre. Comprueba siempre los dos bindings.

No poner control de flujo en el consumidor. Con procesamiento lento, la biblioteca acepta más mensajes de los que puede confirmar, vencen los plazos, se reentregan, y el consumidor se hunde procesando duplicados de su propio atasco.

Crear el cliente de Pub/Sub dentro del manejador de la petición. Es caro (abre conexiones gRPC). Uno por proceso, reutilizado.

Devolver 200 en un except genérico en push. Confirma el mensaje y lo pierde en silencio. Devuelve 5xx en los fallos transitorios.

Activar ordenación sin necesitarla. Reduce el rendimiento y, con un mensaje envenenado, bloquea toda la clave.

Dirigir el dead letter al topic original. Bucle infinito.

Suscripciones huérfanas. Sin consumidor, acumulan mensajes hasta la retención máxima, se factura su almacenamiento y no aportan nada. Revísalas periódicamente.

Notificaciones de bucket que se retroalimentan. Escribir en el mismo bucket que dispara el evento. Usa buckets o prefijos separados y filtra en el consumidor.

Consejo: publica atributos generosamente. Cuestan poco y permiten filtrar del lado del servidor, lo que ahorra entregas y complejidad en el consumidor.

Consejo: un esquema en el topic. Pub/Sub admite registrar un esquema Avro o Protobuf y validar los mensajes al publicar. Convierte los errores de formato en fallos inmediatos del publicador en vez de sorpresas en el consumidor.

Consejo: haz una instantánea antes de cada despliegue arriesgado. Cuesta un comando y da marcha atrás.

Ejercicios

Ejercicio 1: topic de opiniones con filtros y dead letter

Crea el topic opiniones-nuevas con 7 días de retención, su topic de fallidos opiniones-nuevas-fallidos con una suscripción de inspección, y tres suscripciones sobre el topic principal:

  • sub-moderacion: recibe solo las opiniones con attributes.puntuacion_baja = "true", con plazo de 60 s, dead letter tras 5 intentos y retroceso de 10 s a 300 s.
  • sub-analitica-opiniones: recibe todas, con plazo de 120 s.
  • sub-bq-opiniones: escribe directamente en BigQuery en la tabla alpinashop-datos:alpinashop_analitica.opiniones, sin código.

Incluye todos los permisos necesarios y publica dos mensajes de prueba, uno que llegue a moderación y otro que no.

Ejercicio 2: consumidor idempotente de facturación

Escribe el consumidor pull de sub-facturacion que emite la factura de cada pedido. Requisitos: control de flujo de 20 mensajes en vuelo; distinguir errores transitorios (reintentar) de permanentes (descartar con registro); idempotencia basada en el identificador de negocio con registro en Firestore y TTL adecuado; apagado limpio ante SIGTERM; y un contador de facturas emitidas y duplicados detectados. Explica por qué la clave de idempotencia debe ser el pedido_id y no el message_id, con un escenario concreto en el que la elección incorrecta causaría un problema real.

Ejercicio 3: diagnóstico de una cola atascada

A las 09:15 salta la alerta: oldest_unacked_message_age en sub-almacen lleva 47 minutos y sube. num_undelivered_messages ha pasado de 3 a 12.400 desde las 08:20. El servicio de almacén está levantado y responde a su comprobación de estado. En los logs del consumidor aparece, unas cien veces por minuto, el mismo mensaje de error mencionando el pedido PED-2026-1188. La cola de pedidos-nuevos-fallidos está vacía. La suscripción tiene la ordenación activada por pedido_id.

Diagnostica qué está pasando, explica por qué el dead letter está vacío pese a los reintentos, y da un plan de actuación en tres fases: contención inmediata, corrección y prevención.

Soluciones

Solución 1

gcloud config set project alpinashop-datos
PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROY_NUM}@gcp-sa-pubsub.iam.gserviceaccount.com"

# 1) Topics
gcloud pubsub topics create opiniones-nuevas \
  --message-retention-duration=7d \
  --labels=entorno=produccion,equipo=datos,aplicacion=catalogo

gcloud pubsub topics create opiniones-nuevas-fallidos \
  --labels=entorno=produccion,equipo=datos,aplicacion=catalogo

gcloud pubsub subscriptions create sub-opiniones-fallidas \
  --topic=opiniones-nuevas-fallidos \
  --message-retention-duration=31d --ack-deadline=600

# 2) Permiso de publicacion en el topic de fallidos
gcloud pubsub topics add-iam-policy-binding opiniones-nuevas-fallidos \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.publisher"

# 3) Suscripcion de moderacion, con filtro
gcloud pubsub subscriptions create sub-moderacion \
  --topic=opiniones-nuevas \
  --message-filter='attributes.puntuacion_baja = "true"' \
  --ack-deadline=60 \
  --dead-letter-topic=opiniones-nuevas-fallidos \
  --max-delivery-attempts=5 \
  --min-retry-delay=10s --max-retry-delay=300s \
  --labels=consumidor=moderacion

gcloud pubsub subscriptions add-iam-policy-binding sub-moderacion \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.subscriber"

# 4) Suscripcion de analitica, sin filtro
gcloud pubsub subscriptions create sub-analitica-opiniones \
  --topic=opiniones-nuevas --ack-deadline=120 \
  --labels=consumidor=analitica

# 5) Suscripcion directa a BigQuery
bq add-iam-policy-binding --member="serviceAccount:${SA_PUBSUB}" \
  --role="roles/bigquery.dataEditor" \
  alpinashop-datos:alpinashop_analitica.opiniones

gcloud pubsub subscriptions create sub-bq-opiniones \
  --topic=opiniones-nuevas \
  --bigquery-table=alpinashop-datos:alpinashop_analitica.opiniones \
  --use-table-schema --drop-unknown-fields \
  --dead-letter-topic=opiniones-nuevas-fallidos \
  --max-delivery-attempts=5
# Mensaje que SI llega a moderacion (puntuacion 1)
gcloud pubsub topics publish opiniones-nuevas \
  --message='{"opinion_id":"OPI-1001","sku":"FRON-300L","fecha":"2026-03-14","puntuacion":1,"texto":"La bateria dura muchisimo menos de lo anunciado","pais":"FR"}' \
  --attribute=puntuacion_baja=true,sku=FRON-300L,pais=FR

# Mensaje que NO llega a moderacion (puntuacion 5)
gcloud pubsub topics publish opiniones-nuevas \
  --message='{"opinion_id":"OPI-1002","sku":"MOCH-40L-AZ","fecha":"2026-03-14","puntuacion":5,"texto":"Comodisima en travesias largas","pais":"ES"}' \
  --attribute=puntuacion_baja=false,sku=MOCH-40L-AZ,pais=ES

# Comprobacion: moderacion recibe 1, analitica recibe 2
gcloud pubsub subscriptions pull sub-moderacion --limit=5 --format="value(message.data)"
gcloud pubsub subscriptions pull sub-analitica-opiniones --limit=5 --format="value(message.data)"

El punto que se evalúa aquí es el filtro: sub-moderacion recibe una de las dos publicaciones, porque el filtrado ocurre del lado de Pub/Sub y la opinión de 5 estrellas ni siquiera se entrega ni se factura.

Solución 2

"""consumidor_facturacion.py -- Emision idempotente de facturas."""
import datetime
import json
import logging
import signal
import sys

from google.cloud import firestore, pubsub_v1

PROYECTO = "alpinashop-datos"
SUSCRIPCION = "sub-facturacion"
COLECCION = "eventos_procesados"
TTL_DIAS = 14           # > 7 dias de retencion maxima de la suscripcion
CONSUMIDOR = "facturacion"

logger = logging.getLogger(__name__)
db = firestore.Client(project=PROYECTO)

facturas_emitidas = 0
duplicados_detectados = 0


class ErrorTemporal(Exception):
    """Fallo transitorio: merece reintento."""


class ErrorPermanente(Exception):
    """Fallo que no mejorara reintentando."""


def procesar_una_sola_vez(clave, accion):
    doc = db.collection(COLECCION).document(f"{CONSUMIDOR}__{clave}")

    @firestore.transactional
    def _reservar(tx):
        if doc.get(transaction=tx).exists:
            return False
        tx.set(doc, {
            "consumidor": CONSUMIDOR,
            "clave": clave,
            "procesado_en": firestore.SERVER_TIMESTAMP,
            "caduca_en": datetime.datetime.utcnow()
                         + datetime.timedelta(days=TTL_DIAS),
        })
        return True

    if not _reservar(db.transaction()):
        return False
    accion()
    return True


def procesar(mensaje):
    global facturas_emitidas, duplicados_detectados
    try:
        pedido = json.loads(mensaje.data.decode("utf-8"))
    except json.JSONDecodeError:
        logger.error("Mensaje %s ilegible, descartado", mensaje.message_id)
        mensaje.ack()                      # permanente: no reintentar
        return

    if not pedido.get("pedido_id"):
        logger.error("Mensaje %s sin pedido_id, descartado", mensaje.message_id)
        mensaje.ack()
        return

    try:
        emitida = procesar_una_sola_vez(
            clave=pedido["pedido_id"],
            accion=lambda: emitir_factura(pedido),
        )
        if emitida:
            facturas_emitidas += 1
            logger.info("Factura emitida para %s", pedido["pedido_id"])
        else:
            duplicados_detectados += 1
            logger.info("Duplicado ignorado: %s", pedido["pedido_id"])
        mensaje.ack()

    except ErrorTemporal as exc:
        logger.warning("Fallo temporal en %s: %s", pedido["pedido_id"], exc)
        mensaje.nack()                     # reintento con retroceso

    except ErrorPermanente as exc:
        logger.error("Fallo permanente en %s: %s", pedido["pedido_id"], exc)
        mensaje.ack()                      # no tiene arreglo: se retira

    except Exception:
        logger.exception("Fallo desconocido en %s", pedido["pedido_id"])
        mensaje.nack()                     # tras 5 intentos, al dead letter


def main():
    suscriptor = pubsub_v1.SubscriberClient()
    ruta = suscriptor.subscription_path(PROYECTO, SUSCRIPCION)
    control = pubsub_v1.types.FlowControl(max_messages=20)

    futuro = suscriptor.subscribe(ruta, callback=procesar, flow_control=control)
    logger.info("Facturacion escuchando en %s", ruta)

    def apagar(signum, frame):
        logger.info("Cerrando. Emitidas=%s Duplicados=%s",
                    facturas_emitidas, duplicados_detectados)
        futuro.cancel()
        futuro.result(timeout=30)
        sys.exit(0)

    signal.signal(signal.SIGTERM, apagar)
    signal.signal(signal.SIGINT, apagar)

    with suscriptor:
        try:
            futuro.result()
        except Exception:
            logger.exception("Suscriptor caido")
            futuro.cancel()


if __name__ == "__main__":
    logging.basicConfig(level=logging.INFO)
    main()

Por qué pedido_id y no message_id, con escenario concreto:

message_id es único por publicación. Si el mismo hecho de negocio se publica dos veces, cada publicación tendrá un message_id distinto, y un registro basado en él no detectaría nada.

El escenario real: un martes, un fallo en el consumidor de facturación provocó que 300 pedidos de la mañana no se facturaran. Se corrige el código, se despliega y se ejecuta:

gcloud pubsub subscriptions seek sub-facturacion \
  --time="2026-03-17T08:00:00Z"

Ese seek reentrega los mensajes de la mañana. Ahora bien, entre las 08:00 y el incidente hubo 500 pedidos, de los cuales 200 sí se facturaron correctamente antes de que empezara el fallo.

  • Con clave message_id: la reentrega conserva el message_id original, así que en este caso concreto los 200 buenos sí se detectarían. Pero si en lugar de seek alguien recupera y republica los eventos desde el archivo de Cloud Storage (apartado 13), los message_id serán nuevos y se emitirían 200 facturas duplicadas. Ese es el fallo.
  • Con clave pedido_id: da igual cómo llegue el mensaje —reentrega, seek, republicación desde el archivo, reproceso manual—. Si ese pedido ya se facturó, no se vuelve a facturar. La protección es sobre el hecho de negocio, que es lo que realmente importa.

Doscientas facturas duplicadas enviadas a clientes es un incidente contable, de atención al cliente y de reputación. La diferencia entre las dos opciones es una línea de código.

Solución 3

Diagnóstico: mensaje envenenado bloqueando una clave de ordenación, con el dead letter mal configurado.

Las pistas encajan una a una:

  1. El servicio está vivo y responde. No es una caída: es un atasco lógico.
  2. El mismo error cien veces por minuto sobre PED-2026-1188. Ese mensaje falla siempre. Es un mensaje envenenado: algo en él (un SKU inexistente, un campo nulo, un importe imposible) hace fallar el consumidor de forma determinista.
  3. La ordenación está activada por pedido_id. Aquí está la parte crítica que explica la magnitud: con ordenación, los mensajes de la misma clave se entregan en orden y uno bloqueado impide avanzar. Pero además, en la práctica, la combinación de un consumidor que hace nack() en bucle sobre una clave y el control de flujo saturado por reentregas hace que el rendimiento global se desplome: el consumidor gasta su capacidad reprocesando el mismo mensaje una y otra vez en lugar de atender la cola.
  4. 12.400 mensajes acumulados en 55 minutos cuando AlpinaShop hace ~1.200 pedidos al mes: ese volumen no son pedidos nuevos, son reentregas del mismo mensaje más el resto de la cola sin atender.

Por qué el dead letter está vacío. Es la parte que más enseña. Hay dos causas posibles y ambas son frecuentes:

  • Faltan los permisos. La política de dead letter se configura sin error aunque el agente de servicio de Pub/Sub no tenga roles/pubsub.publisher sobre el topic de fallidos ni roles/pubsub.subscriber sobre la suscripción de origen. Sin esos dos permisos, el desvío nunca ocurre y el mensaje se reintenta indefinidamente. Es exactamente el fallo silencioso que advertimos en el apartado 4.
  • La política no está aplicada realmente. Se creó la suscripción sin --dead-letter-topic y se dio por hecho.

Comprobación:

gcloud pubsub subscriptions describe sub-almacen \
  --format="yaml(deadLetterPolicy, enableMessageOrdering, ackDeadlineSeconds)"

gcloud pubsub topics get-iam-policy pedidos-nuevos-fallidos
gcloud pubsub subscriptions get-iam-policy sub-almacen

Plan de actuación en tres fases:

Fase 1 — Contención (minutos). Arreglar los permisos del dead letter, que es lo que desatasca sin perder nada:

PROY_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROY_NUM}@gcp-sa-pubsub.iam.gserviceaccount.com"

gcloud pubsub topics add-iam-policy-binding pedidos-nuevos-fallidos \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.publisher"
gcloud pubsub subscriptions add-iam-policy-binding sub-almacen \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.subscriber"

gcloud pubsub subscriptions update sub-almacen \
  --dead-letter-topic=pedidos-nuevos-fallidos --max-delivery-attempts=5

En cuanto el desvío funcione, PED-2026-1188 saldrá de la cola tras cinco intentos y el resto de mensajes empezará a fluir. No hacer seek al presente: eso descartaría 12.400 mensajes que en su mayoría son pedidos reales pendientes de preparar.

Si el desvío tardara y hubiera urgencia operativa, el parche alternativo es desplegar en el consumidor una regla temporal que confirme explícitamente el pedido_id problemático, dejándolo registrado para tratarlo a mano.

Fase 2 — Corrección (horas). Inspeccionar el mensaje en la cola de fallidos, entender por qué el consumidor falla con él, y corregir el código para que ese tipo de dato se trate como error permanente (ack() con registro) en lugar de reintentarse. Después, tratar manualmente el pedido PED-2026-1188, que sigue siendo un pedido real de un cliente real y hay que prepararlo. Verificar que la cola vuelve a un oldest_unacked_message_age de segundos.

Fase 3 — Prevención (días).

  • Verificar el dead letter de todas las suscripciones, incluidos los permisos, no solo la política. Automatizarlo como comprobación en el despliegue.
  • Alerta sobre dead_letter_message_count > 0, para enterarse del primer mensaje envenenado en vez del número 12.400.
  • Revisar si la ordenación es realmente necesaria en sub-almacen. Los pedidos son independientes entre sí; el orden entre pedidos distintos no aporta nada y sí añade riesgo de bloqueo. Si solo hace falta ordenar los cambios de estado de un mismo pedido, ese caso puede resolverse con el número de versión y un UPDATE ... WHERE version < %s, eliminando la ordenación por completo.
  • Clasificar errores en el consumidor: transitorio → nack(), permanente → ack() con registro. La ausencia de esa distinción es la causa raíz de que un solo dato malo tumbe una cola.
  • Añadir la alerta de oldest_unacked_message_age a 10 minutos, no a 47.

Conclusión

AlpinaShop ya no es un monolito que llama a todo el mundo por teléfono. En esta lección has visto el problema del acoplamiento con números concretos —una compra que espera 5,4 segundos a cuatro sistemas que al cliente no le importan— y cómo se disuelve publicando un hecho en lugar de dar órdenes.

Has interiorizado el modelo: el topic es el tipo de hecho, la suscripción es donde viven realmente los mensajes y cada una recibe su copia independiente, el mensaje lleva cuerpo y atributos filtrables, y la confirmación es lo único que retira un mensaje del sistema. Y sobre todo has visto las garantías reales, que son las que importan: entrega al menos una vez, con duplicados que ocurrirán aunque tu código sea perfecto; sin orden salvo con clave de ordenación; y exactamente una vez como una reducción del problema, no como su desaparición.

Has creado pedidos-nuevos con retención de 7 días, el topic de fallidos pedidos-nuevos-fallidos con su suscripción de inspección, y las suscripciones sub-almacen, sub-facturacion —con exactamente una vez, porque duplicar una factura es grave— y sub-analitica, cada una con su plazo, su política de reintentos y sus permisos de dead letter, incluidos los dos bindings que todo el mundo olvida y sin los cuales el desvío no funciona en silencio.

Has publicado desde Flask con envío por lotes, retroceso exponencial y un cliente reutilizado, decidiendo conscientemente publicar un evento gordo con las líneas dentro y marcando version_esquema y momento_evento como atributos. Has consumido en pull con control de flujo, apagado limpio y distinción entre errores transitorios y permanentes; y en push hacia un endpoint HTTPS con autenticación OIDC, sabiendo que ahí el código de respuesta HTTP es la confirmación. Conoces el ack deadline, modifyAckDeadline y por qué el plazo se ajusta al percentil 99 del procesamiento.

Has entendido por qué todo consumidor serio necesita un dead letter —el mensaje envenenado que se reintenta para siempre— y por qué la clave de idempotencia debe ser el identificador de negocio y no el message_id: un seek o una republicación desde el archivo emitirían facturas duplicadas con la elección equivocada. Has visto seek e instantáneas como red de seguridad para reprocesos y despliegues arriesgados, los filtros por atributo que descartan del lado del servidor, y las suscripciones de BigQuery y Cloud Storage, que ingieren sin una sola línea de código y que para AlpinaShop se llevan los datos crudos mientras Dataflow se queda con lo que de verdad necesita transformación y ventanas.

Sabes cuándo activar ordenación y a qué precio, que Pub/Sub Lite ya no existe y qué usar en su lugar, qué métrica vigilar por encima de todas —oldest_unacked_message_age— y que el coste, con el volumen de AlpinaShop, es de céntimos. Y has cerrado por fin la promesa de 02-02: el bucket alpinashop-catalogo ya avisa cuando aparece una imagen nueva, con el filtro de prefijo, la comprobación de ruta y la idempotencia por generación de objeto que impiden el bucle infinito que arruina a los incautos.

Fíjate en el mapa que tenemos ahora. BigQuery guarda y responde. Dataflow transforma en lote y en streaming. Dataproc calcula lo algorítmico. Pub/Sub conecta los sistemas sin que se conozcan. Es mucha maquinaria. Y sin embargo, hay dos huecos evidentes.

El primero: los datos que no nacen en Google Cloud. El ERP del almacén de AlpinaShop es un MySQL que corre en un servidor de las oficinas de Sabadell. El transportista envía un CSV mensual por correo. El proveedor de mochilas ofrece una API REST con su catálogo. Nada de eso publica en Pub/Sub ni escribe Parquet en un bucket, y Lucía, que sabe SQL pero no ha escrito Beam en su vida, no puede depender de Dani para cada fichero nuevo.

En 04-05, Cloud Data Fusion, veremos la herramienta pensada exactamente para eso: integración de datos visual, sin código, con conectores para casi cualquier origen, un limpiador interactivo llamado Wrangler donde se corrigen fechas y nulos viendo los datos, linaje a nivel de campo que responde a "¿de dónde sale esta columna?", y replicación con captura de cambios desde ese MySQL de Sabadell. Y también veremos su lado incómodo —cuesta por hora de instancia y eso condiciona su uso—, con una tabla honesta de cuándo Data Fusion, cuándo Dataflow, cuándo Datastream y cuándo basta un simple bq load.

Curso de Google Cloud Platform (GCP)

Módulo 1: Introducción a Google Cloud Platform

Módulo 2: Servicios principales de GCP

Módulo 3: Redes y seguridad

Módulo 4: Datos y análisis

Módulo 5: Aprendizaje automático e IA

Módulo 6: DevOps y monitoreo

Módulo 7: Temas avanzados de GCP

Módulo 8: Proyecto final

© Copyright 2026. Todos los derechos reservados