Let us go back to the AlpinaShop application for a moment. A customer presses "Confirm order". The Flask function handling that request has to do several things: save the order in alpinashop-pedidos, take payment at the payment gateway, tell the warehouse to prepare it, tell billing to issue the invoice, send the confirmation email, decrement stock, and — since module 4 — inform analytics.

Today, that code is a list of calls, one after another. And that list has three problems that do not show up until the worst day of the year.

The customer waits for all of them. The purchase's response time is the sum of the seven steps. If the mail service takes four seconds, the customer watches a spinner for four extra seconds over something they do not care about: they have already bought.

One failure brings it all down. If the warehouse system is restarting, the call fails. And now what? Do you undo the order, which has already been charged? Do you carry on and lose the notification? Do you retry and run the risk of duplicating the invoice, which did go out? There is no good answer.

Every new addition touches the purchase code. Adding analytics means modifying, testing and deploying the shop's most critical function. And then the loyalty programme will arrive, and the fraud checks, and the supplier notification. The purchase function grows without pause and every change puts the till at risk.

Cloud Pub/Sub solves all three at once with a very simple idea: instead of calling seven systems, the shop publishes a fact — "order PED-2026-0042 has been confirmed" — and gets on with its business. Whoever is interested in that fact subscribes. The shop neither knows nor cares how many there are.

In this lesson you will create the pedidos-nuevos topic and its subscriptions, publish from Flask, consume in all four possible ways, understand the service's real guarantees — which are not the ones people assume — and learn why idempotency is not an architectural ornament but an unavoidable requirement.

Contents

  1. The coupling problem, before and after
  2. The publish-subscribe model
  3. Pub/Sub's real guarantees
  4. Creating the pedidos-nuevos topic and its subscriptions
  5. Publishing from the Flask application
  6. Pull consumption with the asynchronous client
  7. Push consumption to an HTTPS endpoint with OIDC
  8. Acknowledgements, deadlines and retries
  9. Dead letter topics
  10. Idempotency: the consumer must tolerate duplicates
  11. Retention, seek and snapshots
  12. Subscription filters by attribute
  13. BigQuery and Cloud Storage subscriptions
  14. Ordering by key and its implications
  15. Pub/Sub Lite and other alternatives
  16. Monitoring and cost
  17. The alpinashop-catalogo bucket notifications

  1. The coupling problem, before and after

flowchart TD
    subgraph Antes["BEFORE: chained synchronous calls"]
        W1["Flask: confirm_order()"]
        A1["Warehouse"]
        F1["Billing"]
        M1["Mail"]
        AN1["Analytics"]
        W1 -->|"waits 200 ms"| A1
        W1 -->|"waits 350 ms"| F1
        W1 -->|"waits 4 s"| M1
        W1 -->|"waits 800 ms"| AN1
    end
flowchart TD
    subgraph Despues["AFTER: one published fact, N interested parties"]
        W2["Flask: confirm_order()"]
        T["Topic pedidos-nuevos"]
        S1["sub-almacen"]
        S2["sub-facturacion"]
        S3["sub-analitica"]
        S4["sub-email"]
        A2["Warehouse service"]
        F2["Billing service"]
        AN2["Dataflow -> BigQuery"]
        M2["Mail service"]

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

What changes concretely:

Aspect Chained calls With Pub/Sub
Purchase latency The sum of every step (~5.4 s) Just the publication (~15 ms)
A consumer that is down Breaks the purchase The message waits in its subscription
Adding a consumer Modify and deploy the shop Create a subscription. Zero changes
Traffic peak Every system has to withstand the peak Pub/Sub absorbs it; each one consumes at its own pace
Reprocessing a day Impossible without ad hoc scripts seek to an earlier instant
Traceability Logs scattered around Metrics per subscription

The technical expression of this is decoupling: the producer does not know the consumers, does not know how many there are, does not wait for their reply and does not fail if they fail. The price is that the system becomes eventually consistent: when the shop replies "order confirmed", the warehouse does not know yet. It will know in a few milliseconds, or in a few seconds if it was busy. That nuance has to be accepted consciously, because it changes how the screens are designed: you cannot show "preparing shipment" immediately after the purchase if the warehouse has not found out yet.

  1. The publish-subscribe model

Four concepts, and it is worth being precise about them because the vocabulary is often misused.

Topic. The channel where things are published. It represents one type of fact: pedidos-nuevos, imagenes-subidas, stock-agotado. A topic stores nothing by itself; it is an entry point.

Subscription. The queue of one specific consumer over a topic. This is where the messages actually live. Each subscription receives an independent copy of every published message, and keeps its own accounting of what it has acknowledged and what it has not.

That last point is the key that is hardest to internalise:

1 message is published to pedidos-nuevos
   └── sub-almacen      receives its copy  → acknowledges it after 0.2 s
   └── sub-facturacion  receives its copy  → acknowledges it after 0.5 s
   └── sub-analitica    receives its copy  → its consumer is down: it waits 6 hours

The warehouse acknowledging has no effect whatsoever on analytics' copy. They are separate queues fed by the same topic.

And the practical corollary: if you connect two instances of your warehouse service to the same sub-almacen subscription, each message will go to one of the two. That is load sharing, and it is what you want in order to scale. If instead you create two different subscriptions for the same service, each instance will process every message and the work will be done twice. One subscription per logical consumer, as many instances inside as you like.

Message. It has three parts:

  • data: the body, bytes (typically UTF-8 encoded JSON). Maximum 10 MB.
  • attributes: text key-value pairs, up to 100. They are metadata, and their great virtue is that they can be filtered on without opening the body.
  • System fields: messageId (unique, assigned by Pub/Sub), publishTime, optional orderingKey.

Acknowledgement (ack). The consumer says "processed, do not send it to me again". Until that acknowledgement arrives, Pub/Sub considers the message outstanding and will redeliver it.

  1. Pub/Sub's real guarantees

This section is the most important in the lesson, because most production mistakes with messaging come from assuming guarantees that do not exist.

Guarantee Does Pub/Sub give it? Practical consequence
At-least-once delivery Yes, by default Your consumer will receive duplicates. You have to accept it
At-most-once delivery No A message is never lost by design
Exactly once Yes, if enabled on the subscription (with conditions) It reduces, but does not remove, the need for idempotency
Arrival order No, except with an ordering key Messages can arrive out of order
Durability Yes, replicated across several zones They are not lost even if a zone goes down
Retention 7 days by default, up to 31 They can be replayed
Ordered delivery across different topics No, under no circumstances Do not assume a temporal relationship between topics

The first two rows deserve elaboration, because they are counter-intuitive.

Why there will be duplicates. The scenario is simple: your consumer receives the message, processes it correctly, sends the acknowledgement, and the acknowledgement is lost on the network. Pub/Sub does not receive it, considers the message still outstanding and redelivers it. Your consumer processes it a second time. There is no fault in your code and it has happened anyway.

It also happens if the consumer takes longer than the ack deadline, if it restarts mid-processing, or if autoscaling reassigns the message. It is normal, it is to be expected and it is not an error. That is why section 10 exists.

Exactly once. Pub/Sub offers subscriptions with this semantics (--enable-exactly-once-delivery), with two important caveats: it only applies within a region, and it guarantees that there will not be a second acknowledged delivery, not that your processing is atomic. If your consumer writes to the database and then crashes before acknowledging, the work is already done and the message will come back. Idempotency is still necessary; exactly once only reduces how often you need it.

