How We Handle 50K OpenAI Requests/Minute Without Getting Rate Limited
Back to all articles
AI Engineering
15 min read8 min read

How We Handle 50K OpenAI Requests/Minute Without Getting Rate Limited

Real infrastructure patterns for high-volume LLM applications: queue management, intelligent retries, request batching, and graceful degradation.

Debasish Maji
Debasish Maji
AI Engineering Lead
January 18, 2026
Rate LimitingInfrastructureScalingProduction

The Problem at Scale

Our AI writing assistant hit product-market fit. Great news, except:

  • Day 1: 500 requests/minute, everything works
  • Day 30: 5,000 requests/minute, occasional 429s
  • Day 90: 50,000 requests/minute, 30% of requests failing

OpenAI's rate limits aren't just about requests per minute. They have:

  • Requests per minute (RPM) limits
  • Tokens per minute (TPM) limits
  • Requests per day limits
  • Different limits per model
  • Limits that change based on your tier

A single retry strategy won't work.

•••

The Architecture That Actually Works

Diagram
flowchart TB A[User Request] --> B[Request Queue<br/>Priority-based] B --> C[Rate Limiter<br/>Token Bucket] B --> D[Smart Router] B --> E[Request Batcher] C --> F[Provider Pool] D --> F E --> F F --> G[OpenAI] F --> H[Azure OpenAI] F --> I[Anthropic] G --> J[Response Handler] H --> J I --> J J --> K[Return to User] J -.->|429/500| L[Retry Queue] L -.-> B style A fill:#14b8a6,color:#fff style F fill:#0ea5e9,color:#fff style K fill:#22c55e,color:#fff style L fill:#f59e0b,color:#fff
•••

Component 1: Intelligent Request Queue

Not all requests are equal. A user waiting for a response needs priority over a background task.

Python
import asyncio
from dataclasses import dataclass, field
from typing import Optional
from enum import IntEnum
import heapq
import time

class Priority(IntEnum):
    CRITICAL = 0      # User is actively waiting, timeout imminent
    HIGH = 1          # User is waiting, interactive request
    NORMAL = 2        # User-initiated but can tolerate delay
    LOW = 3           # Background job, batch processing
    BULK = 4          # Analytics, non-time-sensitive

@dataclass(order=True)
class QueuedRequest:
    priority: Priority
    timestamp: float = field(compare=False)
    request_id: str = field(compare=False)
    payload: dict = field(compare=False)
    callback: asyncio.Future = field(compare=False)
    deadline: Optional[float] = field(compare=False, default=None)
    estimated_tokens: int = field(compare=False, default=0)

class PriorityRequestQueue:
    def __init__(self, max_size: int = 10000):
        self.queue: list = []
        self.max_size = max_size
        self.lock = asyncio.Lock()
        self.not_empty = asyncio.Condition()
        
    async def enqueue(self, request: QueuedRequest) -> bool:
        async with self.lock:
            if len(self.queue) >= self.max_size:
                # Shed load: reject low priority requests when full
                if request.priority >= Priority.LOW:
                    return False
                # Or: remove lowest priority item
                self._remove_lowest_priority()
            
            heapq.heappush(self.queue, request)
            
        async with self.not_empty:
            self.not_empty.notify()
        return True
    
    async def dequeue(self) -> QueuedRequest:
        async with self.not_empty:
            while not self.queue:
                await self.not_empty.wait()
        
        async with self.lock:
            # Remove expired requests
            self._cleanup_expired()
            return heapq.heappop(self.queue)
    
    def _cleanup_expired(self):
        now = time.time()
        self.queue = [
            r for r in self.queue 
            if r.deadline is None or r.deadline > now
        ]
        heapq.heapify(self.queue)
•••

Component 2: Token-Aware Rate Limiter

OpenAI limits both requests AND tokens. You need to track both.

Python
import asyncio
from collections import deque
import time

