Rate-Limited Task Queues with Cloud Tasks: Protecting Third-Party APIs at Scale

Using GCP Cloud Tasks for rate-limited, retry-safe integrations with third-party APIs, including queue configuration, dead-letter handling, and backpressure patterns.

#gcp#cloud-tasks#queues#async
Cover image for the article: Rate-Limited Task Queues with Cloud Tasks: Protecting Third-Party APIs at Scale

Third-party APIs have rate limits. Your system generates bursts. The naive solution — retry with exponential backoff — doesn't work when you're generating 10,000 webhook deliveries in a minute but the destination API allows 100 requests per second. You need a queue that respects external rate limits as a first-class constraint, not an afterthought.

Cloud Tasks provides exactly this: managed queues with configurable dispatch rates, automatic retries with backoff, and dead-letter handling — purpose-built for integrating with rate-limited external services.

The Architecture Problem

Our platform sends notifications to 14 different third-party services (Slack, email providers, webhook endpoints, CRM systems, analytics platforms). Each has different rate limits:

ServiceRate LimitOur Peak VolumeOverload Factor
Slack API50 req/min per workspace2,000/min40x
SendGrid600 req/sec10,000/sec (during campaigns)17x
Salesforce100 req/sec per org800/sec8x
Custom webhooksVaries (10-1000/sec)UnpredictableUnknown
HubSpot100 req/10sec500/10sec5x

Without rate-limited queuing, burst traffic triggers 429 responses, retries amplify the problem (retry storms), and notifications are lost or significantly delayed.

Rate Limit Problem Without Queuing

Cloud Tasks Architecture

Cloud Tasks sits between your event producers and the rate-limited consumers:

┌──────────────────┐     ┌─────────────────────┐     ┌──────────────────┐
│  Event Producers │────▶│    Cloud Tasks       │────▶│  Handler Service │
│  (Cloud Functions│     │  (Rate-Limited Queue) │     │  (Cloud Run)     │
│   Pub/Sub, etc.) │     │                       │     │                   │
└──────────────────┘     │  • 100 dispatches/sec │     │  Makes API calls  │
                         │  • Retry with backoff  │     │  to 3rd parties   │
                         │  • Dead letter queue   │     │                   │
                         └─────────────────────┘     └──────────────────┘

Queue Configuration with Terraform

Each third-party integration gets its own queue with tailored rate limits:

# terraform/cloud-tasks-queues.tf

locals {
  task_queues = {
    slack-notifications = {
      max_dispatches_per_second = 0.8    # 48/min (under 50/min limit)
      max_concurrent_dispatches = 5
      max_attempts              = 10
      min_backoff               = "5s"
      max_backoff               = "300s"
      max_doublings             = 5
    }
    sendgrid-emails = {
      max_dispatches_per_second = 500    # Under 600/sec limit
      max_concurrent_dispatches = 100
      max_attempts              = 5
      min_backoff               = "1s"
      max_backoff               = "60s"
      max_doublings             = 4
    }
    salesforce-sync = {
      max_dispatches_per_second = 80     # Under 100/sec limit
      max_concurrent_dispatches = 20
      max_attempts              = 8
      min_backoff               = "10s"
      max_backoff               = "600s"
      max_doublings             = 5
    }
    webhooks-generic = {
      max_dispatches_per_second = 50     # Conservative default
      max_concurrent_dispatches = 10
      max_attempts              = 15
      min_backoff               = "10s"
      max_backoff               = "3600s"
      max_doublings             = 8
    }
    hubspot-crm = {
      max_dispatches_per_second = 8      # Under 100/10sec limit
      max_concurrent_dispatches = 5
      max_attempts              = 10
      min_backoff               = "10s"
      max_backoff               = "600s"
      max_doublings             = 5
    }
  }
}

resource "google_cloud_tasks_queue" "integration_queues" {
  for_each = local.task_queues

  name     = each.key
  location = "us-central1"
  project  = var.project_id

  rate_limits {
    max_dispatches_per_second = each.value.max_dispatches_per_second
    max_concurrent_dispatches = each.value.max_concurrent_dispatches
  }

  retry_config {
    max_attempts       = each.value.max_attempts
    min_backoff        = each.value.min_backoff
    max_backoff        = each.value.max_backoff
    max_doublings      = each.value.max_doublings
    max_retry_duration = "86400s"  # 24 hours max retry window
  }

  stackdriver_logging_config {
    sampling_ratio = 0.1  # Log 10% of tasks for debugging
  }
}

Task Creation with Deduplication

Creating tasks with built-in deduplication prevents duplicate notifications during retries:

// services/task-creator/src/create-task.ts
import { CloudTasksClient, protos } from '@google-cloud/tasks';