Why there is no order. Pub/Sub spreads the messages across many servers in order to scale. If message A is published a millisecond before B, they can end up on different servers and arrive in any order. For AlpinaShop this matters little in pedidos-nuevos — each order is independent — but it would matter enormously if we published state changes for the same order: "confirmado", "enviado", "entregado" out of order would leave the order marked as confirmed after being delivered. That case is solved with an ordering key (section 14).

  1. Creating the pedidos-nuevos topic and its subscriptions

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

# The topic: one type of business fact
gcloud pubsub topics create pedidos-nuevos \
  --message-retention-duration=7d \
  --labels=entorno=produccion,equipo=desarrollo,centro-coste=plataforma,aplicacion=tienda

--message-retention-duration on the topic enables replay: it lets you seek to an instant in the past even for subscriptions created afterwards. Without it, you can only go back within whatever each subscription retains.

Now the dead letter topic, which we will create before the subscriptions because they are going to need it:

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

# And a subscription over it, so we can inspect whatever lands there
gcloud pubsub subscriptions create sub-pedidos-fallidos \
  --topic=pedidos-nuevos-fallidos \
  --message-retention-duration=31d \
  --ack-deadline=600

The three consumption subscriptions:

# 1) WAREHOUSE: pull, fast processing, tolerates duplicates through idempotency
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) BILLING: exactly once, because duplicating an invoice is serious
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) ANALYTICS: it will be consumed by the Dataflow pipeline from 04-02
gcloud pubsub subscriptions create sub-analitica \
  --topic=pedidos-nuevos \
  --ack-deadline=120 \
  --message-retention-duration=7d \
  --labels=consumidor=analitica

Notes on the differences, which are deliberate:

  • sub-facturacion with exactly once. Issuing two invoices for the same order is a real accounting problem. Although the consumer must be idempotent anyway, this guarantee reduces the exposure. It has a cost: a lower maximum throughput per subscription.
  • sub-analitica with no dead letter. The Dataflow pipeline already has its own quarantine mechanism (04-02): malformed messages go to pedidos_streaming_errores. Duplicating the mechanism would complicate diagnosis.
  • Different ack deadline values. 30 s for the warehouse (a fast operation), 120 s for analytics (Dataflow processes in internal batches).

The dead letter needs explicit permissions, and this is always forgotten:

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

# Pub/Sub needs to be able to PUBLISH to the dead letter topic...
gcloud pubsub topics add-iam-policy-binding pedidos-nuevos-fallidos \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.publisher"

# ...and to ACKNOWLEDGE on the source subscription in order to remove the message
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"

Without these two permissions, the dead letter configuration is accepted without complaint and does not work: the messages are retried forever instead of being diverted. It is a classic silent failure.

Verification:

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

  1. Publishing from the Flask application

Now Dani's part, the backend developer. Publishing from the catalogue application:

"""publicador.py -- Publishing order events to Pub/Sub."""
import json
import logging
from concurrent import futures

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

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

# 1) BATCH CONFIGURATION.
#    The client accumulates messages and sends them together: fewer network
#    calls, less cost. It publishes when the FIRST of the three limits is met.
batch_config = pubsub_v1.types.BatchSettings(
    max_messages=100,        # either 100 messages...
    max_bytes=1024 * 1024,   # ...or 1 MB accumulated...
    max_latency=0.05,        # ...or 50 ms of waiting. Whichever comes first.
)

# 2) RETRY CONFIGURATION with exponential backoff.
retry_config = retry.Retry(
    initial=0.1,     # first retry after 100 ms
    maximum=60.0,    # never wait more than 60 s between attempts
    multiplier=2.0,  # 0.1 -> 0.2 -> 0.4 -> 0.8 ...
    deadline=600.0,  # it gives up after 10 minutes
)

# 3) The client is EXPENSIVE to create: one per process, reused.
#    Creating it inside the purchase function would be a serious performance bug.
publisher = pubsub_v1.PublisherClient(batch_settings=batch_config)
topic_path = publisher.topic_path(PROJECT, TOPIC)

logger = logging.getLogger(__name__)


def publish_order(order: dict) -> futures.Future:
    """Publishes a confirmed-order event. It does NOT block.

    It returns a Future. The application can carry on replying to the customer
    without waiting for Pub/Sub's acknowledgement.
    """
    body = json.dumps(order, ensure_ascii=False).encode("utf-8")

    future = publisher.publish(
        topic_path,
        data=body,
        # ATTRIBUTES: metadata that can be filtered without opening the body
        # (section 12). Every value must be a STRING.
        tipo_evento="pedido_confirmado",
        origen=order.get("canal", "web"),
        pais=order["envio"]["pais"],
        version_esquema="1",
        momento_evento=order["momento"],   # Dataflow will use this (04-02)
        pedido_id=order["pedido_id"],
    )

    def _on_complete(fut):
        try:
            message_id = fut.result()
            logger.info("Published %s as %s", order["pedido_id"], message_id)
        except Exception:
            # CRITICAL: if publication fails for good, it has to be
            # logged so it can be recovered. An error log here
            # must trigger an alert (06-04).
            logger.exception("FAILED TO PUBLISH order %s",
                             order["pedido_id"])

    future.add_done_callback(_on_complete)
    return future

And how it is used in the Flask view:

from flask import Flask, jsonify, request

app = Flask(__name__)


@app.post("/api/pedidos")
def confirm_order():
    data = request.get_json()

    # 1) The ESSENTIAL and synchronous part: persist and take payment.
    #    If this fails, the customer must find out.
    order = save_to_cloud_sql(data)
    charge_payment_gateway(order)

    # 2) The DERIVED part: the fact is published and we move on.
    #    Warehouse, billing, mail and analytics will find out on their own.
    publish_order({
        "pedido_id": order.id,
        "momento": order.creado_en.isoformat(),
        "cliente_id": order.cliente_id,
        "canal": order.canal,
        "total": str(order.total),
        "envio": {"pais": order.pais, "ciudad": order.ciudad},
        "lineas": [
            {"sku": l.sku, "cantidad": l.cantidad, "precio": str(l.precio)}
            for l in order.lineas
        ],
    })

    # 3) Immediate reply: we do not wait for Pub/Sub or for the consumers.
    return jsonify({"pedido_id": order.id, "estado": "confirmado"}), 201

Four design questions that deserve to be made explicit:

What goes inside the message. Here we publish a fat event, with the lines included. The alternative is a thin event with just the pedido_id, forcing every consumer to query the database. The fat one avoids that read load but couples the message schema; the thin one is more flexible but multiplies the queries against Cloud SQL. For AlpinaShop, with ~1,200 orders a month, the fat event is clearly better: fewer calls to the operational database and self-sufficient consumers.

version_esquema as an attribute. The day the message structure changes, old consumers will be able to recognise and reject what they do not understand instead of breaking. It costs one attribute and it saves an incident.

The momento_evento. It is the attribute the Dataflow pipeline from 04-02 uses as its timestamp_attribute. Without it, the event-time windows do not work.

