Skip to main content
Create your own

Implementing a Capacity-Limited Request Queue with Graceful Degradation

Introduction

In the previous lesson, we built a token-bucket rate limiter to protect our service on a per-client basis. This is a crucial defense against individual clients who might intentionally or unintentionally overwhelm the system. However, what happens when the aggregate traffic from many well-behaved clients simultaneously exceeds your system's total processing capacity? This scenario, common during viral traffic spikes, can lead to cascading failures if not managed correctly.

Today, we shift our focus from per-client policies to global system health. Your learning outcome is to implement a request queue with a defined capacity and a graceful degradation policy (e.g., 503 Service Unavailable) for overload scenarios. This is a form of load shedding, a fundamental pattern for building robust, scalable systems that can survive under extreme pressure. You'll learn how to make your service predictably say "no" to new work in order to protect the quality of service for the work it has already accepted.

The Problem of Overload: Backpressure and Load Shedding

Even with rate limiting in place, a sudden surge in legitimate users can saturate your inference engine. If your API server continues to accept requests faster than the GPU can process them, these pending requests will pile up in memory, leading to two critical problems:

  1. Unbounded Latency: Requests at the back of an ever-growing queue will eventually time out, but only after consuming server resources for an extended period.
  2. Resource Exhaustion: The queue itself consumes memory. An unbounded queue will eventually cause your API server to run out of memory and crash (an OOM error).

The core concept for managing this is backpressure. To get a clear intuition for this idea, please watch the following video.

Backpressure in Software Development simply explained

The video 'Backpressure in Software Development simply explained' by Software Developer Diaries uses a helpful fluid dynamics analogy to introduce backpressure and discusses the primary strategies for handling it: controlling, buffering, and dropping.

Watch the video from the beginning until 06:06. Pay close attention to the producer/consumer problem and the three strategies discussed. Our implementation will combine 'buffering' (with a limit) and 'dropping'.

As the video explains, when a producer (your API server) is faster than a consumer (your inference engine), you must decide what to do with the excess. A key strategy is load shedding: consciously refusing to accept new work to protect system stability.

To explore this concept in more detail, please read the selected sections from the article "Scalable Systems: Load Control and Backpressure Patterns."

Scalable Systems: Load Control and Backpressure Patterns

This article from Krun.pro provides an excellent, concise overview of load control strategies, emphasizing why saying 'no' is a critical feature of a scalable system.

Please read the following parts: The introduction, under the heading 'Scalable Systems: Load Control'. Focus on the idea of 'disciplined admission control'. The section 'System Load Shedding Strategies'. This directly addresses our goal of returning a 503 error. The FAQ question #6, which clearly distinguishes between the rate limiting we did last lesson and the load shedding we are doing now.

The key takeaway is that it's far better to immediately reject a request with a 503 Service Unavailable error than to accept it and let it time out after 30 seconds. This preserves resources for the requests already in progress and provides a clear, fast signal to the client.

Architectural Design: The Bounded Request Queue

The most direct way to implement load shedding is by placing a bounded queue between your API ingestion layer and your core processing logic.

LLM Inference Request Flow Sequence Diagram
This sequence diagram illustrates the flow of a request. Our bounded queue fits right after the 'API queues the request' step. The API server will attempt to place the request in a queue for the Scheduler. If this queue is full, the API server will immediately reject the request.

By setting a maximum size on this queue, you are effectively capping the maximum amount of work your system will ever have pending. This has two benefits:

  • It puts a ceiling on the memory used by queued requests.
  • It caps the maximum possible wait time for any given request.

The article "Async & Batching in RAG" has a section on "Backpressure Management" that formalizes this pattern as a best practice for high-throughput AI systems. Although you don't need to read the whole article, its recommendation is directly on point: implement backpressure by limiting queue sizes and rejecting new requests with a 503 error code when queues are full.

Implementing a Bounded Queue in FastAPI

Let's translate this design into a practical FastAPI implementation. We will use asyncio.Queue from Python's standard library, which allows us to create a queue with a maxsize. This structure is perfect for integrating with FastAPI's asynchronous request handling.

The architecture will be:

  1. API Endpoint (Producer): The FastAPI endpoint that receives the HTTP request. It will try to put the request into our bounded queue. If the queue is full, it immediately returns a 503 response.
  2. Bounded Queue: An asyncio.Queue instance with a defined maxsize. This queue will live for the duration of the application's lifecycle.
  3. Engine Worker (Consumer): A background task that continuously pulls requests from the queue and "processes" them. In our full system, this worker would be the entry point to your custom inference engine's scheduler.

Here is the implementation sketch:

import asyncio
import time
import random
from fastapi import FastAPI, Request, HTTPException
from contextlib import asynccontextmanager