class DualTokenBucket:
    """Rate limiter that tracks both requests and tokens."""
    
    def __init__(
        self,
        requests_per_minute: int,
        tokens_per_minute: int,
        burst_multiplier: float = 1.5
    ):
        self.rpm_limit = requests_per_minute
        self.tpm_limit = tokens_per_minute
        self.burst_multiplier = burst_multiplier
        
        # Sliding window for accurate tracking
        self.request_timestamps: deque = deque()
        self.token_usage: deque = deque()  # (timestamp, token_count)
        
        self.lock = asyncio.Lock()
    
    async def acquire(self, estimated_tokens: int, timeout: float = 30.0) -> bool:
        """Try to acquire capacity for a request. Returns False if timeout exceeded."""
        deadline = time.time() + timeout
        
        while time.time() < deadline:
            async with self.lock:
                self._cleanup_old_entries()
                
                current_rpm = len(self.request_timestamps)
                current_tpm = sum(t[1] for t in self.token_usage)
                
                # Check if we have capacity
                rpm_ok = current_rpm < self.rpm_limit * self.burst_multiplier
                tpm_ok = current_tpm + estimated_tokens < self.tpm_limit * self.burst_multiplier
                
                if rpm_ok and tpm_ok:
                    now = time.time()
                    self.request_timestamps.append(now)
                    self.token_usage.append((now, estimated_tokens))
                    return True
                
                # Calculate wait time
                wait_time = self._calculate_wait_time(estimated_tokens)
            
            # Wait outside the lock
            await asyncio.sleep(min(wait_time, deadline - time.time()))
        
        return False
    
    def record_actual_usage(self, actual_tokens: int, estimated_tokens: int):
        """Update token count with actual usage after response."""
        async def _update():
            async with self.lock:
                # Find and update the most recent estimate
                for i in range(len(self.token_usage) - 1, -1, -1):
                    if self.token_usage[i][1] == estimated_tokens:
                        timestamp = self.token_usage[i][0]
                        self.token_usage[i] = (timestamp, actual_tokens)
                        break
        
        asyncio.create_task(_update())
    
    def _cleanup_old_entries(self):
        cutoff = time.time() - 60  # 1 minute window
        while self.request_timestamps and self.request_timestamps[0] < cutoff:
            self.request_timestamps.popleft()
        while self.token_usage and self.token_usage[0][0] < cutoff:
            self.token_usage.popleft()
    
    def _calculate_wait_time(self, needed_tokens: int) -> float:
        if not self.request_timestamps:
            return 0
        
        # When will enough capacity free up?
        oldest_request = self.request_timestamps[0]
        return max(0, oldest_request + 60 - time.time() + 0.1)
•••

Component 3: Multi-Provider Router

Don't put all your eggs in one basket. We use multiple providers with automatic failover.

Python
from dataclasses import dataclass
from typing import Dict, List, Optional
import random

@dataclass
class ProviderConfig:
    name: str
    endpoint: str
    api_key: str
    rate_limiter: DualTokenBucket
    models: List[str]
    cost_multiplier: float = 1.0
    latency_p50_ms: float = 500
    is_healthy: bool = True
    consecutive_failures: int = 0

class MultiProviderRouter:
    def __init__(self, providers: List[ProviderConfig]):
        self.providers = {p.name: p for p in providers}
        self.model_to_providers: Dict[str, List[str]] = {}
        
        for provider in providers:
            for model in provider.models:
                if model not in self.model_to_providers:
                    self.model_to_providers[model] = []
                self.model_to_providers[model].append(provider.name)
    
    async def route(self, model: str, estimated_tokens: int, 
                    priority: Priority) -> Optional[ProviderConfig]:
        """Select the best available provider for this request."""
        
        available_providers = self.model_to_providers.get(model, [])
        if not available_providers:
            raise ValueError(f"No providers available for model: {model}")
        
        # Filter to healthy providers
        healthy = [
            self.providers[name] for name in available_providers
            if self.providers[name].is_healthy
        ]
        
        if not healthy:
            # All providers unhealthy, try the one with fewest consecutive failures
            healthy = sorted(
                [self.providers[name] for name in available_providers],
                key=lambda p: p.consecutive_failures
            )[:1]
        
        # Try providers in order of preference
        for provider in self._rank_providers(healthy, priority):
            can_acquire = await provider.rate_limiter.acquire(
                estimated_tokens,
                timeout=5.0 if priority <= Priority.HIGH else 30.0
            )
            if can_acquire:
                return provider
        
        return None
    
    def _rank_providers(self, providers: List[ProviderConfig], 
                        priority: Priority) -> List[ProviderConfig]:
        """Rank providers based on priority requirements."""
        
        if priority <= Priority.HIGH:
            # Prioritize latency for interactive requests
            return sorted(providers, key=lambda p: p.latency_p50_ms)
        else:
            # Prioritize cost for background jobs
            return sorted(providers, key=lambda p: p.cost_multiplier)
    
    def report_success(self, provider_name: str, latency_ms: float):
        provider = self.providers[provider_name]
        provider.is_healthy = True
        provider.consecutive_failures = 0
        # Exponential moving average for latency
        provider.latency_p50_ms = 0.9 * provider.latency_p50_ms + 0.1 * latency_ms
    
    def report_failure(self, provider_name: str, is_rate_limit: bool):
        provider = self.providers[provider_name]
        provider.consecutive_failures += 1
        
        if provider.consecutive_failures >= 5:
            provider.is_healthy = False
            # Schedule health check
            asyncio.create_task(self._health_check(provider_name))
•••

Component 4: Request Batching

For non-interactive workloads, batching dramatically improves throughput.