What happens if the publication fails. It is the real risk of this design: the order is charged and the notification never went out. For AlpinaShop, logging with an alert is enough given the volume. In more demanding systems the outbox pattern is used: the event is written to a table in the same database within the order's own transaction, and a separate process publishes it. That way the event and the order are atomic.

A quick test from the command line:

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. Pull consumption with the asynchronous client

In pull mode, the consumer asks for messages. The Python library uses streaming pull: it keeps a connection open and receives messages as they arrive, with very little latency.

"""consumidor_almacen.py -- AlpinaShop warehouse service."""
import json
import logging
import signal
import sys

from google.cloud import pubsub_v1

PROJECT = "alpinashop-datos"
SUBSCRIPTION = "sub-almacen"

logger = logging.getLogger(__name__)

# FLOW CONTROL: limits how many unacknowledged messages you hold at once.
# Without this, the client can accept thousands of messages, fail to get
# through them in time, let the ack deadline expire and cause mass redeliveries.
flow_control = pubsub_v1.types.FlowControl(
    max_messages=50,
    max_bytes=10 * 1024 * 1024,
)


def process(message):
    """Callback run by the library on a thread from the pool."""
    try:
        order = json.loads(message.data.decode("utf-8"))
    except json.JSONDecodeError:
        # Corrupt message: retries will NOT fix it.
        # It is acknowledged to remove it and logged for investigation.
        logger.error("Unreadable message, discarded: %s", message.message_id)
        message.ack()
        return

    event_type = message.attributes.get("tipo_evento")
    if event_type != "pedido_confirmado":
        # Not for us; we acknowledge it without doing anything.
        message.ack()
        return

    try:
        # Idempotency is implemented in here (section 10)
        create_picking_order(order)
        message.ack()
        logger.info("Picking order created for %s", order["pedido_id"])

    except TemporaryError as exc:
        # Transient failure (saturated DB, network): NACK to retry now.
        logger.warning("Temporary failure on %s: %s", order["pedido_id"], exc)
        message.nack()

    except Exception:
        # Unknown failure: NACK. After 5 attempts it will go to the dead letter.
        logger.exception("Failure processing %s", order.get("pedido_id"))
        message.nack()


def main():
    subscriber = pubsub_v1.SubscriberClient()
    path = subscriber.subscription_path(PROJECT, SUBSCRIPTION)

    future = subscriber.subscribe(path, callback=process,
                                  flow_control=flow_control)
    logger.info("Listening on %s...", path)

    # Clean shutdown: on receiving SIGTERM (Kubernetes, Cloud Run),
    # it stops accepting new messages and finishes the ones in hand.
    def shutdown(signum, frame):
        logger.info("Signal %s received, closing...", signum)
        future.cancel()
        future.result(timeout=30)
        sys.exit(0)

    signal.signal(signal.SIGTERM, shutdown)
    signal.signal(signal.SIGINT, shutdown)

    with subscriber:
        try:
            future.result()
        except Exception:
            logger.exception("The subscriber has failed")
            future.cancel()


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

Three points that make the difference between a toy consumer and a production one:

Flow control. Without FlowControl, the library accepts as many messages as it is sent. If your processing is slow, many will pass their deadline, be redelivered, be processed again by your consumer, and you will enter a spiral in which there is more and more duplicated work. Limiting the messages in flight is what avoids that spiral.

ack() versus nack(). ack() removes the message for good. nack() returns it for immediate redelivery. If you do neither, the message is redelivered when the deadline expires, which delays the retry but works just the same.

Telling recoverable errors from unrecoverable ones. A corrupt JSON will not improve by retrying it five times: it is discarded with a log entry. A saturated database will: nack(). Confusing them means either that poisoned messages consume resources indefinitely or that recoverable data is lost.

There is also synchronous pull, useful for batch processes:

# Read up to 10 messages without acknowledging them (to inspect)
gcloud pubsub subscriptions pull sub-almacen --limit=10 --format=json

# Read them and acknowledge them
gcloud pubsub subscriptions pull sub-almacen --limit=10 --auto-ack

  1. Push consumption to an HTTPS endpoint with OIDC

In push mode, Pub/Sub makes an HTTPS POST to a URL of yours. You keep no process listening: the service calls you.

Aspect Pull Push
Who initiates The consumer Pub/Sub
Infrastructure A process always running An HTTPS endpoint; it can scale to zero
Pace control Total, with FlowControl Limited (automatic sliding window)
Acknowledgement Explicit (ack()) Implicit: HTTP 2xx acknowledges, any other code is a nack
Ideal for High volume, continuous processing Cloud Run, Cloud Functions, moderate volume

Push fits perfectly with what is coming in the course: Cloud Run (07-02, where the catalogue will end up according to decision DA-001) and Cloud Functions (06-03) scale to zero, so there is no sense in having a process sitting waiting for messages.

The secure configuration uses OIDC authentication: Pub/Sub attaches a token signed by Google identifying a service account, and your endpoint verifies it. Without this, anyone who discovered your URL could inject fake messages.

# 1) Service account that will represent Pub/Sub to your endpoint
gcloud iam service-accounts create sa-pubsub-invocador \
  --display-name="Pub/Sub identity for invoking push services"

SA_INV="[email protected]"

# 2) Allow it to invoke the warehouse's Cloud Run service
gcloud run services add-iam-policy-binding svc-almacen \
  --region=europe-west1 \
  --member="serviceAccount:${SA_INV}" \
  --role="roles/run.invoker"

# 3) Allow Pub/Sub to mint tokens on behalf of that account
PROJ_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
gcloud iam service-accounts add-iam-policy-binding "$SA_INV" \
  --member="serviceAccount:service-${PROJ_NUM}@gcp-sa-pubsub.iam.gserviceaccount.com" \
  --role="roles/iam.serviceAccountTokenCreator"

# 4) Push subscription with 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

The receiving endpoint:

import base64
import json

from flask import Flask, request

app = Flask(__name__)


@app.post("/eventos/pedidos")
def receive_event():
    """Pub/Sub push endpoint.

    Cloud Run has already validated the OIDC token before it reaches here,
    because the service requires authentication and the invoking account
    has run.invoker. If the service were public, the Authorization
    header would have to be verified manually.
    """
    envelope = request.get_json(silent=True)
    if not envelope or "message" not in envelope:
        # 400: it will not be retried. It is a format error, not a transient one.
        return "malformed request", 400

    message = envelope["message"]

    # The body arrives base64-encoded inside the envelope
    body = base64.b64decode(message.get("data", "")).decode("utf-8")
    attributes = message.get("attributes", {})
    message_id = message["messageId"]
    attempt = int(envelope.get("deliveryAttempt", 1))

    try:
        order = json.loads(body)
    except json.JSONDecodeError:
        # 200 without processing: we acknowledge so that a message
        # that will never be valid is NOT retried.
        app.logger.error("Message %s unreadable, discarded", message_id)
        return "", 204

    try:
        create_picking_order(order)
    except TemporaryError:
        app.logger.warning("Temporary failure on %s, attempt %s",
                           message_id, attempt)
        # 500: Pub/Sub will retry with exponential backoff
        return "retry", 500

    # 204: acknowledged
    return "", 204

The golden rule of push: the HTTP response code is the acknowledgement. 2xx removes the message; anything else (or a timeout) puts it back in the queue. Returning 200 from a generic except "so it stops bothering us" is the fastest way to lose data silently.

  1. Acknowledgements, deadlines and retries

