Production AI & LLM Task Pipelines: Managing Token Budgets, Backpressure, and Asynchronous Queue Architecture
Originally published at ctousman.com.
Integrating LLM capabilities into enterprise SaaS requires moving beyond simple synchronous HTTP calls. When hit with spikes in demand, vendor rate limits (Tokens Per Minute & Requests Per Minute) and unpredictable inference latencies can cripple your backend.
Here is how to design asynchronous LLM queue pipelines with token-budget throttling, worker backpressure, and graceful degradation.
1. The Problem with Synchronous LLM Integration
Executing model calls (e.g., GPT-4o, Claude 3.5 Sonnet, or self-hosted vLLM instances) directly inside HTTP request handlers causes systemic vulnerabilities in web application servers:
- Thread Pool Exhaustion: Long-running inference queries (spanning 5 to 30 seconds) block web server connection workers, rapidly exhausting socket pools under concurrent traffic.
- Vendor Rate-Limit Cascades: Unthrottled bursts breach provider Requests Per Minute (RPM) or Tokens Per Minute (TPM) caps, triggering HTTP 429 errors and cascading failures across client applications.
- Unpredictable Financial Spikes: Lack of centralized queue concurrency management makes it difficult to enforce global billing guardrails and token budget safety nets across tenants.
2. Architectural Blueprint: The Token-Aware Queue Pipeline
A production-ready pipeline decouples payload submission from model execution using an asynchronous message broker (such as Redis with BullMQ) coupled with a dynamic rate-limiting layer that tracks estimated token weights before invoking vendor APIs.
Token Weighting and Dynamic Throttling
Unlike standard API rate limiters that count discrete HTTP requests, LLM limiters must track dynamic token consumption. Before dispatching a job, the worker estimates prompt tokens and checks the sliding window token bucket.
typescript
import { Worker, Job } from 'bullmq';
import { Redis } from 'ioredis';
import { getEncoding } from 'js-tiktoken';
const redis = new Redis(process.env.REDIS_URL);
const tokenizer = getEncoding('cl100k_base');
interface LLMJobPayload {
tenantId: string;
prompt: string;
maxTokens: number;
}
export const llmWorker = new Worker<LLMJobPayload>(
'llm-generation-queue',
async (job: Job<LLMJobPayload>) => {
const { tenantId, prompt, maxTokens } = job.data;
// 1. Calculate estimated prompt token weight
const promptTokens = tokenizer.encode(prompt).length;
const estimatedTotalTokens = promptTokens + maxTokens;
// 2. Evaluate current global token budget (Sliding Window in Redis)
const allowed = await checkTokenBudget(redis, estimatedTotalTokens);
if (!allowed) {
// Delay job execution dynamically without failing the task
await job.moveToDelayed(Date.now() + 5000, job.token);
throw new Error('RATE_LIMIT_DELAY: Dynamic backpressure applied');
}
// 3. Execute model API call with timeout protection
const response = await executeModelCall({
prompt,
maxTokens,
timeoutMs: 45000
});
// 4. Reconcile exact consumed tokens from API metadata
await reconcileActualTokenUsage(redis, response.usage.total_tokens);
return response.content;
},
{
concurrency: 10,
limiter: {
max: 500,
duration: 60000
}
}
);
Top comments (0)