Python
class RequestBatcher:
    def __init__(self, max_batch_size: int = 20, max_wait_ms: float = 100):
        self.max_batch_size = max_batch_size
        self.max_wait_ms = max_wait_ms
        self.pending: Dict[str, List[QueuedRequest]] = {}  # model -> requests
        self.locks: Dict[str, asyncio.Lock] = {}
    
    async def add_request(self, model: str, request: QueuedRequest) -> dict:
        """Add request to batch and wait for result."""
        
        if model not in self.locks:
            self.locks[model] = asyncio.Lock()
            self.pending[model] = []
        
        async with self.locks[model]:
            self.pending[model].append(request)
            
            # Check if we should flush
            should_flush = (
                len(self.pending[model]) >= self.max_batch_size or
                request.priority <= Priority.HIGH  # Don't batch high priority
            )
        
        if should_flush:
            return await self._flush_batch(model)
        else:
            # Wait for batch to fill or timeout
            await asyncio.sleep(self.max_wait_ms / 1000)
            return await self._flush_batch(model)
    
    async def _flush_batch(self, model: str) -> List[dict]:
        async with self.locks[model]:
            if not self.pending[model]:
                return []
            
            batch = self.pending[model]
            self.pending[model] = []
        
        # Execute batch request
        # OpenAI doesn't have true batching, but we can parallelize
        results = await asyncio.gather(*[
            self._execute_single(model, req) for req in batch
        ], return_exceptions=True)
        
        # Resolve futures
        for req, result in zip(batch, results):
            if isinstance(result, Exception):
                req.callback.set_exception(result)
            else:
                req.callback.set_result(result)
        
        return results
•••

Component 5: Graceful Degradation

When all else fails, degrade gracefully instead of erroring.

Python
class GracefulDegradation:
    def __init__(self):
        self.degradation_level = 0  # 0 = normal, 1 = degraded, 2 = minimal
        self.fallback_responses: Dict[str, str] = {}
    
    async def handle_request(self, request: QueuedRequest, 
                            router: MultiProviderRouter) -> dict:
        try:
            # Try primary path
            provider = await router.route(
                request.payload["model"],
                request.estimated_tokens,
                request.priority
            )
            
            if provider:
                return await self._execute_request(provider, request)
            
            # No provider available
            return await self._handle_no_capacity(request)
            
        except Exception as e:
            return await self._handle_failure(request, e)
    
    async def _handle_no_capacity(self, request: QueuedRequest) -> dict:
        if self.degradation_level == 0:
            # Level 0: Queue and retry
            await asyncio.sleep(1)
            raise RetryableError("No capacity, please retry")
        
        elif self.degradation_level == 1:
            # Level 1: Try smaller/faster model
            fallback_model = self._get_fallback_model(request.payload["model"])
            if fallback_model:
                request.payload["model"] = fallback_model
                return await self.handle_request(request)
            raise RetryableError("No capacity")
        
        else:
            # Level 2: Return cached/default response
            cache_key = self._generate_cache_key(request)
            if cache_key in self.fallback_responses:
                return {
                    "content": self.fallback_responses[cache_key],
                    "degraded": True,
                    "model": "cached"
                }
            raise ServiceUnavailableError("System at capacity")
    
    def _get_fallback_model(self, model: str) -> Optional[str]:
        fallbacks = {
            "gpt-4": "gpt-4-turbo",
            "gpt-4-turbo": "gpt-3.5-turbo",
            "claude-3-opus": "claude-3-sonnet",
            "claude-3-sonnet": "claude-3-haiku",
        }
        return fallbacks.get(model)
•••

Monitoring: Know Before Users Complain

Python
class LLMInfraMetrics:
    def __init__(self):
        self.prometheus = PrometheusClient()
    
    def record_request(self, request_id: str, model: str, provider: str,
                       latency_ms: float, tokens: int, status: str,
                       priority: Priority, queued_time_ms: float):
        
        # Latency histogram
        self.prometheus.histogram(
            "llm_request_latency_ms",
            latency_ms,
            labels={"model": model, "provider": provider, "priority": priority.name}
        )
        
        # Queue time
        self.prometheus.histogram(
            "llm_queue_time_ms",
            queued_time_ms,
            labels={"priority": priority.name}
        )
        
        # Token usage
        self.prometheus.counter(
            "llm_tokens_total",
            tokens,
            labels={"model": model, "provider": provider}
        )
        
        # Success/failure rate
        self.prometheus.counter(
            "llm_requests_total",
            1,
            labels={"model": model, "provider": provider, "status": status}
        )
        
        # Alert on degradation
        if status == "rate_limited":
            self.alert_if_threshold_exceeded(
                "rate_limit_percentage",
                window_minutes=5,
                threshold=0.05  # Alert if >5% rate limited
            )
•••

Results

After implementing this architecture:

MetricBeforeAfter
Success rate at 50K RPM70%99.7%
p50 latency800ms450ms
p99 latency12s2.1s
Cost per 1M requests$2,400$1,850
Incidents/month80
•••

Key Takeaways

  1. 1Rate limiting is multidimensional. Track requests AND tokens across time windows.
  1. 2Not all requests are equal. Priority queuing prevents background jobs from affecting user experience.
  1. 3Multi-provider is mandatory at scale. Azure OpenAI, direct OpenAI, and Anthropic as fallbacks.
  1. 4Degrade gracefully. A slower response beats an error. A cached response beats no response.
  1. 5Monitor everything. You need to know you're approaching limits before you hit them.

Found this helpful?

Share it with others who might benefit

TweetShare

Related articles

📚 Continue Learning