Tornem un moment a l'aplicació d'AlpinaShop. Un client prem "Confirmar comanda". La funció de Flask que atén aquesta petició ha de fer diverses coses: desar la comanda a alpinashop-pedidos, cobrar a la passarel·la de pagament, avisar el magatzem perquè la prepari, avisar facturació perquè emeti la factura, enviar el correu de confirmació, descomptar estoc, i —des del mòdul 4— informar l'analítica.

Avui, aquest codi és una llista de crides, una darrere l'altra. I aquesta llista té tres problemes que no es veuen fins al pitjor dia de l'any.

El client espera tothom. El temps de resposta de la compra és la suma dels set passos. Si el servei de correu triga quatre segons, el client veu una roda girant quatre segons de més per una cosa que no li importa: ell ja ha comprat.

Una fallada ho tomba tot. Si el sistema del magatzem s'està reiniciant, la crida falla. I ara què? Es desfà la comanda, que ja està cobrada? Es continua i es perd l'avís? Es reintenta i es corre el risc de duplicar la factura, que sí que es va enviar? No hi ha cap resposta bona.

Cada novetat toca el codi de la compra. Afegir l'analítica significa modificar, provar i desplegar la funció més crítica de la botiga. I després arribarà el programa de fidelització, i l'antifrau, i l'avís al proveïdor. La funció de compra creix sense parar i cada canvi arrisca la caixa registradora.

Cloud Pub/Sub resol els tres d'una vegada amb una idea molt simple: en comptes de cridar set sistemes, la botiga publica un fet —"s'ha confirmat la comanda PED-2026-0042"— i continua amb el seu. Qui estigui interessat en aquell fet, s'hi subscriu. La botiga no sap ni li importa quants són.

En aquesta lliçó crearàs el topic pedidos-nuevos i les seves subscripcions, publicaràs des de Flask, consumiràs de les quatre maneres possibles, entendràs les garanties reals del servei —que no són les que la gent suposa— i aprendràs per què la idempotència no és un adorn d'arquitectura sinó un requisit ineludible.

Contingut

  1. El problema de l'acoblament, abans i després
  2. El model publicació-subscripció
  3. Les garanties reals de Pub/Sub
  4. Crear el topic pedidos-nuevos i les seves subscripcions
  5. Publicar des de l'aplicació Flask
  6. Consum pull amb el client asíncron
  7. Consum push a un endpoint HTTPS amb OIDC
  8. Confirmacions, terminis i reintents
  9. Temes de missatges fallits (dead letter)
  10. Idempotència: el consumidor ha de tolerar duplicats
  11. Retenció, seek i instantànies
  12. Filtres de subscripció per atribut
  13. Subscripcions de BigQuery i de Cloud Storage
  14. Ordenació per clau i les seves implicacions
  15. Pub/Sub Lite i altres alternatives
  16. Monitoratge i cost
  17. Les notificacions del bucket alpinashop-catalogo

  1. El problema de l'acoblament, abans i després

flowchart TD
    subgraph Abans["ABANS: crides sincrones encadenades"]
        W1["Flask: confirmar_comanda()"]
        A1["Magatzem"]
        F1["Facturacio"]
        M1["Correu"]
        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 Despres["DESPRES: un fet publicat, N interessats"]
        W2["Flask: confirmar_comanda()"]
        T["Topic pedidos-nuevos"]
        S1["sub-almacen"]
        S2["sub-facturacion"]
        S3["sub-analitica"]
        S4["sub-email"]
        A2["Servei de magatzem"]
        F2["Servei de facturacio"]
        AN2["Dataflow -> BigQuery"]
        M2["Servei de correu"]

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

El que canvia en concret:

Aspecte Crides encadenades Amb Pub/Sub
Latència de la compra Suma de tots els passos (~5,4 s) Només la publicació (~15 ms)
Un consumidor caigut Trenca la compra El missatge espera a la seva subscripció
Afegir un consumidor Modificar i desplegar la botiga Crear una subscripció. Zero canvis
Pic de trànsit Cada sistema ha d'aguantar el pic Pub/Sub l'absorbeix; cadascú consumeix al seu ritme
Reprocessar un dia Impossible sense scripts ad hoc seek a un instant anterior
Traçabilitat Registres repartits Mètriques per subscripció

L'expressió tècnica d'això és desacoblament: el productor no coneix els consumidors, no sap quants n'hi ha, no espera la seva resposta i no falla si fallen. El preu és que el sistema passa a ser eventualment consistent: quan la botiga respon "comanda confirmada", el magatzem encara no ho sap. Ho sabrà en uns mil·lisegons, o en uns segons si estava ocupat. Aquest matís cal acceptar-lo conscientment, perquè canvia com es dissenyen les pantalles: no pots mostrar "preparant enviament" immediatament després de comprar si el magatzem encara no se n'ha assabentat.

  1. El model publicació-subscripció

Quatre conceptes, i convé precisar-los perquè el vocabulari es fa servir malament sovint.

Topic (tema). El canal on es publica. Representa un tipus de fet: pedidos-nuevos, imagenes-subidas, stock-agotado. Un topic no desa res per si mateix; és un punt d'entrada.

Subscripció. La cua d'un consumidor concret sobre un topic. Aquí és on viuen realment els missatges. Cada subscripció rep una còpia independent de cada missatge publicat, i porta la seva pròpia comptabilitat de què ha confirmat i què no.

Això últim és la clau que més costa interioritzar:

Es publica 1 missatge a pedidos-nuevos
   └── sub-almacen      rep la seva copia  → la confirma als 0,2 s
   └── sub-facturacion  rep la seva copia  → la confirma als 0,5 s
   └── sub-analitica    rep la seva copia  → el seu consumidor es caigut: espera 6 hores

Que el magatzem confirmi no afecta gens la còpia d'analítica. Són cues separades alimentades pel mateix topic.

I el corol·lari pràctic: si connectes dues instàncies del teu servei de magatzem a la mateixa subscripció sub-almacen, cada missatge anirà a una de les dues. Això és repartiment de càrrega, i és el que vols per escalar. Si en canvi crees dues subscripcions diferents per al mateix servei, cada instància processarà tots els missatges i la feina es farà dues vegades. Una subscripció per consumidor lògic, tantes instàncies com vulguis a dins.

Missatge. Té tres parts:

  • data: el cos, bytes (típicament JSON codificat en UTF-8). Màxim 10 MB.
  • attributes: parells clau-valor de text, fins a 100. Són metadades, i la seva gran virtut és que es poden filtrar sense obrir el cos.
  • Camps del sistema: messageId (únic, assignat per Pub/Sub), publishTime, orderingKey opcional.

Confirmació (ack). El consumidor diu "processat, no me'l tornis a enviar". Fins que arriba aquesta confirmació, Pub/Sub considera el missatge pendent i el reenviarà.

  1. Les garanties reals de Pub/Sub

Aquest apartat és el més important de la lliçó, perquè la majoria dels errors de producció amb missatgeria vénen de suposar garanties que no existeixen.

Garantia La dona Pub/Sub? Conseqüència pràctica
Lliurament almenys una vegada Sí, per defecte El teu consumidor rebrà duplicats. Cal assumir-ho
Lliurament com a màxim una vegada No Mai no es perd un missatge per disseny
Exactament una vegada Sí, si s'activa a la subscripció (amb condicions) Redueix, no elimina, la necessitat d'idempotència
Ordre d'arribada No, llevat que hi hagi clau d'ordenació Els missatges poden arribar desordenats
Durabilitat Sí, replicat en diverses zones No es perden encara que caigui una zona
Retenció 7 dies per defecte, fins a 31 Es poden reproduir
Lliurament en ordre entre topics diferents No, en cap cas No assumeixis relació temporal entre topics

Les dues primeres files mereixen desenvolupament, perquè són contraintuïtives.

Per què hi haurà duplicats. L'escenari és senzill: el teu consumidor rep el missatge, el processa correctament, envia la confirmació, i la confirmació es perd per la xarxa. Pub/Sub no la rep, considera que el missatge continua pendent i el reenvia. El teu consumidor el processa per segona vegada. No hi ha cap fallada al teu codi i tot i així ha passat.

També passa si el consumidor triga més que l'ack deadline, si es reinicia a mig processament, o si l'escalat automàtic reassigna el missatge. És normal, és esperable i no és un error. Per això existeix l'apartat 10.

