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:
- 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.
- 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.

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:
- 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.
- Bounded Queue: An
asyncio.Queueinstance with a definedmaxsize. This queue will live for the duration of the application's lifecycle. - 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
maxsizeon your request queue, you create a simple and effective admission control system. - Graceful Degradation with 503: Immediately returning
503 Service Unavailableis 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.