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
Component 1: Intelligent Request Queue
Not all requests are equal. A user waiting for a response needs priority over a background task.
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.
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.
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.
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.
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
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:
| Metric | Before | After |
|---|---|---|
| Success rate at 50K RPM | 70% | 99.7% |
| p50 latency | 800ms | 450ms |
| p99 latency | 12s | 2.1s |
| Cost per 1M requests | $2,400 | $1,850 |
| Incidents/month | 8 | 0 |
Key Takeaways
- 1Rate limiting is multidimensional. Track requests AND tokens across time windows.
- 2Not all requests are equal. Priority queuing prevents background jobs from affecting user experience.
- 3Multi-provider is mandatory at scale. Azure OpenAI, direct OpenAI, and Anthropic as fallbacks.
- 4Degrade gracefully. A slower response beats an error. A cached response beats no response.
- 5Monitor everything. You need to know you're approaching limits before you hit them.