Exactament una vegada. Pub/Sub ofereix subscripcions amb aquesta semàntica (--enable-exactly-once-delivery), amb dos matisos importants: només s'aplica dins d'una regió, i garanteix que no hi haurà un segon lliurament confirmat, no que el teu processament sigui atòmic. Si el teu consumidor escriu a la base de dades i després cau abans de confirmar, la feina ja està feta i el missatge tornarà. La idempotència continua sent necessària; exactament una vegada només redueix la freqüència amb què la necessites.

Per què no hi ha ordre. Pub/Sub reparteix els missatges entre molts servidors per escalar. Si el missatge A es publica un mil·lisegon abans que el B, poden acabar en servidors diferents i arribar en qualsevol ordre. Per a AlpinaShop això importa poc a pedidos-nuevos —cada comanda és independent—, però importaria moltíssim si publiquéssim canvis d'estat de la mateixa comanda: "confirmado", "enviado", "entregado" fora d'ordre deixaria la comanda marcada com a confirmada després d'entregada. Aquest cas es resol amb clau d'ordenació (apartat 14).

  1. Crear el topic pedidos-nuevos i les seves subscripcions

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

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

--message-retention-duration al topic habilita la reproducció: permet fer seek a un instant del passat fins i tot per a subscripcions creades després. Sense això, només es pot tornar enrere dins del que retingui cada subscripció.

Ara el topic de missatges fallits, que crearem abans de les subscripcions perquè el necessitaran:

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

# I una subscripcio sobre ell, per poder inspeccionar el que hi caigui
gcloud pubsub subscriptions create sub-pedidos-fallidos \
  --topic=pedidos-nuevos-fallidos \
  --message-retention-duration=31d \
  --ack-deadline=600

Les tres subscripcions de consum:

# 1) MAGATZEM: pull, processament rapid, tolera duplicats per 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) FACTURACIO: exactament una vegada, perque duplicar una factura es greu
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

Comentaris sobre les diferències, que són deliberades:

  • sub-facturacion amb exactament una vegada. Emetre dues factures de la mateixa comanda és un problema comptable real. Encara que el consumidor ha de ser idempotent igualment, aquesta garantia redueix l'exposició. Té un cost: menys rendiment màxim per subscripció.
  • sub-analitica sense dead letter. El pipeline de Dataflow ja té el seu propi mecanisme de quarantena (04-02): els missatges mal formats van a pedidos_streaming_errores. Duplicar el mecanisme complicaria el diagnòstic.
  • ack deadline diferents. 30 s per al magatzem (operació ràpida), 120 s per a analítica (Dataflow processa per lots interns).

El dead letter necessita permisos explícits, i això s'oblida sempre:

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

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

# ...i CONFIRMAR a la subscripcio d origen per retirar el missatge
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"

Sense aquests dos permisos, la configuració de dead letter s'accepta sense protestar i no funciona: els missatges es reintenten eternament en comptes de desviar-se. És una fallada silenciosa clàssica.

Verificació:

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

  1. Publicar des de l'aplicació Flask

Ara la part del Dani, el desenvolupador de backend. La publicació des de l'aplicació de catàleg:

"""publicador.py -- Publicacio d esdeveniments de comanda a Pub/Sub."""
import json
import logging
from concurrent import futures

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

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

# 1) CONFIGURACIO DE LOTS.
#    El client acumula missatges i els envia junts: menys crides de xarxa,
#    menys cost. Publica quan es compleix el PRIMER dels tres limits.
config_lot = pubsub_v1.types.BatchSettings(
    max_messages=100,        # o 100 missatges...
    max_bytes=1024 * 1024,   # ...o 1 MB acumulat...
    max_latency=0.05,        # ...o 50 ms d espera. El que passi abans.
)

# 2) CONFIGURACIO DE REINTENTS amb retroces exponencial.
config_reintents = retry.Retry(
    initial=0.1,     # primer reintent als 100 ms
    maximum=60.0,    # mai esperar mes de 60 s entre intents
    multiplier=2.0,  # 0,1 -> 0,2 -> 0,4 -> 0,8 ...
    deadline=600.0,  # es rendeix als 10 minuts
)

# 3) El client es CAR de crear: un per proces, reutilitzat.
#    Crear-lo dins de la funcio de compra seria un error de rendiment greu.
publicador = pubsub_v1.PublisherClient(batch_settings=config_lot)
ruta_topic = publicador.topic_path(PROJECTE, TOPIC)

logger = logging.getLogger(__name__)


def publicar_comanda(comanda: dict) -> futures.Future:
    """Publica un esdeveniment de comanda confirmada. NO bloqueja.

    Retorna un Future. L aplicacio pot continuar responent al client
    sense esperar la confirmacio de Pub/Sub.
    """
    cos = json.dumps(comanda, ensure_ascii=False).encode("utf-8")

    futur = publicador.publish(
        ruta_topic,
        data=cos,
        # ATRIBUTS: metadades filtrables sense obrir el cos (apartat 12).
        # Tots els valors han de ser CADENES.
        tipo_evento="pedido_confirmado",
        origen=comanda.get("canal", "web"),
        pais=comanda["envio"]["pais"],
        version_esquema="1",
        momento_evento=comanda["momento"],   # ho fara servir Dataflow (04-02)
        pedido_id=comanda["pedido_id"],
    )

    def _en_acabar(fut):
        try:
            id_missatge = fut.result()
            logger.info("Publicada %s com a %s", comanda["pedido_id"], id_missatge)
        except Exception:
            # CRITIC: si la publicacio falla definitivament, cal
            # registrar-ho per poder recuperar-ho. Un registre d error aqui
            # ha de disparar una alerta (06-04).
            logger.exception("FALLADA EN PUBLICAR la comanda %s",
                             comanda["pedido_id"])

    futur.add_done_callback(_en_acabar)
    return futur

I el seu ús a la vista de Flask:

from flask import Flask, jsonify, request

app = Flask(__name__)


@app.post("/api/pedidos")
def confirmar_comanda():
    dades = request.get_json()

    # 1) L IMPRESCINDIBLE i sincron: persistir i cobrar.
    #    Si aixo falla, el client se n ha d assabentar.
    comanda = desar_a_cloud_sql(dades)
    cobrar_a_la_passarela(comanda)

    # 2) El DERIVAT: es publica el fet i es continua.
    #    Magatzem, facturacio, correu i analitica se n assabentaran sols.
    publicar_comanda({
        "pedido_id": comanda.id,
        "momento": comanda.creado_en.isoformat(),
        "cliente_id": comanda.cliente_id,
        "canal": comanda.canal,
        "total": str(comanda.total),
        "envio": {"pais": comanda.pais, "ciudad": comanda.ciudad},
        "lineas": [
            {"sku": l.sku, "cantidad": l.cantidad, "precio": str(l.precio)}
            for l in comanda.lineas
        ],
    })

    # 3) Resposta immediata: no esperem Pub/Sub ni els consumidors.
    return jsonify({"pedido_id": comanda.id, "estado": "confirmado"}), 201

Quatre qüestions de disseny que mereixen ser explícites:

Què va dins del missatge. Aquí publiquem un esdeveniment gros, amb les línies incloses. L'alternativa és un esdeveniment prim amb només el pedido_id, obligant cada consumidor a consultar la base de dades. El gros evita aquesta càrrega de lectures però acobla l'esquema del missatge; el prim és més flexible però multiplica les consultes a Cloud SQL. Per a AlpinaShop, amb ~1.200 comandes al mes, l'esdeveniment gros és clarament millor: menys crides a la base de dades operativa i els consumidors són autosuficients.

version_esquema com a atribut. El dia que canviï l'estructura del missatge, els consumidors antics podran reconèixer i rebutjar el que no entenen en comptes de trencar-se. Costa un atribut i estalvia un incident.

El momento_evento. És l'atribut que el pipeline de Dataflow de 04-02 fa servir com a timestamp_attribute. Sense ell, les finestres de temps de l'esdeveniment no funcionen.

Què passa si la publicació falla. És el risc real d'aquest disseny: la comanda està cobrada i l'avís no va sortir. Per a AlpinaShop, el registre amb alerta n'hi ha prou atès el volum. En sistemes de més exigència es fa servir el patró outbox: s'escriu l'esdeveniment en una taula de la mateixa base de dades dins de la mateixa transacció de la comanda, i un procés a part el publica. Així l'esdeveniment i la comanda són atòmics.