const client = new CloudTasksClient();
const PROJECT = process.env.GCP_PROJECT!;
const LOCATION = 'us-central1';

interface NotificationPayload {
  userId: string;
  eventType: string;
  data: Record<string, unknown>;
  idempotencyKey: string;
  createdAt: string;
}

interface QueueConfig {
  queueName: string;
  handlerUrl: string;
  scheduleDelay?: number; // seconds
}

const QUEUE_ROUTING: Record<string, QueueConfig> = {
  'slack.message': {
    queueName: 'slack-notifications',
    handlerUrl: 'https://handlers-xyz.run.app/slack',
  },
  'email.transactional': {
    queueName: 'sendgrid-emails',
    handlerUrl: 'https://handlers-xyz.run.app/email',
  },
  'crm.contact_update': {
    queueName: 'salesforce-sync',
    handlerUrl: 'https://handlers-xyz.run.app/salesforce',
    scheduleDelay: 5, // Buffer for deduplication
  },
  'webhook.generic': {
    queueName: 'webhooks-generic',
    handlerUrl: 'https://handlers-xyz.run.app/webhook',
  },
};

export async function enqueueNotification(
  eventType: string,
  payload: NotificationPayload
): Promise<string> {
  const config = QUEUE_ROUTING[eventType];
  if (!config) {
    throw new Error(`No queue configured for event type: ${eventType}`);
  }

  const parent = client.queuePath(PROJECT, LOCATION, config.queueName);

  // Task name includes idempotency key for deduplication
  // Cloud Tasks rejects duplicate task names within 1 hour
  const taskName = `${parent}/tasks/${payload.idempotencyKey}`;

  const task: protos.google.cloud.tasks.v2.ITask = {
    name: taskName,
    httpRequest: {
      httpMethod: 'POST',
      url: config.handlerUrl,
      headers: {
        'Content-Type': 'application/json',
        'X-Idempotency-Key': payload.idempotencyKey,
        'X-Event-Type': eventType,
      },
      body: Buffer.from(JSON.stringify(payload)).toString('base64'),
      oidcToken: {
        serviceAccountEmail: `task-handler@${PROJECT}.iam.gserviceaccount.com`,
        audience: config.handlerUrl,
      },
    },
  };

  // Optional: schedule for future delivery
  if (config.scheduleDelay) {
    task.scheduleTime = {
      seconds: Math.floor(Date.now() / 1000) + config.scheduleDelay,
    };
  }

  try {
    const [response] = await client.createTask({ parent, task });
    return response.name!;
  } catch (error: any) {
    if (error.code === 6) {
      // ALREADY_EXISTS - task was already created (deduplication working)
      console.log(`Task already exists (deduplicated): ${taskName}`);
      return taskName;
    }
    throw error;
  }
}

Handler Service with 429 Awareness

The handler must communicate rate limit information back to Cloud Tasks via HTTP status codes:

// services/task-handler/src/handlers/slack.ts
import express from 'express';
import { WebClient } from '@slack/web-api';

const router = express.Router();

interface SlackNotification {
  userId: string;
  channel: string;
  message: string;
  blocks?: object[];
}

router.post('/slack', async (req, res) => {
  const payload: SlackNotification = JSON.parse(
    Buffer.from(req.body, 'base64').toString()
  );
  const idempotencyKey = req.headers['x-idempotency-key'] as string;

  // Check if already processed (idempotency)
  const alreadyProcessed = await checkIdempotencyStore(idempotencyKey);
  if (alreadyProcessed) {
    return res.status(200).json({ status: 'already_processed' });
  }

  try {
    const slack = new WebClient(await getSlackToken(payload.userId));

    await slack.chat.postMessage({
      channel: payload.channel,
      text: payload.message,
      blocks: payload.blocks,
    });

    // Mark as processed
    await markProcessed(idempotencyKey);
    return res.status(200).json({ status: 'delivered' });

  } catch (error: any) {
    if (error.code === 'slack_webapi_rate_limited') {
      // Slack returns Retry-After header
      const retryAfter = parseInt(error.headers?.['retry-after'] || '60');

      // Return 429 - Cloud Tasks will respect this and retry later
      res.set('Retry-After', retryAfter.toString());
      return res.status(429).json({
        status: 'rate_limited',
        retryAfter,
      });
    }

    if (error.code === 'slack_webapi_platform_error') {
      // Non-retryable errors (invalid channel, user not found)
      // Return 2xx to prevent retries, log for dead-letter processing
      await logToDeadLetter(idempotencyKey, payload, error);
      return res.status(200).json({ status: 'failed_permanent', error: error.message });
    }

    // Transient errors: return 500 for Cloud Tasks to retry
    console.error(`Transient error for ${idempotencyKey}:`, error);
    return res.status(500).json({ status: 'transient_error' });
  }
});

