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.