Prova ràpida des de la línia de comandes:

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. Consum pull amb el client asíncron

En el mode pull, el consumidor demana missatges. La biblioteca de Python fa servir streaming pull: manté una connexió oberta i rep missatges a mesura que arriben, amb molt poca latència.

"""consumidor_almacen.py -- Servei de magatzem d AlpinaShop."""
import json
import logging
import signal
import sys

from google.cloud import pubsub_v1

PROJECTE = "alpinashop-datos"
SUBSCRIPCIO = "sub-almacen"

logger = logging.getLogger(__name__)

# CONTROL DE FLUX: limita quants missatges es tenen alhora sense confirmar.
# Sense aixo, el client pot acceptar milers de missatges, no donar-los sortida
# a temps, deixar vencer l ack deadline i provocar reentregues massives.
control_flux = pubsub_v1.types.FlowControl(
    max_messages=50,
    max_bytes=10 * 1024 * 1024,
)


def processar(missatge):
    """Callback executat per la biblioteca en un fil del pool."""
    try:
        comanda = json.loads(missatge.data.decode("utf-8"))
    except json.JSONDecodeError:
        # Missatge corrupte: NO te arreglament amb reintents.
        # Es confirma per retirar-lo i es registra per investigar.
        logger.error("Missatge illegible, descartat: %s", missatge.message_id)
        missatge.ack()
        return

    tipus = missatge.attributes.get("tipo_evento")
    if tipus != "pedido_confirmado":
        # No es per a nosaltres; el confirmem sense fer res.
        missatge.ack()
        return

    try:
        # La idempotencia s implementa aqui dins (apartat 10)
        crear_ordre_de_preparacio(comanda)
        missatge.ack()
        logger.info("Preparacio creada per a %s", comanda["pedido_id"])

    except ErrorTemporal as exc:
        # Fallada transitoria (BD saturada, xarxa): NACK per reintentar ja.
        logger.warning("Fallada temporal a %s: %s", comanda["pedido_id"], exc)
        missatge.nack()

    except Exception:
        # Fallada desconeguda: NACK. Despres de 5 intents anira al dead letter.
        logger.exception("Fallada processant %s", comanda.get("pedido_id"))
        missatge.nack()


def main():
    subscriptor = pubsub_v1.SubscriberClient()
    ruta = subscriptor.subscription_path(PROJECTE, SUBSCRIPCIO)

    futur = subscriptor.subscribe(ruta, callback=processar,
                                  flow_control=control_flux)
    logger.info("Escoltant a %s...", ruta)

    # Apagada neta: en rebre SIGTERM (Kubernetes, Cloud Run),
    # es deixen d acceptar missatges nous i s acaben els que hi ha.
    def apagar(signum, frame):
        logger.info("Senyal %s rebut, tancant...", signum)
        futur.cancel()
        futur.result(timeout=30)
        sys.exit(0)

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

    with subscriptor:
        try:
            futur.result()
        except Exception:
            logger.exception("El subscriptor ha fallat")
            futur.cancel()


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

Tres punts que fan la diferència entre un consumidor de joguina i un de producció:

El control de flux. Sense FlowControl, la biblioteca accepta tants missatges com li enviïn. Si el teu processament és lent, molts venceran el seu termini, es reenviaran, el teu consumidor els processarà una altra vegada, i entraràs en una espiral en què cada vegada hi ha més feina duplicada. Limitar els missatges en vol és el que evita aquesta espiral.

ack() davant de nack(). ack() retira el missatge definitivament. nack() el retorna per a reentrega immediata. Si no fas cap de les dues, el missatge es reenvia en vèncer el termini, cosa que endarrereix el reintent però funciona igual.

Distingir errors recuperables d'irrecuperables. Un JSON corrupte no millorarà per reintentar-lo cinc vegades: es descarta amb registre. Una base de dades saturada sí: nack(). Confondre'ls fa que els missatges enverinats consumeixin recursos indefinidament o que es perdin dades recuperables.

També existeix el pull síncron, útil per a processos per lots:

# Llegir fins a 10 missatges sense confirmar-los (per inspeccionar)
gcloud pubsub subscriptions pull sub-almacen --limit=10 --format=json

# Llegir-los i confirmar-los
gcloud pubsub subscriptions pull sub-almacen --limit=10 --auto-ack

  1. Consum push a un endpoint HTTPS amb OIDC

En el mode push, Pub/Sub fa un POST HTTPS a una URL teva. Tu no mantens cap procés escoltant: el servei et crida.

Aspecte Pull Push
Qui inicia El consumidor Pub/Sub
Infraestructura Procés sempre en marxa Endpoint HTTPS; pot escalar a zero
Control del ritme Total, amb FlowControl Limitat (finestra lliscant automàtica)
Confirmació Explícita (ack()) Implícita: HTTP 2xx confirma, un altre codi és nack
Ideal per a Alt volum, processament continu Cloud Run, Cloud Functions, volum moderat

Push encaixa perfectament amb el que ve al curs: Cloud Run (07-02, on acabarà el catàleg segons la decisió DA-001) i Cloud Functions (06-03) escalen a zero, així que no té sentit tenir un procés esperant missatges.

La configuració segura fa servir autenticació OIDC: Pub/Sub adjunta un testimoni signat per Google que identifica un compte de servei, i el teu endpoint el verifica. Sense això, qualsevol que descobrís la teva URL podria injectar missatges falsos.

# 1) Compte de servei que representara Pub/Sub davant del teu endpoint
gcloud iam service-accounts create sa-pubsub-invocador \
  --display-name="Identitat de Pub/Sub per invocar serveis push"

SA_INV="[email protected]"

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

# 3) Permetre que Pub/Sub generi testimonis en nom d aquest compte
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) Subscripcio push amb 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

L'endpoint receptor:

import base64
import json

from flask import Flask, request