The ack deadline is how long Pub/Sub waits for the acknowledgement before redelivering. 10 seconds by default; configurable between 10 and 600.

sequenceDiagram
    participant PS as Pub/Sub
    participant C as Consumer
    PS->>C: delivers message M (deadline 30 s)
    Note over C: processing... 25 s
    C->>PS: modifyAckDeadline(+60 s)
    Note over C: still processing... 40 s
    C->>PS: ack(M)
    Note over PS: M removed from this subscription

If the consumer neither acknowledges nor extends the deadline, the message goes back into the queue. Choosing the deadline well matters:

  • Too short: messages that are being processed correctly get redelivered and the work is duplicated.
  • Too long: if a consumer dies, its messages take a long time to be reassigned to another.

Rule of thumb: the 99th percentile of your processing time, with margin. If 99 % of picking orders are created in under 8 seconds, 30 seconds is reasonable.

The good news is that the Python library extends the deadline automatically while your callback is still running (lease management), up to a configurable maximum. If you need to extend it by hand:

# Extend the deadline from inside the processing
message.modify_ack_deadline(120)

The retries follow the policy you configured:

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

With exponential backoff, the retries are spaced out: 10 s, 20 s, 40 s, 80 s… up to 600 s. This is the right thing when the failure is because a system is saturated: retrying every second would sink it further. Without backoff, a consumer that is down generates a storm of retries that prevents it from recovering.

  1. Dead letter topics

Imagine a message with a sku that does not exist in the catalogue. The warehouse consumer fails. It retries. It fails. It retries. Forever. That message is called poisoned, and without a dead letter it has three consequences: it consumes resources indefinitely, it clutters the logs, and — if ordering is enabled — it blocks every subsequent message with its key.

The dead letter topic solves this: after N attempts, the message is diverted to another topic and removed from the original subscription.

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

When a message reaches the dead letter, it keeps its body and its attributes and gains metadata about the origin and the number of attempts. A review process, typically weekly:

"""revisar_fallidos.py -- Inspecting the dead letter queue."""
import json

from google.cloud import pubsub_v1

subscriber = pubsub_v1.SubscriberClient()
path = subscriber.subscription_path("alpinashop-datos", "sub-pedidos-fallidos")

response = subscriber.pull(
    request={"subscription": path, "max_messages": 100},
    timeout=30,
)

for received in response.received_messages:
    m = received.message
    print("---")
    print("Original ID  :", m.attributes.get(
        "CloudPubSubDeadLetterSourceMessageId", "?"))
    print("Subscription :", m.attributes.get(
        "CloudPubSubDeadLetterSourceSubscription", "?"))
    print("Attempts     :", m.attributes.get(
        "CloudPubSubDeadLetterSourceDeliveryCount", "?"))
    print("Published    :", m.publish_time)
    print("Body         :", m.data.decode("utf-8")[:300])

    # Manual decision: fix and republish, or discard with a log entry.

Dead letter rules worth fixing as team policy:

  1. Every production subscription has one. No exceptions. It is the difference between "something failed and we have it saved" and "something failed and we do not know what it was".
  2. max-delivery-attempts between 5 and 10. Fewer, and a transient failure sends good messages to the dead letter queue. More, and it takes too long to spot the problem.
  3. An alert on the number of messages in the dead letter. A dead letter nobody looks at is a black hole with extra steps.
  4. Do not point a dead letter at the original topic. That is an infinite loop, and it is a mistake made more often than you would think.

  1. Idempotency: the consumer must tolerate duplicates

We already know there will be duplicates. The solution is not to avoid them — you cannot — but to make processing the same message twice have the same effect as processing it once. That is idempotency.

Three strategies, from least to most robust.

A. Naturally idempotent operations

The best one, when it is possible: design the operation so that repeating it changes nothing.

# NOT idempotent: subtracting. Twice subtracts double.
db.execute("UPDATE stock SET unidades = unidades - %s WHERE sku = %s",
           (quantity, sku))

# Idempotent: set an absolute value computed from the source of truth.
db.execute("UPDATE stock SET unidades = %s WHERE sku = %s",
           (computed_units, sku))

B. Conditional insert by business key

Use a uniqueness constraint in the database:

