GCP Pub/Sub Exactly-Once Delivery: Message Guarantees and Deduplication in Production

Implementing exactly-once message processing with Cloud Pub/Sub using deduplication strategies, idempotency patterns, and dead letter queues.

#gcp#pub-sub#messaging#reliability
Cover image for the article: GCP Pub/Sub Exactly-Once Delivery: Message Guarantees and Deduplication in Production

"Exactly-once delivery" is one of distributed systems' most misunderstood promises. Cloud Pub/Sub offers exactly-once delivery semantics, but what it actually provides is exactly-once processing through a combination of server-side deduplication and client acknowledgment. After building an event-driven order processing system handling 2M+ messages per day, here's what that distinction means in practice.

The Problem: Duplicate Processing in Event-Driven Systems

Our order fulfillment pipeline processes payment confirmations, inventory updates, and shipping notifications through Pub/Sub. A single duplicate payment processing event means charging a customer twice. A duplicate inventory decrement means phantom stock shortages.

Before implementing proper deduplication, we observed duplicate processing rates of 0.1-0.3% under normal conditions and up to 2.4% during subscriber restarts or deployment rollouts. At 2M messages per day, that's 2,000-48,000 duplicate operations daily.

Pub/Sub's Delivery Guarantees Explained

Pub/Sub provides three delivery modes:

ModeGuaranteeWhen to Use
At-least-onceMessages delivered 1+ timesDefault, works for idempotent operations
Exactly-onceServer deduplicates before deliveryFinancial transactions, inventory
At-most-onceBest-effort, no redeliveryMetrics, logs, non-critical events

Pub/Sub Message Delivery Flow

Exactly-once delivery in Pub/Sub works through a two-phase mechanism:

  1. Server-side deduplication: Pub/Sub assigns unique message IDs and deduplicates publishes with the same ordering key within a 10-minute window.
  2. Client-side acknowledgment: Subscribers acknowledge messages with exactly-once semantics enabled, and Pub/Sub tracks acknowledgment state.

Enabling Exactly-Once Delivery

# Create subscription with exactly-once delivery
gcloud pubsub subscriptions create order-processing \
  --topic=payment-events \
  --enable-exactly-once-delivery \
  --ack-deadline=60 \
  --dead-letter-topic=payment-events-dlq \
  --max-delivery-attempts=5 \
  --message-retention-duration=7d

On the publisher side, use ordering keys to enable deduplication:

package main

import (
    "context"
    "fmt"
    "cloud.google.com/go/pubsub"
)

func PublishOrderEvent(ctx context.Context, client *pubsub.Client, order *Order) error {
    topic := client.Topic("payment-events")
    
    // Enable message ordering for this topic
    topic.EnableMessageOrdering = true
    
    // Ordering key determines deduplication scope
    // Messages with the same ordering key are deduplicated and delivered in order
    result := topic.Publish(ctx, &pubsub.Message{
        Data: order.ToJSON(),
        Attributes: map[string]string{
            "event_type":    "payment.confirmed",
            "idempotency_key": order.IdempotencyKey,
            "version":       "2",
        },
        OrderingKey: order.CustomerID, // All events for a customer are ordered
    })
    
    id, err := result.Get(ctx)
    if err != nil {
        return fmt.Errorf("publishing order event: %w", err)
    }
    
    fmt.Printf("Published message ID: %s\n", id)
    return nil
}

Application-Level Deduplication

Server-side deduplication has a 10-minute window. For longer processing times or cross-system idempotency, implement application-level deduplication:

type DeduplicatingProcessor struct {
    store  DeduplicationStore
    handler MessageHandler
}

func (p *DeduplicatingProcessor) Process(ctx context.Context, msg *pubsub.Message) error {
    // Extract idempotency key from message attributes
    idempotencyKey := msg.Attributes["idempotency_key"]
    if idempotencyKey == "" {
        idempotencyKey = msg.ID // Fallback to Pub/Sub message ID
    }
    
    // Check if we've already processed this message
    processed, err := p.store.IsProcessed(ctx, idempotencyKey)
    if err != nil {
        return fmt.Errorf("checking deduplication store: %w", err)
    }
    
    if processed {
        // Already processed — acknowledge without reprocessing
        msg.Ack()
        return nil
    }
    
    // Process the message
    if err := p.handler.Handle(ctx, msg); err != nil {
        // Processing failed — nack for retry
        msg.Nack()
        return fmt.Errorf("handling message: %w", err)
    }
    
    // Mark as processed BEFORE acknowledging
    // Use a TTL of 7 days (matching message retention)
    if err := p.store.MarkProcessed(ctx, idempotencyKey, 7*24*time.Hour); err != nil {
        // If we can't mark as processed, nack the message
        // This may cause a reprocessing, but that's safer than missing the dedup record
        msg.Nack()
        return fmt.Errorf("marking processed: %w", err)
    }
    
    msg.Ack()
    return nil
}