app = Flask(__name__)


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

    Cloud Run ja ha validat el testimoni OIDC abans que arribi aqui,
    perque el servei requereix autenticacio i el compte invocador
    te run.invoker. Si el servei fos public, caldria
    verificar la capcalera Authorization manualment.
    """
    sobre = request.get_json(silent=True)
    if not sobre or "message" not in sobre:
        # 400: no es reintentara. Es un error de format, no transitori.
        return "peticio mal formada", 400

    missatge = sobre["message"]

    # El cos arriba en base64 dins del sobre
    cos = base64.b64decode(missatge.get("data", "")).decode("utf-8")
    atributs = missatge.get("attributes", {})
    id_missatge = missatge["messageId"]
    intent = int(sobre.get("deliveryAttempt", 1))

    try:
        comanda = json.loads(cos)
    except json.JSONDecodeError:
        # 200 sense processar: confirmem perque NO es reintenti
        # un missatge que mai no sera valid.
        app.logger.error("Missatge %s illegible, descartat", id_missatge)
        return "", 204

    try:
        crear_ordre_de_preparacio(comanda)
    except ErrorTemporal:
        app.logger.warning("Fallada temporal a %s, intent %s",
                           id_missatge, intent)
        # 500: Pub/Sub reintentara amb retroces exponencial
        return "reintentar", 500

    # 204: confirmat
    return "", 204

La regla d'or del push: el codi de resposta HTTP és la confirmació. 2xx retira el missatge; qualsevol altra cosa (o un temps d'espera esgotat) el retorna a la cua. Retornar 200 en un except genèric "perquè no molesti" és la manera més ràpida de perdre dades en silenci.

  1. Confirmacions, terminis i reintents

L'ack deadline és el temps que Pub/Sub espera la confirmació abans de reenviar. Per defecte 10 segons; configurable entre 10 i 600.

sequenceDiagram
    participant PS as Pub/Sub
    participant C as Consumidor
    PS->>C: entrega missatge M (deadline 30 s)
    Note over C: processant... 25 s
    C->>PS: modifyAckDeadline(+60 s)
    Note over C: continua processant... 40 s
    C->>PS: ack(M)
    Note over PS: M retirat d aquesta subscripcio

Si el consumidor no confirma ni amplia el termini, el missatge torna a la cua. Triar bé el termini és important:

  • Massa curt: missatges que s'estan processant correctament es reenvien i es duplica la feina.
  • Massa llarg: si un consumidor mor, els seus missatges triguen molt a reassignar-se a un altre.

Regla pràctica: el percentil 99 del teu temps de processament, amb marge. Si el 99 % de les ordres de magatzem es creen en menys de 8 segons, 30 segons és raonable.

La bona notícia és que la biblioteca de Python amplia el termini automàticament mentre el teu callback continua executant-se (lease management), fins a un màxim configurable. Si necessites ampliar-lo a mà:

# Ampliar el termini des de dins del processament
missatge.modify_ack_deadline(120)

Els reintents segueixen la política que vas configurar:

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

Amb retrocés exponencial, els reintents s'espaien: 10 s, 20 s, 40 s, 80 s… fins a 600 s. Això és el correcte quan la fallada és perquè un sistema està saturat: reintentar cada segon l'enfonsaria més. Sense retrocés, un consumidor caigut genera una tempesta de reintents que impedeix que es recuperi.

  1. Temes de missatges fallits (dead letter)

Imagina't un missatge amb un sku que no existeix al catàleg. El consumidor del magatzem falla. Reintenta. Falla. Reintenta. Per sempre. Aquest missatge s'anomena enverinat, i sense dead letter té tres conseqüències: consumeix recursos indefinidament, embruta els registres, i —si hi ha ordenació activada— bloqueja tots els missatges posteriors de la seva clau.

El tema de missatges fallits resol això: després de N intents, el missatge es desvia a un altre topic i es retira de la subscripció original.

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

Quan un missatge arriba al dead letter, conserva el seu cos i els seus atributs i afegeix metadades sobre l'origen i el nombre d'intents. Un procés de revisió, típicament setmanal:

"""revisar_fallidos.py -- Inspeccio de la cua de missatges fallits."""
import json

from google.cloud import pubsub_v1

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

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

for rebut in resposta.received_messages:
    m = rebut.message
    print("---")
    print("ID original :", m.attributes.get(
        "CloudPubSubDeadLetterSourceMessageId", "?"))
    print("Subscripcio :", m.attributes.get(
        "CloudPubSubDeadLetterSourceSubscription", "?"))
    print("Intents     :", m.attributes.get(
        "CloudPubSubDeadLetterSourceDeliveryCount", "?"))
    print("Publicat    :", m.publish_time)
    print("Cos         :", m.data.decode("utf-8")[:300])

    # Decisio manual: corregir i republicar, o descartar amb registre.

Regles de dead letter que convé fixar com a política d'equip:

  1. Tota subscripció de producció en té un. Sense excepció. És la diferència entre "alguna cosa ha fallat i la tenim desada" i "alguna cosa ha fallat i no sabem què era".
  2. max-delivery-attempts entre 5 i 10. Menys, i una fallada transitòria envia missatges bons a la cua de fallits. Més, i es triga massa a detectar el problema.
  3. Una alerta sobre el nombre de missatges al dead letter. Un dead letter que ningú no mira és un forat negre amb passos extra.
  4. No dirigeixis un dead letter al topic original. És un bucle infinit, i és un error que es comet més del que sembla.

  1. Idempotència: el consumidor ha de tolerar duplicats

Ja sabem que hi haurà duplicats. La solució no és evitar-los —no es pot— sinó que processar dues vegades el mateix missatge tingui el mateix efecte que processar-lo una vegada. Això és la idempotència.

Tres estratègies, de menys a més robusta.

A. Operacions naturalment idempotents

La millor, quan és possible: dissenyar l'operació perquè repetir-la no canviï res.

# NO idempotent: sumar. Dues vegades resta el doble.
db.execute("UPDATE stock SET unidades = unidades - %s WHERE sku = %s",
           (quantitat, sku))

# Idempotent: fixar un valor absolut calculat des de la font de veritat.
db.execute("UPDATE stock SET unidades = %s WHERE sku = %s",
           (unitats_calculades, sku))

B. Inserció condicional per clau de negoci

Fer servir una restricció d'unicitat de la base de dades:

def crear_ordre_de_preparacio(comanda):
    """L index UNIQUE sobre pedido_id fa la feina per nosaltres."""
    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
            """,
            (comanda["pedido_id"],),
        )
        # Si ja existia, ON CONFLICT no fa res i no hi ha error.
        # El missatge duplicat es processa sense efecte: idempotent.

És l'opció preferible quan existeix una clau de negoci natural, com aquí pedido_id. Zero infraestructura afegida i la garantia la dona la base de dades.

C. Registre de missatges processats

Quan l'operació no és idempotent ni hi ha clau natural (enviar un correu, cridar una API externa), cal portar un registre. Firestore, que AlpinaShop ja fa servir per a la cistella (02-06), és ideal per la seva latència baixa i les seves transaccions:

"""idempotencia.py -- Registre de missatges ja processats."""
import datetime

from google.cloud import firestore

db = firestore.Client(project="alpinashop-datos")
COLECCIO = "eventos_procesados"
TTL_DIES = 14   # mes gran que la retencio maxima de Pub/Sub (7 dies)


def processar_una_sola_vegada(clau_idempotencia: str, consumidor: str, accio):
    """Executa `accio` nomes si aquesta clau no s ha processat abans.

    La clau inclou el consumidor: la mateixa comanda s ha de processar
    una vegada al magatzem I una vegada a facturacio. Son marques diferents.
    """
    doc_id = f"{consumidor}__{clau_idempotencia}"
    ref = db.collection(COLECCIO).document(doc_id)

    @firestore.transactional
    def _reservar(tx):
        instantania = ref.get(transaction=tx)
        if instantania.exists:
            return False          # ja processat: no fem res
        tx.set(ref, {
            "consumidor": consumidor,
            "clave": clau_idempotencia,
            "procesado_en": firestore.SERVER_TIMESTAMP,
            # Camp de TTL: Firestore esborra el document automaticament
            "caduca_en": datetime.datetime.utcnow()
                         + datetime.timedelta(days=TTL_DIES),
        })
        return True

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

    accio()
    return True

I el seu ús:

def processar_facturacio(missatge):
    comanda = json.loads(missatge.data.decode("utf-8"))

    # La clau d idempotencia es l identificador de NEGOCI,
    # no el message_id de Pub/Sub. Motiu: si la comanda es republica
    # per un reproces amb seek, el message_id sera diferent pero
    # la comanda es la mateixa i NO s ha de facturar dues vegades.
    executat = processar_una_sola_vegada(
        clau_idempotencia=comanda["pedido_id"],
        consumidor="facturacion",
        accio=lambda: emetre_factura(comanda),
    )

    if not executat:
        logger.info("Comanda %s ja facturada, s ignora el duplicat",
                    comanda["pedido_id"])

    missatge.ack()

El comentari sobre la clau és la part més important de tot l'apartat. Fer servir message_id protegeix contra les reentregues de Pub/Sub, però no contra un reprocés amb seek ni contra una republicació des de l'aplicació. Fer servir l'identificador de negoci (pedido_id) protegeix contra tots dos. És la diferència entre un sistema que aguanta un reprocés i un que factura dues vegades el dia que algú reprodueix una hora de missatges.

Un detall pràctic: el TTL del registre ha de ser més gran que la retenció màxima de la subscripció. Si Pub/Sub pot reentregar durant 7 dies i el teu registre caduca als 3, un missatge reentregat el dia 5 es processaria de nou.

  1. Retenció, seek i instantànies

Els missatges es retenen a la subscripció entre 10 minuts i 31 dies (7 per defecte). Amb seek es mou el punt de lectura.

# 1) Reprocessar les ultimes 3 hores: util despres de 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 TOT el que hi ha pendent: util en desembussar una cua inservible
gcloud pubsub subscriptions seek sub-almacen --time="$(date -u '+%Y-%m-%dT%H:%M:%SZ')"

Casos reals en què seek salva el dia:

  • Un error al consumidor d'analítica va calcular malament l'IVA durant sis hores. Es corregeix el codi, es desplega, es fa seek a fa sis hores i es reprocessa tot. Amb consumidors idempotents, això és segur. Sense idempotència, és un desastre.
  • Una cua va acumular 400.000 missatges obsolets per un consumidor caigut durant el cap de setmana. seek al present els descarta de cop.

Les instantànies (snapshots) capturen l'estat de confirmació d'una subscripció per poder-hi tornar:

