Skip to content

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
Output
4 answers in 0.3 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)
Output
text/plain; charset=utf-8
Retrieval-augmented generation grounds answers.

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:

curl -N -X POST "http://127.0.0.1:8000/stream?prompt=What+is+RAG"

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)
Output
{'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())
Output
504 {'detail': 'LLM took too long'}

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)
Output
{'reply': "answer to 'hello'"}
[{'user': 'priya', 'tokens': 3}]

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")
Output
echo: hello

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 /multi to use an asyncio.Semaphore(2) so at most two LLM calls run at once.
  • Add an error event to sse_events and send it if the token stream raises an exception.

Next: FastAPI for GenAI — a production-shaped LLM service.