def create_picking_order(order):
    """The UNIQUE index on pedido_id does the work for us."""
    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
            """,
            (order["pedido_id"],),
        )
        # If it already existed, ON CONFLICT does nothing and there is no error.
        # The duplicate message is processed with no effect: idempotent.

This is the preferable option when a natural business key exists, as pedido_id does here. Zero added infrastructure and the database provides the guarantee.

C. A register of processed messages

When the operation is neither idempotent nor has a natural key (sending an email, calling an external API), you have to keep a register. Firestore, which AlpinaShop already uses for the shopping cart (02-06), is ideal thanks to its low latency and its transactions:

"""idempotencia.py -- Register of already processed messages."""
import datetime

from google.cloud import firestore

db = firestore.Client(project="alpinashop-datos")
COLLECTION = "eventos_procesados"
TTL_DAYS = 14   # longer than Pub/Sub's maximum retention (7 days)


def process_once(idempotency_key: str, consumer: str, action):
    """Runs `action` only if this key has not been processed before.

    The key includes the consumer: the same order must be processed
    once in the warehouse AND once in billing. They are different marks.
    """
    doc_id = f"{consumer}__{idempotency_key}"
    ref = db.collection(COLLECTION).document(doc_id)

    @firestore.transactional
    def _reserve(tx):
        snapshot = ref.get(transaction=tx)
        if snapshot.exists:
            return False          # already processed: we do nothing
        tx.set(ref, {
            "consumidor": consumer,
            "clave": idempotency_key,
            "procesado_en": firestore.SERVER_TIMESTAMP,
            # TTL field: Firestore deletes the document automatically
            "caduca_en": datetime.datetime.utcnow()
                         + datetime.timedelta(days=TTL_DAYS),
        })
        return True

    first_time = _reserve(db.transaction())
    if not first_time:
        return False

    action()
    return True

And how it is used:

def process_billing(message):
    order = json.loads(message.data.decode("utf-8"))

    # The idempotency key is the BUSINESS identifier,
    # not Pub/Sub's message_id. The reason: if the order is republished
    # by a reprocessing run with seek, the message_id will be different but
    # the order is the same and it must NOT be invoiced twice.
    executed = process_once(
        idempotency_key=order["pedido_id"],
        consumer="facturacion",
        action=lambda: issue_invoice(order),
    )

    if not executed:
        logger.info("Order %s already invoiced, the duplicate is ignored",
                    order["pedido_id"])

    message.ack()

The comment about the key is the most important part of this whole section. Using message_id protects against Pub/Sub redeliveries, but not against a reprocessing run with seek nor against a republication from the application. Using the business identifier (pedido_id) protects against both. It is the difference between a system that survives a reprocessing run and one that invoices twice the day somebody replays an hour of messages.

A practical detail: the register's TTL must be longer than the subscription's maximum retention. If Pub/Sub can redeliver for 7 days and your register expires after 3, a message redelivered on day 5 would be processed again.

  1. Retention, seek and snapshots

Messages are retained in the subscription for between 10 minutes and 31 days (7 by default). With seek you move the read point.

# 1) Reprocess the last 3 hours: useful after fixing a bug in the consumer
gcloud pubsub subscriptions seek sub-analitica \
  --time="$(date -u -d '3 hours ago' '+%Y-%m-%dT%H:%M:%SZ')"

# 2) Discard EVERYTHING pending: useful when unblocking an unusable queue
gcloud pubsub subscriptions seek sub-almacen --time="$(date -u '+%Y-%m-%dT%H:%M:%SZ')"

Real cases where seek saves the day:

  • A bug in the analytics consumer computed VAT wrongly for six hours. You fix the code, deploy, seek back six hours and reprocess everything. With idempotent consumers this is safe. Without idempotency it is a disaster.
  • A queue accumulated 400,000 obsolete messages because of a consumer that was down over the weekend. seek to the present discards them in one go.

Snapshots capture a subscription's acknowledgement state so you can return to it:

# BEFORE deploying a risky version of the consumer
gcloud pubsub snapshots create snap-analitica-pre-v2 --subscription=sub-analitica

# ... it is deployed, and if it goes badly ...
gcloud pubsub subscriptions seek sub-analitica --snapshot=snap-analitica-pre-v2

# Cleanup (snapshots expire after 7 days, but it is better to be explicit)
gcloud pubsub snapshots delete snap-analitica-pre-v2

It is the equivalent of a restore point before a risky deployment, and it should be part of the deployment procedure for any consumer that writes to business systems.

  1. Subscription filters by attribute

A subscription can filter on the message attributes, so that it only receives what interests it. The filtering happens on Pub/Sub's side: discarded messages are neither delivered nor billed as a delivery.

# Subscription that only receives orders from Spain
gcloud pubsub subscriptions create sub-almacen-es \
  --topic=pedidos-nuevos \
  --message-filter='attributes.pais = "ES"'

# Only large orders, for manual fraud review
gcloud pubsub subscriptions create sub-revision-fraude \
  --topic=pedidos-nuevos \
  --message-filter='attributes.tipo_evento = "pedido_confirmado" AND attributes.importe_alto = "true"'

# Everything except the telephone channel
gcloud pubsub subscriptions create sub-analitica-digital \
  --topic=pedidos-nuevos \
  --message-filter='NOT (attributes.origen = "telefono")'

# Prefix: any event whose type starts with "pedido_"
gcloud pubsub subscriptions create sub-todo-pedidos \
  --topic=eventos-tienda \
  --message-filter='hasPrefix(attributes.tipo_evento, "pedido_")'

Available operators: =, !=, AND, OR, NOT, hasPrefix(), and attributes:key to check for existence.

Two limitations you have to know about:

  • You can only filter on attributes, never on the body. If you want to filter on the amount, you have to publish it as an attribute. That is why the publisher in section 5 includes pais and origen as attributes even though they are also inside the JSON.
  • The filter is immutable. It cannot be modified once the subscription is created; you have to create another one.

Filters allow a very clean pattern: one topic per domain, with many event types, and each consumer filtering out its own. For AlpinaShop, one eventos-tienda topic with pedido_confirmado, pedido_cancelado, carrito_abandonado and stock_bajo is more manageable than four topics, because it preserves the relative order of events from the same domain and simplifies the publisher.

  1. BigQuery and Cloud Storage subscriptions

Here comes one of the service's best features, and the one that saves the most code.

A BigQuery subscription writes the messages straight into a table. No Dataflow, no code, nothing to deploy.

# 1) Destination table with the message's schema
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) Permission for Pub/Sub to write
PROJ_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROJ_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) The subscription
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: JSON is expected whose fields match the columns. The alternative, --use-topic-schema, validates against an Avro or Protobuf schema registered on the topic.
  • --drop-unknown-fields: JSON fields that do not exist as a column are ignored instead of causing a failure. Essential if the message schema can evolve.
  • The dead letter here is fundamental: without it, a message that does not fit the schema is retried indefinitely.

The honest comparison:

Criterion BigQuery subscription Dataflow (04-02)
Code None A pipeline to maintain
Cost Only Pub/Sub's plus the write Workers 24×7: tens of €/month
Transformation None: the JSON goes in as-is Any
Windows and event time No Yes
Enrichment from other sources No Yes
Aggregation in flight No Yes
When to use it The message already has the table's shape You have to transform, aggregate or window

AlpinaShop's decision: to dump the raw events into a landing table, a BigQuery subscription, because it is free in effort and there is nothing to maintain. The Dataflow pipeline is reserved for what genuinely needs transformation: the aggregate per hour and channel with event-time windows and tolerance for late data. It is the usual and correct pattern: raw data by the cheap route, aggregates by the powerful route.

The Cloud Storage subscription does the equivalent with files, batching messages by time or size:

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

It is an excellent way to keep an immutable archive of every event in the lake, later processed by Spark (04-03) or loaded into BigQuery. And it is the ultimate safety net: if one day the whole analytical warehouse has to be rebuilt, the raw events are there.

  1. Ordering by key and its implications

By default there is no order. To guarantee it, you publish with an ordering key:

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

# All the events for the SAME order share a key: they arrive in order
publisher.publish(
    topic_path,
    data=body,
    ordering_key=order["pedido_id"],      # <-- the key
    tipo_evento="pedido_enviado",
)
gcloud pubsub subscriptions create sub-estados-pedido \
  --topic=eventos-tienda \
  --enable-message-ordering

The guarantee is: messages with the same ordering key, published in the same region, are delivered in publication order. Messages with different keys have no relationship with each other, which allows parallelism to continue.

The three costs of enabling ordering, which have to be weighed up:

  1. Lower throughput. The messages for one key are processed sequentially. With pedido_id as the key there is no problem: there are thousands of distinct orders and parallelism is preserved. With pais as the key, every order from Spain would go in series.
  2. Blocking on a poisoned message. If a message with key PED-2026-0042 fails repeatedly, every subsequent message for that order is blocked until it is resolved or diverted to the dead letter. Ordering without a dead letter is a time bomb.
  3. Slower publication. The client has to serialise each key's sends, which reduces the effect of batching.

Golden rule: enable ordering only when the order is semantically necessary, and choose the key with the highest possible cardinality that preserves that order. For AlpinaShop: pedido_id yes (an order's states must go in order); pais no.

And the alternative that is usually better: make the order irrelevant. If each event includes a version number or a timestamp, the consumer can discard those that are older than what it has already processed, and arrival order stops being a problem.

# Consumer that tolerates out-of-order arrival, with no need for ordering_key
def apply_status(pedido_id, new_status, event_version):
    db.execute(
        """
        UPDATE pedidos SET estado = %s, version = %s
        WHERE pedido_id = %s AND version < %s
        """,
        (new_status, event_version, pedido_id, event_version),
    )
    # If an old event arrives, the condition version < %s is not met
    # and the UPDATE affects no rows. Disorder absorbed.

  1. Pub/Sub Lite and other alternatives

Pub/Sub Lite was a lower-cost variant, with capacity provisioned by the user rather than served on demand, intended for very high and predictable volumes. It required sizing partitions and storage by hand, in exchange for a considerably lower price per GB.

Google announced its retirement and the service stopped being available in March 2026. If you find references to Pub/Sub Lite in old documentation or courses, they are obsolete; always verify the status in the official documentation. The current alternatives are standard Pub/Sub or, if you need the Kafka API, Google Cloud Managed Service for Apache Kafka.

The comparison that is still useful:

Service Model When it makes sense
Pub/Sub Global, unprovisioned, pay per use The default: the vast majority of cases
Managed Kafka Managed Kafka, with partitions and consumer groups You already have Kafka code, or you need its log semantics
Cloud Tasks Task queue with scheduling and rate control Queuing work with a single, known recipient
Eventarc Routing of platform events Reacting to Google service events (it uses Pub/Sub underneath)
Memorystore (Redis Pub/Sub) In-memory messaging, no persistence Ephemeral notifications where losing messages is acceptable

The most common confusion is Pub/Sub versus Cloud Tasks. Pub/Sub is for facts that interest N unknown consumers. Cloud Tasks is for assignments addressed to a specific recipient, with fine control of the delivery rate and deferred scheduling. "An order has been confirmed" is Pub/Sub. "Send this email in 30 minutes, at a maximum of 10 per second" is Cloud Tasks.

  1. Monitoring and cost

The metrics to watch, with their alerts:

Metric What it indicates Reasonable alert
subscription/oldest_unacked_message_age Age of the oldest outstanding message The queen of metrics. > 600 s: the consumer cannot keep up or is down
subscription/num_undelivered_messages Messages piled up undelivered Sustained growth: the consumers are falling behind
subscription/dead_letter_message_count Messages diverted to the dead letter > 0 in an hour: it needs looking at
topic/send_request_count Publications A sharp drop: the application has stopped publishing
subscription/push_request_count by code Health of the push endpoint Lots of 5xx: the endpoint is failing
subscription/ack_message_count Acknowledgement rate Compare it with publications

The essential alert:

gcloud alpha monitoring policies create \
  --notification-channels="$INFRA_CHANNEL" \
  --display-name="Pub/Sub: unacknowledged messages on 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 is the best metric because it detects all three possible failures at once: a consumer that is down (it grows indefinitely), a slow consumer (it grows slowly) and a poisoned message blocking an ordered key (it sticks at a high value).

The cost is billed mainly by data volume:

Item Order of magnitude (verify in the official documentation)
Publication + delivery ~$40 per TiB, with a free monthly tier
Storage of retained messages ~$0.27 per GiB per month
Cross-region egress Network rates

A realistic calculation for AlpinaShop: 1,200 orders a month, messages of ~2 KB, with 4 subscriptions. That is 1,200 publications and 4,800 deliveries: around 12 MB a month. Cents, or simply within the free tier. The catalogue browsing events, far more numerous, would still be cheap.

One billing detail that surprises people: each subscription counts as a delivery. Ten subscriptions on the same topic multiply the billed delivery volume by ten. That is not a reason to avoid subscriptions, but it is a reason not to leave orphan subscriptions lying around: a forgotten subscription with no consumer piles up messages, its storage is billed and it serves no purpose.

# Look for subscriptions with no consumption: candidates for deletion
gcloud pubsub subscriptions list --format="table(name, topic)" | while read -r s _; do
  echo "$s"
done

  1. The alpinashop-catalogo bucket notifications

Back in 02-02 a promise was left outstanding: that Cloud Storage can give notice when a new object appears. Now we have the pieces to keep it properly.

AlpinaShop's case: when someone uploads a product's original image to productos/<sku>/original/, the web/ and thumb/ versions have to be generated automatically. Today Marta does that by hand with a script.

# 1) Topic for the bucket's events
gcloud pubsub topics create imagenes-subidas \
  --labels=entorno=produccion,equipo=infra,aplicacion=catalogo

# 2) Permission for the Cloud Storage agent to publish
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) The notification, scoped to the prefix and the event we care about
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) Subscription for the image processor
gcloud pubsub subscriptions create sub-procesar-imagenes \
  --topic=imagenes-subidas \
  --ack-deadline=300 \
  --dead-letter-topic=pedidos-nuevos-fallidos \
  --max-delivery-attempts=5

Available event types:

Event When it fires
OBJECT_FINALIZE A new object is created or overwritten
OBJECT_DELETE It is deleted (or the previous version is overwritten)
OBJECT_ARCHIVE A version becomes archived (with versioning enabled)
OBJECT_METADATA_UPDATE The metadata changes

The consumer:

def process_image(message):
    """Generates the web and thumb versions of a freshly uploaded image."""
    # The object's details come in the ATTRIBUTES; the body is not needed
    bucket = message.attributes["bucketId"]
    name = message.attributes["objectId"]
    generation = message.attributes["objectGeneration"]

    # 1) INFINITE LOOP FILTER: if we do not check this, writing
    #    web/ and thumb/ would fire new notifications that would generate
    #    more images, indefinitely. It is the classic mistake.
    if "/original/" not in name:
        message.ack()
        return

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

    # 2) IDEMPOTENCY: the generation identifies the object's exact version.
    #    If the message is redelivered, the key is the same and nothing repeats.
    key = f"{bucket}/{name}#{generation}"

    process_once(
        idempotency_key=key,
        consumer="procesador-imagenes",
        action=lambda: generate_versions(bucket, name),
    )
    message.ack()

The infinite loop filter deserves emphasis: it is the most expensive failure in this pattern. Writing to the same bucket that fires the notification creates a recursion that is only detected when the invoice arrives. The three defences are: filtering by prefix on the notification (--object-prefix=productos/), checking the path in the consumer, and — better still — writing the output to a different bucket from the input one.

This same pattern applies to the carrier's monthly file and to the nightly exports: as soon as the file lands in exportaciones/<yyyy>/<mm>/<dd>/, a notification triggers the load into BigQuery. No cron, no checking folders every five minutes, no waiting windows. We will put the processing logic in Cloud Functions in 06-03, and who orchestrates it all is the next lesson.

Common Mistakes and Tips

Assuming there will be no duplicates. There will be, with or without exactly once. Every production consumer must be idempotent. It is not optional.

Using message_id as the idempotency key. It protects against redeliveries, but not against seek or republications. Use the business identifier.

Configuring a dead letter and forgetting the permissions. The command is accepted and the feature does not work: messages are retried forever. Always check both bindings.

Not putting flow control in the consumer. With slow processing, the library accepts more messages than it can acknowledge, the deadlines expire, they are redelivered, and the consumer sinks processing duplicates of its own backlog.

Creating the Pub/Sub client inside the request handler. It is expensive (it opens gRPC connections). One per process, reused.

Returning 200 from a generic except in push. It acknowledges the message and loses it silently. Return 5xx on transient failures.

Enabling ordering without needing it. It reduces throughput and, with a poisoned message, blocks the whole key.

Pointing the dead letter at the original topic. Infinite loop.

Orphan subscriptions. With no consumer, they pile up messages until the maximum retention, their storage is billed and they contribute nothing. Review them periodically.

Bucket notifications that feed themselves. Writing to the same bucket that fires the event. Use separate buckets or prefixes and filter in the consumer.

Tip: publish attributes generously. They cost little and they allow server-side filtering, which saves deliveries and complexity in the consumer.

Tip: a schema on the topic. Pub/Sub can register an Avro or Protobuf schema and validate messages at publication time. It turns format errors into immediate publisher failures rather than surprises in the consumer.

Tip: take a snapshot before every risky deployment. It costs one command and it gives you a way back.

Exercises

Exercise 1: a customer review topic with filters and a dead letter

Create the opiniones-nuevas topic with 7 days of retention, its dead letter topic opiniones-nuevas-fallidos with an inspection subscription, and three subscriptions over the main topic:

  • sub-moderacion: receives only the reviews with attributes.puntuacion_baja = "true", with a 60 s deadline, dead letter after 5 attempts and backoff from 10 s to 300 s.
  • sub-analitica-opiniones: receives them all, with a 120 s deadline.
  • sub-bq-opiniones: writes straight into BigQuery in the table alpinashop-datos:alpinashop_analitica.opiniones, with no code.

Include every necessary permission and publish two test messages, one that reaches moderation and one that does not.

Exercise 2: an idempotent billing consumer

Write the pull consumer for sub-facturacion that issues each order's invoice. Requirements: flow control of 20 messages in flight; distinguishing transient errors (retry) from permanent ones (discard with a log entry); idempotency based on the business identifier with a register in Firestore and a suitable TTL; clean shutdown on SIGTERM; and a counter of invoices issued and duplicates detected. Explain why the idempotency key must be the pedido_id and not the message_id, with a concrete scenario in which the wrong choice would cause a real problem.

Exercise 3: diagnosing a stuck queue

At 09:15 the alert fires: oldest_unacked_message_age on sub-almacen has been at 47 minutes and rising. num_undelivered_messages has gone from 3 to 12,400 since 08:20. The warehouse service is up and answers its health check. In the consumer's logs the same error message appears about a hundred times a minute, mentioning order PED-2026-1188. The pedidos-nuevos-fallidos queue is empty. The subscription has ordering enabled by pedido_id.

Diagnose what is happening, explain why the dead letter is empty despite the retries, and give an action plan in three phases: immediate containment, correction and prevention.

Solutions

Solution 1

gcloud config set project alpinashop-datos
PROJ_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROJ_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) Publish permission on the dead letter topic
gcloud pubsub topics add-iam-policy-binding opiniones-nuevas-fallidos \
  --member="serviceAccount:${SA_PUBSUB}" --role="roles/pubsub.publisher"

# 3) Moderation subscription, with a filter
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) Analytics subscription, with no filter
gcloud pubsub subscriptions create sub-analitica-opiniones \
  --topic=opiniones-nuevas --ack-deadline=120 \
  --labels=consumidor=analitica

# 5) Direct BigQuery subscription
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
# Message that DOES reach moderation (score 1)
gcloud pubsub topics publish opiniones-nuevas \
  --message='{"opinion_id":"OPI-1001","sku":"FRON-300L","fecha":"2026-03-14","puntuacion":1,"texto":"The battery lasts far less than advertised","pais":"FR"}' \
  --attribute=puntuacion_baja=true,sku=FRON-300L,pais=FR

# Message that does NOT reach moderation (score 5)
gcloud pubsub topics publish opiniones-nuevas \
  --message='{"opinion_id":"OPI-1002","sku":"MOCH-40L-AZ","fecha":"2026-03-14","puntuacion":5,"texto":"Extremely comfortable on long treks","pais":"ES"}' \
  --attribute=puntuacion_baja=false,sku=MOCH-40L-AZ,pais=ES

# Check: moderation receives 1, analytics receives 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)"

The point being assessed here is the filter: sub-moderacion receives one of the two publications, because the filtering happens on Pub/Sub's side and the five-star review is neither delivered nor billed.

Solution 2

"""consumidor_facturacion.py -- Idempotent invoice issuing."""
import datetime
import json
import logging
import signal
import sys

from google.cloud import firestore, pubsub_v1

PROJECT = "alpinashop-datos"
SUBSCRIPTION = "sub-facturacion"
COLLECTION = "eventos_procesados"
TTL_DAYS = 14           # > the subscription's 7 days of maximum retention
CONSUMER = "facturacion"

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

invoices_issued = 0
duplicates_detected = 0


class TemporaryError(Exception):
    """Transient failure: worth retrying."""


class PermanentError(Exception):
    """Failure that will not improve by retrying."""


def process_once(key, action):
    doc = db.collection(COLLECTION).document(f"{CONSUMER}__{key}")

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

    if not _reserve(db.transaction()):
        return False
    action()
    return True


def process(message):
    global invoices_issued, duplicates_detected
    try:
        order = json.loads(message.data.decode("utf-8"))
    except json.JSONDecodeError:
        logger.error("Message %s unreadable, discarded", message.message_id)
        message.ack()                      # permanent: do not retry
        return

    if not order.get("pedido_id"):
        logger.error("Message %s has no pedido_id, discarded", message.message_id)
        message.ack()
        return

    try:
        issued = process_once(
            key=order["pedido_id"],
            action=lambda: issue_invoice(order),
        )
        if issued:
            invoices_issued += 1
            logger.info("Invoice issued for %s", order["pedido_id"])
        else:
            duplicates_detected += 1
            logger.info("Duplicate ignored: %s", order["pedido_id"])
        message.ack()

    except TemporaryError as exc:
        logger.warning("Temporary failure on %s: %s", order["pedido_id"], exc)
        message.nack()                     # retry with backoff

    except PermanentError as exc:
        logger.error("Permanent failure on %s: %s", order["pedido_id"], exc)
        message.ack()                      # no fix possible: it is removed

    except Exception:
        logger.exception("Unknown failure on %s", order["pedido_id"])
        message.nack()                     # after 5 attempts, to the dead letter


def main():
    subscriber = pubsub_v1.SubscriberClient()
    path = subscriber.subscription_path(PROJECT, SUBSCRIPTION)
    flow = pubsub_v1.types.FlowControl(max_messages=20)

    future = subscriber.subscribe(path, callback=process, flow_control=flow)
    logger.info("Billing listening on %s", path)

    def shutdown(signum, frame):
        logger.info("Closing. Issued=%s Duplicates=%s",
                    invoices_issued, duplicates_detected)
        future.cancel()
        future.result(timeout=30)
        sys.exit(0)

    signal.signal(signal.SIGTERM, shutdown)
    signal.signal(signal.SIGINT, shutdown)

    with subscriber:
        try:
            future.result()
        except Exception:
            logger.exception("Subscriber down")
            future.cancel()


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

Why pedido_id and not message_id, with a concrete scenario:

message_id is unique per publication. If the same business fact is published twice, each publication will have a different message_id, and a register based on it would detect nothing.

The real scenario: on a Tuesday, a bug in the billing consumer meant that 300 of the morning's orders were not invoiced. The code is fixed, deployed, and this is run:

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

That seek redelivers the morning's messages. Now, between 08:00 and the incident there were 500 orders, of which 200 were invoiced correctly before the failure began.

  • With message_id as the key: redelivery preserves the original message_id, so in this specific case the 200 good ones would indeed be detected. But if instead of seek somebody recovers and republishes the events from the Cloud Storage archive (section 13), the message_id values will be new and 200 duplicate invoices would be issued. That is the failure.
  • With pedido_id as the key: it does not matter how the message arrives — redelivery, seek, republication from the archive, manual reprocessing. If that order has already been invoiced, it is not invoiced again. The protection is over the business fact, which is what actually matters.

Two hundred duplicate invoices sent to customers is an incident for accounting, for customer service and for reputation. The difference between the two options is one line of code.

Solution 3

Diagnosis: a poisoned message blocking an ordering key, with the dead letter misconfigured.

The clues fit one by one:

  1. The service is alive and responding. It is not an outage: it is a logical jam.
  2. The same error a hundred times a minute about PED-2026-1188. That message always fails. It is a poisoned message: something in it (a non-existent SKU, a null field, an impossible amount) makes the consumer fail deterministically.
  3. Ordering is enabled by pedido_id. Here is the critical part that explains the scale: with ordering, messages with the same key are delivered in order and a blocked one prevents progress. But on top of that, in practice, the combination of a consumer that nack()s in a loop over one key and flow control saturated by redeliveries makes overall throughput collapse: the consumer spends its capacity reprocessing the same message over and over instead of attending to the queue.
  4. 12,400 messages piled up in 55 minutes when AlpinaShop takes ~1,200 orders a month: that volume is not new orders, it is redeliveries of the same message plus the rest of the queue going unattended.

Why the dead letter is empty. This is the part that teaches the most. There are two possible causes and both are common:

  • The permissions are missing. The dead letter policy is configured without error even though Pub/Sub's service agent has neither roles/pubsub.publisher on the dead letter topic nor roles/pubsub.subscriber on the source subscription. Without those two permissions, the diversion never happens and the message is retried indefinitely. It is exactly the silent failure we warned about in section 4.
  • The policy is not actually applied. The subscription was created without --dead-letter-topic and it was taken for granted.

Check:

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

Three-phase action plan:

Phase 1 — Containment (minutes). Fix the dead letter permissions, which is what unblocks things without losing anything:

PROJ_NUM=$(gcloud projects describe alpinashop-datos --format="value(projectNumber)")
SA_PUBSUB="service-${PROJ_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

As soon as the diversion works, PED-2026-1188 will leave the queue after five attempts and the rest of the messages will start to flow. Do not seek to the present: that would discard 12,400 messages, most of which are real orders waiting to be prepared.

If the diversion were slow and there were operational urgency, the alternative patch is to deploy a temporary rule in the consumer that explicitly acknowledges the problematic pedido_id, logging it so it can be handled by hand.

Phase 2 — Correction (hours). Inspect the message in the dead letter queue, understand why the consumer fails on it, and fix the code so that this kind of data is treated as a permanent error (ack() with a log entry) instead of being retried. Then handle order PED-2026-1188 manually, since it is still a real order from a real customer and it has to be prepared. Verify that the queue returns to an oldest_unacked_message_age measured in seconds.

Phase 3 — Prevention (days).

  • Verify the dead letter on every subscription, including the permissions, not just the policy. Automate it as a check at deployment time.
  • An alert on dead_letter_message_count > 0, so you find out about the first poisoned message instead of number 12,400.
  • Review whether ordering is really necessary on sub-almacen. Orders are independent of each other; the order between different orders adds nothing and does add a blocking risk. If all that is needed is to order the state changes of a single order, that case can be solved with a version number and an UPDATE ... WHERE version < %s, removing ordering entirely.
  • Classify errors in the consumer: transient → nack(), permanent → ack() with a log entry. The absence of that distinction is the root cause of one bad record bringing down a queue.
  • Add the oldest_unacked_message_age alert at 10 minutes, not at 47.

Conclusion

AlpinaShop is no longer a monolith that phones everybody. In this lesson you have seen the coupling problem with concrete numbers — a purchase that waits 5.4 seconds for four systems the customer does not care about — and how it dissolves by publishing a fact instead of giving orders.

You have internalised the model: the topic is the type of fact, the subscription is where the messages actually live and each one receives its own independent copy, the message carries a body and filterable attributes, and the acknowledgement is the only thing that removes a message from the system. And above all you have seen the real guarantees, which are the ones that matter: at-least-once delivery, with duplicates that will happen even if your code is perfect; no order except with an ordering key; and exactly once as a reduction of the problem, not as its disappearance.

You have created pedidos-nuevos with 7 days of retention, the dead letter topic pedidos-nuevos-fallidos with its inspection subscription, and the sub-almacen, sub-facturacion — with exactly once, because duplicating an invoice is serious — and sub-analitica subscriptions, each with its deadline, its retry policy and its dead letter permissions, including the two bindings everybody forgets and without which the diversion silently fails to work.

You have published from Flask with batched sending, exponential backoff and a reused client, consciously deciding to publish a fat event with the lines inside and marking version_esquema and momento_evento as attributes. You have consumed in pull with flow control, clean shutdown and a distinction between transient and permanent errors; and in push towards an HTTPS endpoint with OIDC authentication, knowing that there the HTTP response code is the acknowledgement. You know the ack deadline, modifyAckDeadline and why the deadline is set to the 99th percentile of your processing.

You have understood why every serious consumer needs a dead letter — the poisoned message retried forever — and why the idempotency key must be the business identifier and not the message_id: a seek or a republication from the archive would issue duplicate invoices with the wrong choice. You have seen seek and snapshots as a safety net for reprocessing runs and risky deployments, the attribute filters that discard on the server side, and the BigQuery and Cloud Storage subscriptions, which ingest without a single line of code and which for AlpinaShop take the raw data while Dataflow keeps what genuinely needs transformation and windows.

You know when to enable ordering and at what price, that Pub/Sub Lite no longer exists and what to use instead, which metric to watch above all others — oldest_unacked_message_age — and that the cost, at AlpinaShop's volume, is cents. And you have finally kept the promise from 02-02: the alpinashop-catalogo bucket now gives notice when a new image appears, with the prefix filter, the path check and the idempotency by object generation that prevent the infinite loop that ruins the unwary.

Look at the map we now have. BigQuery stores and answers. Dataflow transforms in batch and in streaming. Dataproc computes the algorithmic things. Pub/Sub connects the systems without them knowing each other. That is a lot of machinery. And yet, there are two obvious gaps.

The first: the data that is not born in Google Cloud. AlpinaShop's warehouse ERP is a MySQL running on a server in the Sabadell offices. The carrier sends a monthly CSV by email. The backpack supplier offers a REST API with its catalogue. None of that publishes to Pub/Sub or writes Parquet into a bucket, and Lucía, who knows SQL but has never written Beam in her life, cannot depend on Dani for every new file.

In 04-05, Cloud Data Fusion, we will look at the tool designed for exactly that: visual, no-code data integration, with connectors for almost any source, an interactive cleaner called Wrangler where dates and nulls are fixed while you look at the data, field-level lineage that answers "where does this column come from?", and replication with change capture from that MySQL in Sabadell. And we will also look at its uncomfortable side — it charges per instance-hour and that shapes how it is used — with an honest table of when Data Fusion, when Dataflow, when Datastream and when a simple bq load will do.

Google Cloud Platform (GCP) Course

Module 1: Introduction to Google Cloud Platform

Module 2: Core GCP Services

Module 3: Networking and Security

Module 4: Data and Analytics

Module 5: Machine Learning and AI

Module 6: DevOps and Monitoring

Module 7: Advanced GCP Topics

Module 8: Final Project

© Copyright 2026. All rights reserved