# ABANS de desplegar una versio arriscada del consumidor
gcloud pubsub snapshots create snap-analitica-pre-v2 --subscription=sub-analitica

# ... es desplega, i si surt malament ...
gcloud pubsub subscriptions seek sub-analitica --snapshot=snap-analitica-pre-v2

# Neteja (les instantanies caduquen als 7 dies, pero millor ser explicit)
gcloud pubsub snapshots delete snap-analitica-pre-v2

És l'equivalent a un punt de restauració abans d'un desplegament de risc, i hauria de formar part del procediment de desplegament de qualsevol consumidor que escrigui en sistemes de negoci.

  1. Filtres de subscripció per atribut

Una subscripció pot filtrar pels atributs del missatge, de manera que només rep el que li interessa. El filtratge passa al costat de Pub/Sub: els missatges descartats ni es lliuren ni es facturen com a lliurament.

# Subscripcio que nomes rep comandes d Espanya
gcloud pubsub subscriptions create sub-almacen-es \
  --topic=pedidos-nuevos \
  --message-filter='attributes.pais = "ES"'

# Nomes comandes grans, per a revisio manual antifrau
gcloud pubsub subscriptions create sub-revision-fraude \
  --topic=pedidos-nuevos \
  --message-filter='attributes.tipo_evento = "pedido_confirmado" AND attributes.importe_alto = "true"'

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

# Prefix: qualsevol esdeveniment el tipus del qual comenci per "pedido_"
gcloud pubsub subscriptions create sub-todo-pedidos \
  --topic=eventos-tienda \
  --message-filter='hasPrefix(attributes.tipo_evento, "pedido_")'

Operadors disponibles: =, !=, AND, OR, NOT, hasPrefix(), i attributes:clau per comprovar existència.

Dues limitacions que cal conèixer:

  • Només es filtra per atributs, mai pel cos. Si vols filtrar per l'import, l'has de publicar com a atribut. D'aquí que el publicador de l'apartat 5 inclogui pais i origen com a atributs encara que també siguin dins del JSON.
  • El filtre és immutable. No es pot modificar un cop creada la subscripció; cal crear-ne una altra.

Els filtres permeten un patró molt net: un topic per domini, amb molts tipus d'esdeveniment, i cada consumidor filtra el seu. Per a AlpinaShop, un topic eventos-tienda amb pedido_confirmado, pedido_cancelado, carrito_abandonado i stock_bajo és més manejable que quatre topics, perquè conserva l'ordre relatiu dels esdeveniments del mateix domini i simplifica el publicador.

  1. Subscripcions de BigQuery i de Cloud Storage

Aquí arriba una de les millors funcions del servei, i la que més codi estalvia.

Una subscripció de BigQuery escriu els missatges directament en una taula. Sense Dataflow, sense codi, sense res per desplegar.

# 1) Taula desti amb l esquema del missatge
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) Permis perque Pub/Sub escrigui
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 subscripcio
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: s'espera JSON els camps del qual coincideixin amb les columnes. L'alternativa, --use-topic-schema, valida contra un esquema Avro o Protobuf registrat al topic.
  • --drop-unknown-fields: els camps del JSON que no existeixin com a columna s'ignoren en comptes de provocar una fallada. Imprescindible si l'esquema del missatge pot evolucionar.
  • El dead letter aquí és fonamental: sense ell, un missatge que no encaixi amb l'esquema es reintenta indefinidament.

La comparació honesta:

Criteri Subscripció de BigQuery Dataflow (04-02)
Codi Cap Un pipeline que cal mantenir
Cost Només el de Pub/Sub i l'escriptura Treballadors 24×7: desenes de €/mes
Transformació Cap: el JSON va tal qual Qualsevol
Finestres i temps de l'esdeveniment No Sí
Enriquiment amb altres fonts No Sí
Agregació en vol No Sí
Quan fer-la servir El missatge ja té la forma de la taula Cal transformar, agregar o fer finestres

La decisió d'AlpinaShop: per bolcar els esdeveniments crus en una taula d'aterratge, subscripció de BigQuery, perquè és gratis en esforç i no hi ha res per mantenir. El pipeline de Dataflow es reserva per al que sí que necessita transformació: l'agregat per hora i canal amb finestres de temps de l'esdeveniment i tolerància a dades tardanes. És el patró habitual i correcte: dades crues per la via barata, agregats per la via potent.

La subscripció de Cloud Storage fa l'equivalent amb fitxers, agrupant missatges per temps o mida:

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

És una manera excel·lent de tenir un arxiu immutable de tots els esdeveniments al llac, que després processa Spark (04-03) o carrega BigQuery. I és la xarxa de seguretat definitiva: si un dia cal reconstruir tot el magatzem analític, els esdeveniments crus són allà.

  1. Ordenació per clau i les seves implicacions

Per defecte no hi ha ordre. Per garantir-lo, es publica amb clau d'ordenació:

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

# Tots els esdeveniments de la MATEIXA comanda comparteixen clau: arriben en ordre
publicador.publish(
    ruta_topic,
    data=cos,
    ordering_key=comanda["pedido_id"],    # <-- la clau
    tipo_evento="pedido_enviado",
)
gcloud pubsub subscriptions create sub-estados-pedido \
  --topic=eventos-tienda \
  --enable-message-ordering

La garantia és: els missatges amb la mateixa clau d'ordenació, publicats a la mateixa regió, es lliuren en ordre de publicació. Missatges amb claus diferents no tenen relació entre si, cosa que permet continuar paral·lelitzant.

Els tres costos d'activar l'ordenació, que cal sospesar:

  1. Menys rendiment. Els missatges d'una clau es processen seqüencialment. Amb pedido_id com a clau no hi ha problema: hi ha milers de comandes diferents i el paral·lelisme es manté. Amb pais com a clau, totes les comandes d'Espanya anirien en sèrie.
  2. Bloqueig per missatge enverinat. Si un missatge de la clau PED-2026-0042 falla repetidament, tots els missatges posteriors d'aquella comanda queden bloquejats fins que es resolgui o es desviï al dead letter. Ordenació sense dead letter és una bomba de rellotgeria.
  3. Publicació més lenta. El client ha de serialitzar els enviaments de cada clau, cosa que redueix l'efecte de l'enviament per lots.

Regla d'or: activa l'ordenació només quan l'ordre sigui semànticament necessari, i tria la clau amb la màxima cardinalitat possible que preservi aquest ordre. Per a AlpinaShop: pedido_id sí (els estats d'una comanda han d'anar en ordre); pais no.

I l'alternativa que sol ser millor: fer que l'ordre no importi. Si cada esdeveniment inclou un número de versió o una marca de temps, el consumidor pot descartar els que siguin anteriors al que ja va processar, i l'ordre d'arribada deixa de ser un problema.

# Consumidor tolerant al desordre, sense necessitat d ordering_key
def aplicar_estat(pedido_id, nou_estat, versio_esdeveniment):
    db.execute(
        """
        UPDATE pedidos SET estado = %s, version = %s
        WHERE pedido_id = %s AND version < %s
        """,
        (nou_estat, versio_esdeveniment, pedido_id, versio_esdeveniment),
    )
    # Si arriba un esdeveniment antic, la condicio version < %s no es compleix
    # i l UPDATE no afecta cap fila. Desordre absorbit.

  1. Pub/Sub Lite i altres alternatives

Pub/Sub Lite va ser una variant de menys cost, amb capacitat aprovisionada per l'usuari en comptes de servida sota demanda, pensada per a volums molt alts i predictibles. Requeria dimensionar particions i emmagatzematge a mà, a canvi d'un preu per GB considerablement inferior.

Google en va anunciar la retirada i el servei va deixar d'estar disponible el març de 2026. Si trobes referències a Pub/Sub Lite en documentació o cursos antics, són obsoletes; verifica sempre l'estat a la documentació oficial. Les alternatives actuals són Pub/Sub estàndard o, si necessites l'API de Kafka, Google Cloud Managed Service for Apache Kafka.

La comparació que sí que continua sent útil:

