4. Async & streaming¶
Intermediate · 11 min read
An LLM call takes seconds. While one request waits for OpenAI, the server should keep serving other
users — that's what async gives you. And instead of making the user stare at a spinner for 10 seconds,
you stream tokens as they are generated.
New to async/await? Read Async programming first.
4.1 async def or def?¶
| Your endpoint does… | Write | Why |
|---|---|---|
await calls (async LLM SDK, httpx.AsyncClient, async DB) |
async def |
runs on the event loop; thousands of waits at once |
blocking calls (requests, sync SDK, time.sleep, CPU work) |
def |
FastAPI runs it in a thread pool so it doesn't block others |
| nothing slow | either | — |
The one mistake to avoid
A blocking call inside async def (e.g. requests.post(...) or the sync OpenAI() client)
freezes the whole server for every user until it finishes. Use the async client
(AsyncOpenAI, AsyncAnthropic, httpx.AsyncClient) — or make the endpoint a plain def.
4.2 Concurrent LLM calls in one request¶
Ask several things at once with asyncio.gather — total time ≈ the slowest call, not the sum:
import asyncio, time
from fastapi import FastAPI
from fastapi.testclient import TestClient
from pydantic import BaseModel
app = FastAPI()
client = TestClient(app)
async def fake_llm(prompt: str, delay: float = 0.3) -> str:
await asyncio.sleep(delay) # stands in for `await client.chat.completions.create(...)`
return f"answer to {prompt!r}"
class MultiQuery(BaseModel):
questions: list[str]
@app.post("/multi")
async def multi(req: MultiQuery):
start = time.perf_counter()
answers = await asyncio.gather(*(fake_llm(q) for q in req.questions))
return {"answers": answers, "seconds": round(time.perf_counter() - start, 1)}
r = client.post("/multi", json={"questions": ["a", "b", "c", "d"]}).json()
print(len(r["answers"]), "answers in", r["seconds"], "s") # not 4 × 0.3 = 1.2 s
This is how you run multi-query retrieval, parallel tool calls or several LLM judges in one request.
4.3 Streaming a response — StreamingResponse¶
Return a generator; each yield is sent to the client immediately.
from fastapi.responses import StreamingResponse
async def fake_token_stream(prompt: str):
for token in ["Retrieval", "-augmented", " generation", " grounds", " answers", "."]:
await asyncio.sleep(0.01)
yield token
@app.post("/stream")
async def stream(prompt: str):
return StreamingResponse(fake_token_stream(prompt), media_type="text/plain")
r = client.post("/stream", params={"prompt": "What is RAG?"})
print(r.headers["content-type"])
print(r.text)
TestClient collects the whole body before returning, so you see the joined text. To watch the tokens
arrive one by one, run the app and use curl's -N (no buffering) flag:
With a real SDK the generator wraps the provider's stream:
# no-run — needs an API key
from openai import AsyncOpenAI
llm = AsyncOpenAI()
async def openai_tokens(prompt: str):
stream = await llm.chat.completions.create(
model="gpt-4o-mini", stream=True,
messages=[{"role": "user", "content": prompt}],
)
async for chunk in stream:
delta = chunk.choices[0].delta.content
if delta:
yield delta
4.4 Server-Sent Events (SSE)¶
Plain text streams are hard to extend. SSE is a simple standard: each message is a data: ... line
followed by a blank line. Browsers read it with EventSource, and you can send JSON events — tokens,
sources, usage, errors, a final [DONE]:
import json
async def sse_events(question: str):
yield f"data: {json.dumps({'type': 'sources', 'sources': ['policy.pdf#p3']})}\n\n"
async for token in fake_token_stream(question):
yield f"data: {json.dumps({'type': 'token', 'text': token})}\n\n"
yield f"data: {json.dumps({'type': 'usage', 'completion_tokens': 6})}\n\n"
yield "data: [DONE]\n\n"
@app.post("/chat/stream")
async def chat_stream(question: str):
return StreamingResponse(sse_events(question), media_type="text/event-stream",
headers={"Cache-Control": "no-cache"})
answer = ""
with client.stream("POST", "/chat/stream", params={"question": "What is RAG?"}) as r:
for line in r.iter_lines():
if not line.startswith("data: ") or line == "data: [DONE]":
continue
event = json.loads(line[6:])
if event["type"] == "token":
answer += event["text"]
else:
print(event)
print(answer)
{'type': 'sources', 'sources': ['policy.pdf#p3']}
{'type': 'usage', 'completion_tokens': 6}
Retrieval-augmented generation grounds answers.
This is the same format the OpenAI and Anthropic APIs use to stream to you.
Behind a proxy?
Nginx and some load balancers buffer responses, which kills streaming. Send the header
X-Accel-Buffering: no and check your platform's docs for SSE support.
4.5 Timeouts¶
Never wait forever on an upstream LLM:
from fastapi import HTTPException
@app.post("/chat/timeout")
async def chat_with_timeout(prompt: str):
try:
return {"reply": await asyncio.wait_for(fake_llm(prompt, delay=2), timeout=0.2)}
except asyncio.TimeoutError:
raise HTTPException(status_code=504, detail="LLM took too long")
r = client.post("/chat/timeout", params={"prompt": "hi"})
print(r.status_code, r.json())
4.6 Background tasks — work after the response¶
Logging usage, saving the conversation, sending feedback to an eval store — the user shouldn't wait for these:
from fastapi import BackgroundTasks
usage_log: list[dict] = []
def record_usage(user: str, tokens: int):
usage_log.append({"user": user, "tokens": tokens}) # e.g. write to a database
@app.post("/chat/logged")
async def chat_logged(prompt: str, background: BackgroundTasks):
reply = await fake_llm(prompt, delay=0)
background.add_task(record_usage, user="priya", tokens=len(reply.split()))
return {"reply": reply}
print(client.post("/chat/logged", params={"prompt": "hello"}).json())
print(usage_log)
Background tasks are not a job queue
They run in the same process and are lost if it restarts. For long jobs — ingesting 1,000 PDFs, re-embedding a corpus — use a real queue (Celery, RQ, Arq, Cloud Tasks) and return a job ID the client can poll.
4.7 WebSockets (two-way)¶
SSE is server → client only. For voice agents or live collaboration you need both directions:
from fastapi import WebSocket
@app.websocket("/ws")
async def ws_chat(ws: WebSocket):
await ws.accept()
while True:
msg = await ws.receive_text()
if msg == "bye":
await ws.close()
break
await ws.send_text(f"echo: {msg}")
with client.websocket_connect("/ws") as ws:
ws.send_text("hello")
print(ws.receive_text())
ws.send_text("bye")
For a normal chat UI, SSE is simpler and works through more proxies — pick WebSockets only when the client must also stream (audio, interruptions).
Practice¶
- Change
/multito use anasyncio.Semaphore(2)so at most two LLM calls run at once. - Add an
errorevent tosse_eventsand send it if the token stream raises an exception.
Next: FastAPI for GenAI — a production-shaped LLM service.