Building Event-Driven Microservices with Cloud Run and Pub/Sub
Event-driven architecture decouples services and makes your system more resilient to failure. This guide shows how to build production event-driven microservices on GCP using Cloud Run and Pub/Sub, including dead-letter handling, exactly-once processing, and end-to-end tracing.
Microservices that call each other synchronously over HTTP introduce a fragility problem: if any service in a chain is slow or unavailable, the calling service blocks and eventually times out. In a system with 20 services where each has 99.9% availability, the probability that all 20 are available simultaneously is only 98% — meaning roughly 14 hours of downtime per month just from normally operating services failing at different times.
Event-driven architecture addresses this by making services communicate asynchronously through a message bus. Publisher services emit events and move on; subscriber services consume events when they're ready. Services don't know about or depend on each other directly — they only know about events. This loose coupling is what makes systems resilient, independently deployable, and easier to evolve.
GCP's Cloud Run and Pub/Sub combination makes this architecture practical for teams that can't afford the operational overhead of running Kafka, RabbitMQ, or another message broker. Pub/Sub handles message durability, ordering, and delivery guarantees. Cloud Run handles containerized service execution at any scale. Together they give you a production-grade event-driven system with minimal operational footprint.
Architecture Overview
A typical event-driven system on Cloud Run and Pub/Sub looks like this:
Order API (Cloud Run)
|
| publishes: order.created
v
Pub/Sub Topic: orders
|
+---> Subscription: order-processor --> Order Processor Service (Cloud Run)
| publishes: order.fulfilled, order.failed
|
+---> Subscription: analytics --> Analytics Service (Cloud Run)
|
+---> Dead Letter Topic: orders-dead-letter
|
v
Dead Letter Subscription --> Alert + retry logic
Events flow through topics. Services subscribe to the topics they care about and process events independently. The failure of the analytics service doesn't affect order processing. Order processing retries automatically if it fails.
Setting Up the Pub/Sub Infrastructure
# Create the orders topic
gcloud pubsub topics create orders --project=my-project --message-retention-duration=7d
# Create a dead-letter topic for failed messages
gcloud pubsub topics create orders-dead-letter --project=my-project
# Create a push subscription that delivers to Cloud Run
gcloud pubsub subscriptions create order-processor-sub --topic=orders --push-endpoint=https://order-processor-xxxxx-ew.a.run.app/events --push-auth-service-account=pubsub-invoker@my-project.iam.gserviceaccount.com --ack-deadline=60 --dead-letter-topic=projects/my-project/topics/orders-dead-letter --max-delivery-attempts=5 --min-retry-delay=10s --max-retry-delay=600s
The --push-auth-service-account parameter makes Pub/Sub authenticate to Cloud Run using OIDC — your Cloud Run service only accepts requests from this specific service account, preventing unauthorized event injection.
Grant the Pub/Sub service account permission to invoke your Cloud Run service:
# Create service account for Pub/Sub to use when calling Cloud Run
gcloud iam service-accounts create pubsub-invoker --display-name="Pub/Sub Cloud Run Invoker"
# Grant it permission to invoke Cloud Run
gcloud run services add-iam-policy-binding order-processor --region=europe-west4 --member=serviceAccount:pubsub-invoker@my-project.iam.gserviceaccount.com --role=roles/run.invoker
# Grant Pub/Sub service account token creation rights
gcloud projects add-iam-policy-binding my-project --member=serviceAccount:service-PROJECT_NUMBER@gcp-sa-pubsub.iam.gserviceaccount.com --role=roles/iam.serviceAccountTokenCreator
Implementing the Event Publisher
The Order API publishes events to Pub/Sub after every state change. Use structured event envelopes that all services agree on:
# events.py — shared event schema
import json
import uuid
from dataclasses import dataclass, asdict
from datetime import datetime, UTC
from typing import Any
from google.cloud import pubsub_v1
@dataclass
class Event:
event_id: str
event_type: str
source: str
subject: str
timestamp: str
data: dict
@classmethod
def create(cls, event_type: str, source: str, subject: str, data: dict) -> "Event":
return cls(
event_id=str(uuid.uuid4()),
event_type=event_type,
source=source,
subject=subject,
timestamp=datetime.now(UTC).isoformat(),
data=data,
)
class EventPublisher:
def __init__(self, project_id: str):
self.publisher = pubsub_v1.PublisherClient()
self.project_id = project_id
def publish(self, topic_id: str, event: Event) -> str:
topic_path = self.publisher.topic_path(self.project_id, topic_id)
# Serialize event to JSON
data = json.dumps(asdict(event)).encode("utf-8")
# Include event_type as an attribute for Pub/Sub filter subscriptions
future = self.publisher.publish(
topic_path,
data,
event_type=event.event_type,
source=event.source,
)
return future.result() # Block until published (or raise on error)
# order_api.py — publishing events from the Order API
from flask import Flask, request, jsonify
from events import EventPublisher, Event
app = Flask(__name__)
publisher = EventPublisher(project_id="my-project")
@app.route("/orders", methods=["POST"])
def create_order():
order_data = request.get_json()
# Persist the order
order = save_order(order_data)
# Publish event
event = Event.create(
event_type="order.created",
source="order-api",
subject=f"orders/{order.id}",
data={
"order_id": order.id,
"customer_id": order.customer_id,
"items": order.items,
"total_amount": order.total_amount,
},
)
publisher.publish("orders", event)
return jsonify({"order_id": order.id, "status": "created"}), 201
Implementing the Event Consumer
Cloud Run receives Pub/Sub push messages as HTTP POST requests. The message arrives in the Pub/Sub push format:
# order_processor.py
import base64
import json
import logging
from flask import Flask, request, jsonify, abort
from events import Event
app = Flask(__name__)
logger = logging.getLogger(__name__)
@app.route("/events", methods=["POST"])
def handle_event():
# Validate it's a Pub/Sub message
envelope = request.get_json(silent=True)
if not envelope or "message" not in envelope:
abort(400, "Invalid Pub/Sub message format")
# Decode the message
message = envelope["message"]
try:
raw_data = base64.b64decode(message["data"]).decode("utf-8")
event_data = json.loads(raw_data)
except Exception as e:
logger.error(f"Failed to decode message: {e}")
# Return 200 to ack the message — we can't process it and don't want retries
# Log to dead letter manually or rely on max delivery attempts
return jsonify({"status": "discarded", "reason": "malformed"}), 200
event = Event(**event_data)
try:
process_event(event)
# Return 200 to acknowledge successful processing
return jsonify({"status": "processed"}), 200
except RetryableError as e:
logger.warning(f"Retryable error: {e}. Returning 500 to trigger retry.")
# Return non-2xx to cause Pub/Sub to redeliver
return jsonify({"status": "retry"}), 500
except Exception as e:
logger.error(f"Non-retryable error: {e}. Acking to prevent infinite retry.")
# Return 200 to avoid infinite retry loop for non-retryable errors
# The event will be handled via the dead-letter topic
return jsonify({"status": "failed", "error": str(e)}), 200
def process_event(event: Event):
if event.event_type == "order.created":
handle_order_created(event.data)
elif event.event_type == "order.cancelled":
handle_order_cancelled(event.data)
else:
logger.warning(f"Unknown event type: {event.event_type}")
Idempotent Event Handling
Pub/Sub guarantees at-least-once delivery. Under network conditions or retries, your service may receive the same event more than once. Every event handler must be idempotent — processing the same event twice produces the same result as processing it once.
import redis
from contextlib import contextmanager
redis_client = redis.Redis.from_url("redis://...")
@contextmanager
def idempotency_guard(event_id: str, ttl_seconds: int = 86400):
"""Context manager that ensures an event is processed at most once."""
key = f"processed_event:{event_id}"
# Try to set the key with NX (only if not exists)
was_set = redis_client.set(key, "1", nx=True, ex=ttl_seconds)
if not was_set:
# Event was already processed
yield False
return
try:
yield True
except Exception:
# Processing failed — remove the key so it can be retried
redis_client.delete(key)
raise
def handle_order_created(data: dict):
order_id = data["order_id"]
with idempotency_guard(f"order.created:{order_id}") as should_process:
if not should_process:
logger.info(f"Order {order_id} already processed, skipping")
return
# Actual processing logic — safe to run knowing it's the first time
reserve_inventory(data["items"])
charge_payment(data["customer_id"], data["total_amount"])
send_confirmation_email(data["customer_id"], order_id)
Using Eventarc for Tighter Cloud Run Integration
Eventarc is GCP's managed eventing service that routes events from 90+ GCP sources to Cloud Run without manual Pub/Sub subscription management:
# Trigger Cloud Run when a file is uploaded to Cloud Storage
gcloud eventarc triggers create process-uploaded-file --destination-run-service=file-processor --destination-run-region=europe-west4 --event-filters="type=google.cloud.storage.object.v1.finalized" --event-filters="bucket=my-upload-bucket" --service-account=eventarc-invoker@my-project.iam.gserviceaccount.com --location=europe-west4
# Trigger Cloud Run from Pub/Sub (Eventarc wrapping Pub/Sub)
gcloud eventarc triggers create orders-processor --destination-run-service=order-processor --destination-run-region=europe-west4 --event-filters="type=google.cloud.pubsub.topic.v1.messagePublished" --transport-topic=orders --service-account=eventarc-invoker@my-project.iam.gserviceaccount.com --location=europe-west4
With Eventarc, the event arrives at your Cloud Run service in CloudEvents format:
from cloudevents.http import from_http
from flask import Flask, request
app = Flask(__name__)
@app.route("/", methods=["POST"])
def handle_cloud_event():
event = from_http(request.headers, request.get_data())
if event.type == "google.cloud.storage.object.v1.finalized":
bucket = event.data["bucket"]
name = event.data["name"]
process_uploaded_file(bucket, name)
return "OK", 200
Dead Letter Processing and Monitoring
Events that fail after max retry attempts land in the dead-letter topic. Consume these with a separate monitoring service:
# dead_letter_processor.py
@app.route("/dead-letter", methods=["POST"])
def handle_dead_letter():
envelope = request.get_json()
message = envelope["message"]
attributes = message.get("attributes", {})
dead_letter_event = {
"original_subscription": attributes.get("CloudPubSubDeadLetterSourceSubscription"),
"delivery_attempt": attributes.get("CloudPubSubDeadLetterSourceDeliveryAttempt"),
"message_id": message["messageId"],
"data": base64.b64decode(message["data"]).decode("utf-8"),
}
# Alert the on-call team
send_alert(f"Dead letter received: {dead_letter_event['original_subscription']}")
# Store for manual investigation
save_to_bigquery("dead_letter_events", dead_letter_event)
return "OK", 200
Monitor dead letter queue depth:
# Alert if dead letter queue grows beyond threshold
gcloud alpha monitoring policies create --display-name="Dead Letter Queue Backlog" --condition-filter='resource.type="pubsub_subscription" AND resource.label.subscription_id="orders-dead-letter-sub" AND metric.type="pubsub.googleapis.com/subscription/num_undelivered_messages"' --condition-threshold-value=10 --condition-threshold-comparison=COMPARISON_GT --notification-channels=$CHANNEL_ID
Distributed Tracing Across Services
In an event-driven system, a single user action triggers a chain of events across many services. Cloud Trace lets you follow a request through the entire chain:
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.exporter.cloud_trace import CloudTraceSpanExporter
from opentelemetry.propagate import inject, extract
from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator
# Set up tracer
provider = TracerProvider()
provider.add_span_processor(BatchSpanProcessor(CloudTraceSpanExporter()))
trace.set_tracer_provider(provider)
tracer = trace.get_tracer(__name__)
# Publisher: inject trace context into event
def publish_with_trace(event: Event, topic_id: str):
carrier = {}
inject(carrier) # Adds traceparent header
publisher.publish(
topic_path,
data=json.dumps(asdict(event)).encode("utf-8"),
**carrier, # Pass trace context as Pub/Sub attributes
)
# Consumer: extract trace context and continue the trace
@app.route("/events", methods=["POST"])
def handle_event():
message = request.get_json()["message"]
attributes = message.get("attributes", {})
# Extract trace context from Pub/Sub attributes
ctx = extract(attributes)
with tracer.start_as_current_span("handle_event", context=ctx):
process_event(message)
return "OK", 200
This connects spans across Cloud Run services into a single trace, so you can see the full journey of an order from creation through fulfillment in Cloud Trace.
Deployment with Cloud Build
Deploy the entire event-driven system with a Cloud Build pipeline:
# cloudbuild.yaml
steps:
# Build and push all service images
- name: gcr.io/cloud-builders/docker
args: [build, -t, europe-west4-docker.pkg.dev/$PROJECT_ID/services/order-api:$COMMIT_SHA, ./order-api]
- name: gcr.io/cloud-builders/docker
args: [push, europe-west4-docker.pkg.dev/$PROJECT_ID/services/order-api:$COMMIT_SHA]
- name: gcr.io/cloud-builders/docker
args: [build, -t, europe-west4-docker.pkg.dev/$PROJECT_ID/services/order-processor:$COMMIT_SHA, ./order-processor]
- name: gcr.io/cloud-builders/docker
args: [push, europe-west4-docker.pkg.dev/$PROJECT_ID/services/order-processor:$COMMIT_SHA]
# Deploy services to Cloud Run
- name: gcr.io/google.com/cloudsdktool/cloud-sdk
args:
- gcloud
- run
- deploy
- order-api
- --image=europe-west4-docker.pkg.dev/$PROJECT_ID/services/order-api:$COMMIT_SHA
- --region=europe-west4
- --no-allow-unauthenticated
- name: gcr.io/google.com/cloudsdktool/cloud-sdk
args:
- gcloud
- run
- deploy
- order-processor
- --image=europe-west4-docker.pkg.dev/$PROJECT_ID/services/order-processor:$COMMIT_SHA
- --region=europe-west4
- --no-allow-unauthenticated
The combination of Cloud Run's automatic scaling, Pub/Sub's durability and retry semantics, and dead-letter queues for handling failures gives you a resilient event-driven system that handles traffic spikes, service failures, and deployment updates gracefully.
For performance tuning of Cloud Run services, see our Cloud Run concurrency and cold start guide. For securing these services, see our GCP serverless security guide.