Servei Model Quan té sentit
Pub/Sub Global, sense aprovisionar, pagament per ús El valor per defecte: la immensa majoria dels casos
Managed Kafka Kafka gestionat, amb particions i grups de consumidors Ja tens codi Kafka, o necessites la seva semàntica de log
Cloud Tasks Cua de tasques amb programació i control de ritme Encuar feina amb un destinatari únic i conegut
Eventarc Encaminament d'esdeveniments de la plataforma Reaccionar a esdeveniments de serveis de Google (fa servir Pub/Sub per sota)
Memorystore (Redis Pub/Sub) Missatgeria en memòria, sense persistència Notificacions efímeres on perdre missatges és acceptable

La confusió més habitual és Pub/Sub davant de Cloud Tasks. Pub/Sub és per a fets que interessen a N consumidors desconeguts. Cloud Tasks és per a encàrrecs dirigits a un destinatari concret, amb control fi del ritme de lliurament i programació diferida. "S'ha confirmat una comanda" és Pub/Sub. "Envia aquest correu d'aquí a 30 minuts, amb un màxim de 10 per segon" és Cloud Tasks.

  1. Monitoratge i cost

Les mètriques que cal vigilar, amb les seves alertes:

Mètrica Què indica Alerta raonable
subscription/oldest_unacked_message_age Antiguitat del missatge pendent més antic La mètrica reina. > 600 s: el consumidor no dona l'abast o és caigut
subscription/num_undelivered_messages Missatges acumulats sense lliurar Creixement sostingut: els consumidors van endarrerits
subscription/dead_letter_message_count Missatges desviats a fallits > 0 en una hora: cal mirar-ho
topic/send_request_count Publicacions Caiguda brusca: l'aplicació ha deixat de publicar
subscription/push_request_count per codi Salut de l'endpoint push Molts 5xx: l'endpoint falla
subscription/ack_message_count Ritme de confirmació Comparar amb publicacions

L'alerta imprescindible:

gcloud alpha monitoring policies create \
  --notification-channels="$CANAL_INFRA" \
  --display-name="Pub/Sub: missatges sense confirmar a 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 és la millor mètrica perquè detecta alhora les tres fallades possibles: consumidor caigut (creix indefinidament), consumidor lent (creix a poc a poc) i missatge enverinat bloquejant una clau ordenada (es queda clavat en un valor alt).

El cost es factura principalment per volum de dades:

Concepte Ordre de magnitud (verificar a la documentació oficial)
Publicació + lliurament ~40 $ per TiB, amb una franja mensual gratuïta
Emmagatzematge de missatges retinguts ~0,27 $ per GiB i mes
Sortida entre regions Tarifes de xarxa

Un càlcul realista per a AlpinaShop: 1.200 comandes al mes, missatges de ~2 KB, amb 4 subscripcions. Això són 1.200 publicacions i 4.800 lliuraments: uns 12 MB mensuals. Cèntims, o directament dins de la franja gratuïta. Els esdeveniments de navegació del catàleg, molt més nombrosos, continuarien sent barats.

Un detall de facturació que sorprèn: cada subscripció compta com un lliurament. Deu subscripcions sobre el mateix topic multipliquen per deu el volum facturat de lliurament. No és motiu per no fer servir subscripcions, però sí per no deixar subscripcions òrfenes: una subscripció oblidada sense consumidor acumula missatges, se'n factura l'emmagatzematge i no serveix de res.

# Buscar subscripcions sense consum: candidates a esborrar
gcloud pubsub subscriptions list --format="table(name, topic)" | while read -r s _; do
  echo "$s"
done

  1. Les notificacions del bucket alpinashop-catalogo

A 02-02 va quedar pendent una promesa: que Cloud Storage pot avisar quan apareix un objecte nou. Ara tenim les peces per tancar-la de debò.

El cas d'AlpinaShop: quan algú puja la imatge original d'un producte a productos/<sku>/original/, cal generar automàticament les versions web/ i thumb/. Avui això ho fa la Marta a mà amb un script.

# 1) Topic per als esdeveniments del bucket
gcloud pubsub topics create imagenes-subidas \
  --labels=entorno=produccion,equipo=infra,aplicacion=catalogo

# 2) Permis perque l agent de Cloud Storage publiqui
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 notificacio, acotada al prefix i a l esdeveniment que interessen
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) Subscripcio per al processador d imatges
gcloud pubsub subscriptions create sub-procesar-imagenes \
  --topic=imagenes-subidas \
  --ack-deadline=300 \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5

Tipus d'esdeveniment disponibles:

Esdeveniment Quan es dispara
OBJECT_FINALIZE Es crea un objecte nou o se sobreescriu
OBJECT_DELETE S'esborra (o se sobreescriu la versió anterior)
OBJECT_ARCHIVE Una versió passa a arxivada (amb versionat actiu)
OBJECT_METADATA_UPDATE Canvien les metadades

El consumidor:

def processar_imatge(missatge):
    """Genera les versions web i thumb d una imatge acabada de pujar."""
    # Les dades de l objecte arriben als ATRIBUTS, no cal el cos
    bucket = missatge.attributes["bucketId"]
    nom = missatge.attributes["objectId"]
    generacio = missatge.attributes["objectGeneration"]

    # 1) FILTRE DE BUCLE INFINIT: si no comprovem aixo, en escriure
    #    web/ i thumb/ es dispararien noves notificacions que generarien
    #    mes imatges, indefinidament. Es l error classic.
    if "/original/" not in nom:
        missatge.ack()
        return

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

    # 2) IDEMPOTENCIA: la generacio identifica la versio exacta de l objecte.
    #    Si el missatge es reentrega, la clau es la mateixa i no es repeteix.
    clau = f"{bucket}/{nom}#{generacio}"

    processar_una_sola_vegada(
        clau_idempotencia=clau,
        consumidor="procesador-imagenes",
        accio=lambda: generar_versions(bucket, nom),
    )
    missatge.ack()

El filtre del bucle infinit mereix subratllar-se: és la fallada més cara d'aquest patró. Escriure al mateix bucket que dispara la notificació genera una recursió que només es detecta quan arriba la factura. Les tres defenses són: filtrar per prefix a la notificació (--object-prefix=productos/), comprovar la ruta al consumidor, i —millor encara— escriure la sortida en un bucket diferent del d'entrada.

Aquest mateix patró s'aplica al fitxer mensual del transportista i a les exportacions nocturnes: així que el fitxer aterra a exportaciones/<yyyy>/<mm>/<dd>/, una notificació dispara la càrrega a BigQuery. Sense cron, sense comprovar carpetes cada cinc minuts, sense finestres d'espera. La lògica del processament la posarem a Cloud Functions a 06-03, i qui ho orquestra tot és la lliçó següent.

Errors habituals i consells

Suposar que no hi haurà duplicats. N'hi haurà, amb o sense exactament una vegada. Tot consumidor de producció ha de ser idempotent. No és opcional.

Fer servir message_id com a clau d'idempotència. Protegeix contra reentregues, però no contra seek ni contra republicacions. Fes servir l'identificador de negoci.

Configurar dead letter i oblidar els permisos. La comanda s'accepta i la funció no opera: els missatges es reintenten per sempre. Comprova sempre els dos bindings.

No posar control de flux al consumidor. Amb processament lent, la biblioteca accepta més missatges dels que pot confirmar, vencen els terminis, es reentreguen, i el consumidor s'enfonsa processant duplicats del seu propi embús.

Crear el client de Pub/Sub dins del gestor de la petició. És car (obre connexions gRPC). Un per procés, reutilitzat.

Retornar 200 en un except genèric en push. Confirma el missatge i el perd en silenci. Retorna 5xx en les fallades transitòries.

Activar ordenació sense necessitar-la. Redueix el rendiment i, amb un missatge enverinat, bloqueja tota la clau.

Dirigir el dead letter al topic original. Bucle infinit.

Subscripcions òrfenes. Sense consumidor, acumulen missatges fins a la retenció màxima, se'n factura l'emmagatzematge i no aporten res. Revisa-les periòdicament.

Notificacions de bucket que es retroalimenten. Escriure al mateix bucket que dispara l'esdeveniment. Fes servir buckets o prefixos separats i filtra al consumidor.

Consell: publica atributs generosament. Costen poc i permeten filtrar del costat del servidor, cosa que estalvia lliuraments i complexitat al consumidor.

Consell: un esquema al topic. Pub/Sub admet registrar un esquema Avro o Protobuf i validar els missatges en publicar. Converteix els errors de format en fallades immediates del publicador en comptes de sorpreses al consumidor.