# --- Configuration ---
MAX_QUEUE_SIZE = 50  # Max number of requests to hold before rejecting new ones




# --- Application State ---
# This dictionary will hold the application's state, including the request queue.
# It's populated by the lifespan manager.
app_state = {}

async def inference_worker():
    """
    This worker simulates our inference engine. It pulls requests from the queue
    and processes them one by one.
    """
    print("Inference worker started.")
    while True:
        try:



            # Wait for a request to appear in the queue
            request_data, received_at = await app_state["request_queue"].get()
            
            queue_wait_time = time.monotonic() - received_at
            print(f"Worker pulled request. Waited in queue for {queue_wait_time:.4f}s.")




            # Simulate the work of LLM inference
            processing_time = 0.5 + random.uniform(-0.2, 0.2)
            await asyncio.sleep(processing_time)

            print(f"Finished processing request after {processing_time:.4f}s.")
            



            # Mark the task as done
            app_state["request_queue"].task_done()
        
        except asyncio.CancelledError:
            print("Inference worker shutting down.")
            break
        except Exception as e:
            print(f"Error in inference worker: {e}")



            # Avoid rapid-fire error loops
            await asyncio.sleep(1)


@asynccontextmanager
async def lifespan(app: FastAPI):
    """
    Manages the application's startup and shutdown logic.
    """



    # Startup: Create the bounded queue and start the worker task
    print("Application startup...")
    app_state["request_queue"] = asyncio.Queue(maxsize=MAX_QUEUE_SIZE)
    worker_task = asyncio.create_task(inference_worker())
    app_state["worker_task"] = worker_task
    
    yield # The application is now running
    



    # Shutdown: Gracefully stop the worker task
    print("Application shutdown...")
    worker_task.cancel()
    try:
        await worker_task
    except asyncio.CancelledError:
        print("Worker task was successfully cancelled.")


app = FastAPI(lifespan=lifespan)

@app.post("/predict")
async def predict(request: Request):
    """
    This endpoint acts as the producer. It receives requests and tries to
    enqueue them.
    """
    request_data = await request.json()
    
    try:



        # Get the current time to measure queue wait
        received_at = time.monotonic()
        



        # Put the request into the queue without waiting.
        # If the queue is full, this will raise asyncio.QueueFull.
        app_state["request_queue"].put_nowait((request_data, received_at))
        
        qsize = app_state["request_queue"].qsize()
        print(f"Request enqueued. Current queue size: {qsize}/{MAX_QUEUE_SIZE}")
        
        return {"message": "Request accepted and is being processed.", "queue_position": qsize}

    except asyncio.QueueFull:



        # If the queue is full, reject the request immediately.
        print(f"Queue is full. Rejecting request. Queue size: {app_state['request_queue'].qsize()}/{MAX_QUEUE_SIZE}")
        raise HTTPException(
            status_code=503,
            detail="Service Unavailable: Server is at capacity. Please try again later.",
            headers={"Retry-After": "10"} # Suggest client to wait 10s
        )

How to Test It

To see the load shedding in action, you can use a simple bash script to send a burst of concurrent requests. Save the Python code as main.py and run it with uvicorn main:app. Then, run this script in your terminal:




#!/bin/bash
# Send 100 concurrent requests to the endpoint
for i in {1..100}; do
    curl -X POST -H "Content-Type: application/json" -d '{"prompt": "hello"}' http://127.0.0.1:8000/predict &
done

wait

You will see in your uvicorn logs that the first ~50 requests are accepted, and the rest are immediately rejected with a 503 error, demonstrating that your system is successfully protecting itself from being overwhelmed.

Conclusion

You have now implemented a critical defense mechanism for a production-grade service. While rate limiting protects you from individual bad actors, load shedding with a bounded queue protects your system's core stability from aggregate, legitimate traffic spikes.

Key Takeaways:

  • Load Shedding vs. Rate Limiting: Rate limiting is a per-client policy; load shedding is a global system health mechanism. You need both.
  • The Danger of Unbounded Queues: They lead to resource exhaustion (OOM errors) and infinite request latencies.
  • Bounded Queues as Admission Control: By setting a maxsize on your request queue, you create a simple and effective admission control system.
  • Graceful Degradation with 503: Immediately returning 503 Service Unavailable is better for both the server and the client than accepting a request that is destined to fail after a long wait.

Preview of the Next Lesson:
We've just introduced a new component—the request queue—that adds latency to our request path. But how does this queuing time compare to other sources of latency, like network I/O, tokenization, or model execution? In the next lesson, you will learn to profile the full end-to-end API request path to identify and optimize non-inference latency bottlenecks, giving you a complete picture of your system's performance.

Can't find a good explanation? Sign up and we'll make it for you

Sign up