For the deduplication store, we use Cloud Memorystore (Redis) for speed with Spanner as a durable fallback:

type RedisDeduplicationStore struct {
    client *redis.Client
}

func (s *RedisDeduplicationStore) IsProcessed(ctx context.Context, key string) (bool, error) {
    exists, err := s.client.Exists(ctx, "dedup:"+key).Result()
    if err != nil {
        return false, err
    }
    return exists > 0, nil
}

func (s *RedisDeduplicationStore) MarkProcessed(ctx context.Context, key string, ttl time.Duration) error {
    return s.client.Set(ctx, "dedup:"+key, "1", ttl).Err()
}

Dead Letter Queue Configuration

Messages that fail after max delivery attempts need a safety net:

func SetupDLQMonitoring(ctx context.Context, client *pubsub.Client) {
    dlqSub := client.Subscription("payment-events-dlq-sub")
    
    err := dlqSub.Receive(ctx, func(ctx context.Context, msg *pubsub.Message) {
        // Log detailed information for investigation
        log.Error("Message moved to DLQ",
            "message_id", msg.ID,
            "publish_time", msg.PublishTime,
            "delivery_attempt", msg.DeliveryAttempt,
            "attributes", msg.Attributes,
            "data_preview", string(msg.Data[:min(200, len(msg.Data))]),
        )
        
        // Alert on-call if DLQ rate exceeds threshold
        metrics.IncrementCounter("dlq_messages_total", map[string]string{
            "topic": "payment-events",
            "event_type": msg.Attributes["event_type"],
        })
        
        msg.Ack()
    })
    
    if err != nil {
        log.Fatal("DLQ receiver failed", "error", err)
    }
}

Performance Impact of Exactly-Once

Enabling exactly-once delivery has measurable performance implications:

MetricAt-Least-OnceExactly-OnceImpact
Publish latency (p50)12ms14ms+17%
Publish latency (p99)45ms68ms+51%
Subscribe latency (p50)8ms11ms+38%
Ack latency (p50)5ms18ms+260%
Max throughput100K msg/s30K msg/s-70%
Duplicate rate0.1-0.3%<0.001%~eliminated

Pub/Sub Throughput Comparison

The throughput reduction is significant. For high-volume, latency-sensitive topics where occasional duplicates are acceptable (metrics, logging), use at-least-once with application-level idempotency.

Ordering Guarantees and Parallelism

Ordering keys guarantee in-order delivery for messages with the same key but limit parallelism:

# Monitor subscription backlog by ordering key
gcloud pubsub subscriptions describe order-processing \
  --format="json" | jq '.messageRetentionDuration'

# Check for hot ordering keys causing backlog
gcloud monitoring read \
  "pubsub.googleapis.com/subscription/num_unacked_messages_by_region" \
  --filter='resource.labels.subscription_id="order-processing"'

Our solution: use customer IDs as ordering keys (guaranteeing per-customer ordering) while allowing parallelism across customers. With 500K active customers, this gives us 500K parallel streams.

Key Takeaways

  1. Exactly-once in Pub/Sub means exactly-once processing, not delivery. The server deduplicates, but your subscriber logic must still be idempotent as a defense-in-depth measure.
  2. The 10-minute deduplication window is sufficient for most workloads. For longer processing, add application-level deduplication with Redis/Memorystore.
  3. Exactly-once has real performance costs. Ack latency increases 260% and max throughput drops 70%. Use it only where duplicates cause business-level harm.
  4. Dead letter queues are non-negotiable. Every subscription should have a DLQ with monitoring and alerting.
  5. Ordering keys are your parallelism knob. Too few keys (e.g., one global key) serializes everything. Too many keys (e.g., message ID) provides no ordering. Entity IDs are the sweet spot.

The safest architecture combines Pub/Sub's exactly-once delivery with application-level idempotency. Belt and suspenders isn't over-engineering when customer money is on the line.

Comments

    No comments yet. Be the first to share your thoughts.