Consell: fes una instantània abans de cada desplegament arriscat. Costa una comanda i dona marxa enrere.

Exercicis

Exercici 1: topic de ressenyes amb filtres i dead letter

Crea el topic opiniones-nuevas amb 7 dies de retenció, el seu topic de fallits opiniones-nuevas-fallidos amb una subscripció d'inspecció, i tres subscripcions sobre el topic principal:

  • sub-moderacion: rep només les ressenyes amb attributes.puntuacion_baja = "true", amb termini de 60 s, dead letter després de 5 intents i retrocés de 10 s a 300 s.
  • sub-analitica-opiniones: les rep totes, amb termini de 120 s.
  • sub-bq-opiniones: escriu directament a BigQuery a la taula alpinashop-datos:alpinashop_analitica.opiniones, sense codi.

Inclou tots els permisos necessaris i publica dos missatges de prova, un que arribi a moderació i un altre que no.

Exercici 2: consumidor idempotent de facturació

Escriu el consumidor pull de sub-facturacion que emet la factura de cada comanda. Requisits: control de flux de 20 missatges en vol; distingir errors transitoris (reintentar) de permanents (descartar amb registre); idempotència basada en l'identificador de negoci amb registre a Firestore i TTL adequat; apagada neta davant de SIGTERM; i un comptador de factures emeses i duplicats detectats. Explica per què la clau d'idempotència ha de ser el pedido_id i no el message_id, amb un escenari concret en què l'elecció incorrecta causaria un problema real.

Exercici 3: diagnòstic d'una cua embussada

A les 09:15 salta l'alerta: oldest_unacked_message_age a sub-almacen porta 47 minuts i puja. num_undelivered_messages ha passat de 3 a 12.400 des de les 08:20. El servei de magatzem està aixecat i respon a la seva comprovació d'estat. Als registres del consumidor apareix, unes cent vegades per minut, el mateix missatge d'error esmentant la comanda PED-2026-1188. La cua de pedidos-nuevos-fallidos és buida. La subscripció té l'ordenació activada per pedido_id.

Diagnostica què està passant, explica per què el dead letter és buit malgrat els reintents, i dona un pla d'actuació en tres fases: contenció immediata, correcció i prevenció.

Solucions

Solució 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) Permis de publicacio al topic de fallits
gcloud pubsub topics add-iam-policy-binding opiniones-nuevas-fallidos \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.publisher"

# 3) Subscripcio de moderacio, amb filtre
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) Subscripcio d analitica, sense filtre
gcloud pubsub subscriptions create sub-analitica-opiniones \
  --topic=opiniones-nuevas --ack-deadline=120 \
  --labels=consumidor=analitica

# 5) Subscripcio 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
# Missatge que SI arriba a moderacio (puntuacio 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 moltissim menys del que s anuncia","pais":"FR"}' \
  --attribute=puntuacion_baja=true,sku=FRON-300L,pais=FR

# Missatge que NO arriba a moderacio (puntuacio 5)
gcloud pubsub topics publish opiniones-nuevas \
  --message='{"opinion_id":"OPI-1002","sku":"MOCH-40L-AZ","fecha":"2026-03-14","puntuacion":5,"texto":"Comodissima en travessies llargues","pais":"ES"}' \
  --attribute=puntuacion_baja=false,sku=MOCH-40L-AZ,pais=ES

# Comprovacio: moderacio rep 1, analitica rep 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 punt que s'avalua aquí és el filtre: sub-moderacion rep una de les dues publicacions, perquè el filtratge passa al costat de Pub/Sub i la ressenya de 5 estrelles ni tan sols es lliura ni es factura.

Solució 2

"""consumidor_facturacion.py -- Emissio idempotent de factures."""
import datetime
import json
import logging
import signal
import sys

from google.cloud import firestore, pubsub_v1

PROJECTE = "alpinashop-datos"
SUBSCRIPCIO = "sub-facturacion"
COLECCIO = "eventos_procesados"
TTL_DIES = 14           # > 7 dies de retencio maxima de la subscripcio
CONSUMIDOR = "facturacion"

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

factures_emeses = 0
duplicats_detectats = 0


class ErrorTemporal(Exception):
    """Fallada transitoria: mereix reintent."""


class ErrorPermanent(Exception):
    """Fallada que no millorara reintentant."""


def processar_una_sola_vegada(clau, accio):
    doc = db.collection(COLECCIO).document(f"{CONSUMIDOR}__{clau}")

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

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


def processar(missatge):
    global factures_emeses, duplicats_detectats
    try:
        comanda = json.loads(missatge.data.decode("utf-8"))
    except json.JSONDecodeError:
        logger.error("Missatge %s illegible, descartat", missatge.message_id)
        missatge.ack()                     # permanent: no reintentar
        return

    if not comanda.get("pedido_id"):
        logger.error("Missatge %s sense pedido_id, descartat", missatge.message_id)
        missatge.ack()
        return

    try:
        emesa = processar_una_sola_vegada(
            clau=comanda["pedido_id"],
            accio=lambda: emetre_factura(comanda),
        )
        if emesa:
            factures_emeses += 1
            logger.info("Factura emesa per a %s", comanda["pedido_id"])
        else:
            duplicats_detectats += 1
            logger.info("Duplicat ignorat: %s", comanda["pedido_id"])
        missatge.ack()

    except ErrorTemporal as exc:
        logger.warning("Fallada temporal a %s: %s", comanda["pedido_id"], exc)
        missatge.nack()                    # reintent amb retroces

    except ErrorPermanent as exc:
        logger.error("Fallada permanent a %s: %s", comanda["pedido_id"], exc)
        missatge.ack()                     # no te arreglament: es retira

    except Exception:
        logger.exception("Fallada desconeguda a %s", comanda["pedido_id"])
        missatge.nack()                    # despres de 5 intents, al dead letter


def main():
    subscriptor = pubsub_v1.SubscriberClient()
    ruta = subscriptor.subscription_path(PROJECTE, SUBSCRIPCIO)
    control = pubsub_v1.types.FlowControl(max_messages=20)

    futur = subscriptor.subscribe(ruta, callback=processar, flow_control=control)
    logger.info("Facturacio escoltant a %s", ruta)

    def apagar(signum, frame):
        logger.info("Tancant. Emeses=%s Duplicats=%s",
                    factures_emeses, duplicats_detectats)
        futur.cancel()
        futur.result(timeout=30)
        sys.exit(0)

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

    with subscriptor:
        try:
            futur.result()
        except Exception:
            logger.exception("Subscriptor caigut")
            futur.cancel()


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

Per què pedido_id i no message_id, amb escenari concret:

message_id és únic per publicació. Si el mateix fet de negoci es publica dues vegades, cada publicació tindrà un message_id diferent, i un registre basat en ell no detectaria res.

L'escenari real: un dimarts, una fallada al consumidor de facturació va provocar que 300 comandes del matí no es facturessin. Es corregeix el codi, es desplega i s'executa:

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

Aquest seek reentrega els missatges del matí. Ara bé, entre les 08:00 i l'incident hi va haver 500 comandes, de les quals 200 sí que es van facturar correctament abans que comencés la fallada.

  • Amb clau message_id: la reentrega conserva el message_id original, així que en aquest cas concret les 200 bones sí que es detectarien. Però si en comptes de seek algú recupera i republica els esdeveniments des de l'arxiu de Cloud Storage (apartat 13), els message_id seran nous i s'emetrien 200 factures duplicades. Aquesta és la fallada.
  • Amb clau pedido_id: tant li fa com arribi el missatge —reentrega, seek, republicació des de l'arxiu, reprocés manual—. Si aquella comanda ja es va facturar, no es torna a facturar. La protecció és sobre el fet de negoci, que és el que realment importa.

Dues-centes factures duplicades enviades a clients és un incident comptable, d'atenció al client i de reputació. La diferència entre les dues opcions és una línia de codi.

Solució 3

Diagnòstic: missatge enverinat bloquejant una clau d'ordenació, amb el dead letter mal configurat.

