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
- The coupling problem, before and after
- The publish-subscribe model
- Pub/Sub's real guarantees
- Creating the
pedidos-nuevostopic and its subscriptions - Publishing from the Flask application
- Pull consumption with the asynchronous client
- Push consumption to an HTTPS endpoint with OIDC
- Acknowledgements, deadlines and retries
- Dead letter topics
- Idempotency: the consumer must tolerate duplicates
- Retention,
seekand snapshots - Subscription filters by attribute
- BigQuery and Cloud Storage subscriptions
- Ordering by key and its implications
- Pub/Sub Lite and other alternatives
- Monitoring and cost
- The
alpinashop-catalogobucket notifications
- 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.
- 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, optionalorderingKey.
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.
- 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).
- Creating the
pedidos-nuevos topic and its subscriptions
pedidos-nuevos topic and its subscriptionsgcloud 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=600The 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=analiticaNotes on the differences, which are deliberate:
sub-facturacionwith 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-analiticawith no dead letter. The Dataflow pipeline already has its own quarantine mechanism (04-02): malformed messages go topedidos_streaming_errores. Duplicating the mechanism would complicate diagnosis.- Different
ack deadlinevalues. 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)"
- 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 futureAnd 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"}), 201Four 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
- 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
- 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=5The 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 "", 204The 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.
- 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:
The retries follow the policy you configured:
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.
- 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=5When 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:
- 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".
max-delivery-attemptsbetween 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.- An alert on the number of messages in the dead letter. A dead letter nobody looks at is a black hole with extra steps.
- 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.
- 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 TrueAnd 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.
- Retention,
seek and snapshots
seek and snapshotsMessages 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,
seekback 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.
seekto 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-v2It 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.
- 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
paisandorigenas 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.
- 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=jsonIt 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.
- 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-orderingThe 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:
- Lower throughput. The messages for one key are processed sequentially. With
pedido_idas the key there is no problem: there are thousands of distinct orders and parallelism is preserved. Withpaisas the key, every order from Spain would go in series. - Blocking on a poisoned message. If a message with key
PED-2026-0042fails 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. - 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.
- 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.
- 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
- The
alpinashop-catalogo bucket notifications
alpinashop-catalogo bucket notificationsBack 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=5Available 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 withattributes.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 tablealpinashop-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:
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_idas the key: redelivery preserves the originalmessage_id, so in this specific case the 200 good ones would indeed be detected. But if instead ofseeksomebody recovers and republishes the events from the Cloud Storage archive (section 13), themessage_idvalues will be new and 200 duplicate invoices would be issued. That is the failure. - With
pedido_idas 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:
- The service is alive and responding. It is not an outage: it is a logical jam.
- 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. - 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 thatnack()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. - 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.publisheron the dead letter topic norroles/pubsub.subscriberon 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-topicand 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-almacenThree-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=5As 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 anUPDATE ... 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_agealert 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
- What is Google Cloud Platform?
- Setting Up Your GCP Account
- A Tour of the GCP Console
- Projects, Resource Hierarchy and Billing
- Regions, Zones and the Shared Responsibility Model
- Cloud Shell and the gcloud CLI
Module 2: Core GCP Services
- Compute Engine: Virtual Machines on Google Cloud
- Cloud Storage: Object Storage
- Cloud SQL: Managed Relational Databases
- App Engine: Platform as a Service
- Google Kubernetes Engine (GKE)
- NoSQL Databases: Firestore, Bigtable and Spanner
- How to Choose the Right Compute Service
Module 3: Networking and Security
- VPC Networks
- Cloud Load Balancing
- Cloud CDN
- Identity and Access Management (IAM)
- Cloud Armor
- Secrets and Encryption: Secret Manager and Cloud KMS
- Cloud DNS, TLS Certificates and Publishing Services Securely
Module 4: Data and Analytics
- BigQuery: The Analytical Data Warehouse
- Cloud Dataflow: Batch and Streaming Data Processing
- Cloud Dataproc: Managed Spark and Hadoop
- Cloud Pub/Sub: Asynchronous Messaging
- Cloud Data Fusion: Code-Free Data Integration
- Orchestrating Pipelines with Cloud Composer and Workflows
- Data Governance and Dashboards with Dataplex and Looker Studio
Module 5: Machine Learning and AI
- Vertex AI: The Machine Learning Platform on GCP
- AutoML: Custom Models Without Writing Code
- TensorFlow on GCP: Training and Serving Models
- Natural Language API
- Vision API
- Generative AI on Vertex AI: Gemini Models and Embeddings
- MLOps: From Model to Product with Vertex AI Pipelines
Module 6: DevOps and Monitoring
- Cloud Build: Continuous Integration on GCP
- Cloud Source Repositories and Source Code Management
- Cloud Functions: Serverless Functions
- Cloud Monitoring (formerly Stackdriver): Metrics, Dashboards and Alerts
- Cloud Deployment Manager and Native Infrastructure as Code
- Cloud Logging and Cloud Trace: Logs, Traces and Diagnostics
- Terraform on GCP: Infrastructure as Code in Practice
Module 7: Advanced GCP Topics
- Hybrid and Multicloud with Anthos
- Serverless Computing with Cloud Run
- Advanced Networking: Shared VPC, Peering and Hybrid Connectivity
- Security Best Practices
- Cost Management and Optimization
- Reliability: SLOs, High Availability and Disaster Recovery
- Governance at Scale: Organization, Policies and Auditing