Dead Letter Queue Pattern

Tasks that exhaust all retries need a dead-letter path for investigation and manual replay:

// services/dead-letter-processor/src/index.ts
import { PubSub } from '@google-cloud/pubsub';
import { Firestore } from '@google-cloud/firestore';

const pubsub = new PubSub();
const db = new Firestore();

interface DeadLetterTask {
  originalQueue: string;
  taskName: string;
  payload: unknown;
  failureReason: string;
  attempts: number;
  firstAttempt: string;
  lastAttempt: string;
}

export async function processDeadLetter(message: any) {
  const task: DeadLetterTask = JSON.parse(
    Buffer.from(message.data, 'base64').toString()
  );

  // Store in Firestore for investigation dashboard
  await db.collection('dead-letter-tasks').doc(task.taskName).set({
    ...task,
    status: 'pending_review',
    createdAt: new Date(),
  });

  // Alert if dead letter rate exceeds threshold
  const recentDeadLetters = await db
    .collection('dead-letter-tasks')
    .where('originalQueue', '==', task.originalQueue)
    .where('createdAt', '>=', new Date(Date.now() - 3600000))
    .count()
    .get();

  if (recentDeadLetters.data().count > 50) {
    await publishAlert({
      severity: 'HIGH',
      queue: task.originalQueue,
      message: `Dead letter rate exceeded threshold: ${recentDeadLetters.data().count} tasks/hour`,
    });
  }

  message.ack();
}

Monitoring and Observability

MetricQueueAlert ThresholdAction
Queue depthAll> 10,000 tasksScale handler instances
Oldest task ageAll> 1 hourInvestigate handler errors
Dead letter rateAll> 1% of dispatchesCheck API availability
429 response ratePer queue> 10%Reduce dispatch rate
Success ratePer queue< 95%Investigate failures

Cloud Tasks Monitoring Dashboard

Performance Results

MetricBefore (Direct Calls)After (Cloud Tasks)Change
429 errors/day12,4000-100%
Failed notifications/day3408 (dead-lettered)-97%
Retry storm incidents/month40-100%
Notification delivery P9945s (retry delays)3s (queue dispatch)-93%
API ban incidents/quarter20-100%
Burst handling capacity100/sec (API limit)Unlimited (queue absorbs)Unlimited

Advanced Pattern: Dynamic Rate Adjustment

For APIs with variable rate limits (quota increases during off-peak hours):

// services/rate-adjuster/src/index.ts
import { CloudTasksClient } from '@google-cloud/tasks';

const client = new CloudTasksClient();

interface RateSchedule {
  queue: string;
  peakRate: number;
  offPeakRate: number;
  peakHoursUTC: [number, number]; // [start, end]
}

const RATE_SCHEDULES: RateSchedule[] = [
  {
    queue: 'salesforce-sync',
    peakRate: 40,        // Conservative during business hours
    offPeakRate: 95,     // Near-limit during off-hours
    peakHoursUTC: [13, 22], // 9am-6pm ET
  },
];

export async function adjustRates() {
  const currentHour = new Date().getUTCHours();

  for (const schedule of RATE_SCHEDULES) {
    const isPeak = currentHour >= schedule.peakHoursUTC[0] &#x26;&#x26;
                   currentHour &#x3C; schedule.peakHoursUTC[1];

    const targetRate = isPeak ? schedule.peakRate : schedule.offPeakRate;
    const queuePath = client.queuePath(PROJECT, LOCATION, schedule.queue);

    await client.updateQueue({
      queue: {
        name: queuePath,
        rateLimits: {
          maxDispatchesPerSecond: targetRate,
        },
      },
      updateMask: { paths: ['rate_limits.max_dispatches_per_second'] },
    });

    console.log(`Adjusted ${schedule.queue} to ${targetRate} dispatches/sec (${isPeak ? 'peak' : 'off-peak'})`);
  }
}

When NOT to Use Cloud Tasks

  • Fan-out/fan-in patterns: Use Pub/Sub with pull subscriptions instead
  • Real-time processing: Cloud Tasks adds minimum ~100ms latency; use direct calls if latency matters
  • Within-GCP service-to-service: Use Pub/Sub push subscriptions (similar but better integrated)
  • Exactly-once processing required: Cloud Tasks guarantees at-least-once; build idempotency yourself

Conclusion

Cloud Tasks transformed our third-party integration reliability from 97.6% to 99.98% delivery rate while completely eliminating rate limit violations and API bans. The key insight: rate limiting isn't about being polite to third-party APIs — it's about building a system that handles burst traffic gracefully without cascading failures.

Design your queues around external constraints, build idempotency into handlers, implement dead-letter processing from day one, and monitor queue depth as your primary health signal. The 429 error should never reach your application code — the queue should absorb it.

Comments

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