Les pistes encaixen una a una:

  1. El servei és viu i respon. No és una caiguda: és un embús lògic.
  2. El mateix error cent vegades per minut sobre PED-2026-1188. Aquest missatge falla sempre. És un missatge enverinat: alguna cosa que hi ha (un SKU inexistent, un camp nul, un import impossible) fa fallar el consumidor de manera determinista.
  3. L'ordenació està activada per pedido_id. Aquí hi ha la part crítica que explica la magnitud: amb ordenació, els missatges de la mateixa clau es lliuren en ordre i un de bloquejat impedeix avançar. Però a més, a la pràctica, la combinació d'un consumidor que fa nack() en bucle sobre una clau i el control de flux saturat per reentregues fa que el rendiment global s'ensorri: el consumidor gasta la seva capacitat reprocessant el mateix missatge una vegada i una altra en comptes d'atendre la cua.
  4. 12.400 missatges acumulats en 55 minuts quan AlpinaShop fa ~1.200 comandes al mes: aquest volum no són comandes noves, són reentregues del mateix missatge més la resta de la cua sense atendre.

Per què el dead letter és buit. És la part que més ensenya. Hi ha dues causes possibles i totes dues són freqüents:

  • Falten els permisos. La política de dead letter es configura sense error encara que l'agent de servei de Pub/Sub no tingui roles/pubsub.publisher sobre el topic de fallits ni roles/pubsub.subscriber sobre la subscripció d'origen. Sense aquests dos permisos, el desviament no passa mai i el missatge es reintenta indefinidament. És exactament la fallada silenciosa que advertíem a l'apartat 4.
  • La política no està aplicada realment. Es va crear la subscripció sense --dead-letter-topic i es va donar per fet.

Comprovació:

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

Pla d'actuació en tres fases:

Fase 1 — Contenció (minuts). Arreglar els permisos del dead letter, que és el que desembussa sense perdre res:

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

Tan bon punt el desviament funcioni, PED-2026-1188 sortirà de la cua després de cinc intents i la resta de missatges començarà a fluir. No fer seek al present: això descartaria 12.400 missatges que en la seva majoria són comandes reals pendents de preparar.

Si el desviament trigués i hi hagués urgència operativa, el pedaç alternatiu és desplegar al consumidor una regla temporal que confirmi explícitament el pedido_id problemàtic, deixant-lo registrat per tractar-lo a mà.

Fase 2 — Correcció (hores). Inspeccionar el missatge a la cua de fallits, entendre per què el consumidor hi falla, i corregir el codi perquè aquell tipus de dada es tracti com a error permanent (ack() amb registre) en comptes de reintentar-se. Després, tractar manualment la comanda PED-2026-1188, que continua sent una comanda real d'un client real i cal preparar-la. Verificar que la cua torna a un oldest_unacked_message_age de segons.

Fase 3 — Prevenció (dies).

  • Verificar el dead letter de totes les subscripcions, inclosos els permisos, no només la política. Automatitzar-ho com a comprovació al desplegament.
  • Alerta sobre dead_letter_message_count > 0, per assabentar-se del primer missatge enverinat en comptes del número 12.400.
  • Revisar si l'ordenació és realment necessària a sub-almacen. Les comandes són independents entre si; l'ordre entre comandes diferents no aporta res i sí que afegeix risc de bloqueig. Si només cal ordenar els canvis d'estat d'una mateixa comanda, aquest cas es pot resoldre amb el número de versió i un UPDATE ... WHERE version < %s, eliminant l'ordenació completament.
  • Classificar errors al consumidor: transitori → nack(), permanent → ack() amb registre. L'absència d'aquesta distinció és la causa arrel que una sola dada dolenta tombi una cua.
  • Afegir l'alerta d'oldest_unacked_message_age a 10 minuts, no a 47.

Conclusió

AlpinaShop ja no és un monòlit que truca a tothom per telèfon. En aquesta lliçó has vist el problema de l'acoblament amb números concrets —una compra que espera 5,4 segons quatre sistemes que al client no li importen— i com es dissol publicant un fet en comptes de donar ordres.

Has interioritzat el model: el topic és el tipus de fet, la subscripció és on viuen realment els missatges i cadascuna rep la seva còpia independent, el missatge porta cos i atributs filtrables, i la confirmació és l'única cosa que retira un missatge del sistema. I sobretot has vist les garanties reals, que són les que importen: lliurament almenys una vegada, amb duplicats que passaran encara que el teu codi sigui perfecte; sense ordre llevat que hi hagi clau d'ordenació; i exactament una vegada com una reducció del problema, no com la seva desaparició.

Has creat pedidos-nuevos amb retenció de 7 dies, el topic de fallits pedidos-nuevos-fallidos amb la seva subscripció d'inspecció, i les subscripcions sub-almacen, sub-facturacion —amb exactament una vegada, perquè duplicar una factura és greu— i sub-analitica, cadascuna amb el seu termini, la seva política de reintents i els seus permisos de dead letter, inclosos els dos bindings que tothom oblida i sense els quals el desviament no funciona en silenci.

Has publicat des de Flask amb enviament per lots, retrocés exponencial i un client reutilitzat, decidint conscientment publicar un esdeveniment gros amb les línies a dins i marcant version_esquema i momento_evento com a atributs. Has consumit en pull amb control de flux, apagada neta i distinció entre errors transitoris i permanents; i en push cap a un endpoint HTTPS amb autenticació OIDC, sabent que allà el codi de resposta HTTP és la confirmació. Coneixes l'ack deadline, modifyAckDeadline i per què el termini s'ajusta al percentil 99 del processament.

Has entès per què tot consumidor seriós necessita un dead letter —el missatge enverinat que es reintenta per sempre— i per què la clau d'idempotència ha de ser l'identificador de negoci i no el message_id: un seek o una republicació des de l'arxiu emetrien factures duplicades amb l'elecció equivocada. Has vist seek i instantànies com a xarxa de seguretat per a reprocessos i desplegaments arriscats, els filtres per atribut que descarten del costat del servidor, i les subscripcions de BigQuery i Cloud Storage, que ingereixen sense ni una sola línia de codi i que per a AlpinaShop s'emporten les dades crues mentre Dataflow es queda amb el que de debò necessita transformació i finestres.

Saps quan activar l'ordenació i a quin preu, que Pub/Sub Lite ja no existeix i què fer servir en el seu lloc, quina mètrica vigilar per damunt de totes —oldest_unacked_message_age— i que el cost, amb el volum d'AlpinaShop, és de cèntims. I has tancat per fi la promesa de 02-02: el bucket alpinashop-catalogo ja avisa quan apareix una imatge nova, amb el filtre de prefix, la comprovació de ruta i la idempotència per generació d'objecte que impedeixen el bucle infinit que arruïna els incauts.

Fixa't en el mapa que tenim ara. BigQuery desa i respon. Dataflow transforma per lots i en temps real. Dataproc calcula el que és algorítmic. Pub/Sub connecta els sistemes sense que es coneguin. És molta maquinària. I tanmateix, hi ha dos buits evidents.

El primer: les dades que no neixen a Google Cloud. L'ERP del magatzem d'AlpinaShop és un MySQL que corre en un servidor de les oficines de Sabadell. El transportista envia un CSV mensual per correu. El proveïdor de motxilles ofereix una API REST amb el seu catàleg. Res d'això no publica a Pub/Sub ni escriu Parquet en un bucket, i la Lucía, que sap SQL però no ha escrit Beam a la vida, no pot dependre del Dani per a cada fitxer nou.

A 04-05, Cloud Data Fusion, veurem l'eina pensada exactament per a això: integració de dades visual, sense codi, amb connectors per a gairebé qualsevol origen, un netejador interactiu anomenat Wrangler on es corregeixen dates i nuls veient les dades, llinatge a nivell de camp que respon a "d'on surt aquesta columna?", i replicació amb captura de canvis des d'aquell MySQL de Sabadell. I també veurem el seu costat incòmode —costa per hora d'instància i això en condiciona l'ús—, amb una taula honesta de quan Data Fusion, quan Dataflow, quan Datastream i quan n'hi ha prou amb un simple bq load.

Curs de Google Cloud Platform (GCP)

Mòdul 1: Introducció a Google Cloud Platform

Mòdul 2: Serveis principals de GCP

Mòdul 3: Xarxes i seguretat

Mòdul 4: Dades i anàlisi

Mòdul 5: Aprenentatge automàtic i IA

Mòdul 6: DevOps i monitoratge

Mòdul 7: Temes avançats de GCP

Mòdul 8: Projecte final

© Copyright 2026. Tots els drets reservats