DEV Community

Usman Khan
Usman Khan

Posted on Originally published at ctousman.com

Production AI & LLM Task Pipelines: Managing Token Budgets, Backpressure, and Asynchronous Queue Architecture

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
    }
  }
);
Enter fullscreen mode Exit fullscreen mode

Top comments (0)