06 · WebSockets & Streaming Responses¶
A normal endpoint computes a whole response and sends it. Some features don't fit that shape: a chat, live progress for a long job, a feed of price updates. FastAPI offers two families of tools:
- Streaming responses — one HTTP request, a response that arrives in pieces over time (JSON Lines, Server-Sent Events). Server → client only, plain HTTP, works through most proxies.
- WebSockets — an upgraded connection where both sides can send messages at any time. Needed when the client also talks continuously.
Prefer streaming when data flows one way; it's simpler to deploy, cache-friendly to
reason about, and reconnects automatically in the browser's EventSource.
Streaming JSON Lines from a yield endpoint¶
On the version used here (FastAPI 0.143.0), an endpoint can simply be an async generator with a declared item type:
import asyncio
from collections.abc import AsyncIterable
from fastapi import FastAPI
from pydantic import BaseModel
app = FastAPI()
class Progress(BaseModel):
step: int
of: int
@app.get("/jobs/42/lines")
async def lines() -> AsyncIterable[Progress]:
for i in range(1, 4):
await asyncio.sleep(0.2)
yield Progress(step=i, of=3)
Under Uvicorn, with curl -N (no buffering):
HTTP/1.1 200 OK
server: uvicorn
content-type: application/jsonl
transfer-encoding: chunked
{"step":1,"of":3}
{"step":2,"of":3}
{"step":3,"of":3}
Each item was validated against Progress, serialized, and sent as one line as soon as
it was yielded. transfer-encoding: chunked means there's no Content-Length — the
server doesn't know the total size in advance.
Version-specific
Yield-based endpoints and fastapi.sse are recent additions to FastAPI. On older
versions you return a StreamingResponse wrapping a generator yourself, and use a
third-party package (such as sse-starlette) for SSE. Check
python -c "import fastapi.sse" to see whether your version has it.
Server-Sent Events¶
SSE is a standard text format for event streams that browsers consume with the built-in
EventSource API:
from fastapi.sse import EventSourceResponse, ServerSentEvent
@app.get("/jobs/42/events", response_class=EventSourceResponse)
async def events() -> AsyncIterable[ServerSentEvent]:
for i in range(1, 4):
await asyncio.sleep(0.2)
yield ServerSentEvent(data={"step": i, "of": 3}, event="progress", id=str(i))
yield ServerSentEvent(data={"done": True}, event="done")
HTTP/1.1 200 OK
server: uvicorn
content-type: text/event-stream; charset=utf-8
cache-control: no-cache
x-accel-buffering: no
transfer-encoding: chunked
event: progress
data: {"step": 1, "of": 3}
id: 1
event: progress
data: {"step": 2, "of": 3}
id: 2
...
event: done
data: {"done": true}
Each event is a block of field: value lines ended by a blank line. Two headers were
added for you: cache-control: no-cache, and x-accel-buffering: no, which tells nginx
not to buffer the stream (a buffering proxy would hold events back until the response
ended, defeating the point). Timing confirmed the stream: the first bytes arrived after
0.0015 s, the whole response after 0.61 s — three 0.2 s sleeps.
In a browser:
const es = new EventSource("/jobs/42/events");
es.addEventListener("progress", (e) => console.log(JSON.parse(e.data)));
es.addEventListener("done", () => es.close());
If the connection drops, EventSource reconnects by itself and sends a
Last-Event-ID header with the last id it saw, so the server can resume.
(EventSource reconnects after a done too, unless you call close().)
WebSockets¶
from fastapi import WebSocket, WebSocketDisconnect
@app.websocket("/ws/echo")
async def echo(ws: WebSocket):
await ws.accept()
try:
while True:
text = await ws.receive_text()
await ws.send_text(f"echo: {text}")
except WebSocketDisconnect as e:
print(f"echo client left, code={e.code}")
The test client can drive WebSockets too:
Code 1000 is a normal close (the with block ending). Every WebSocket endpoint needs
that try/except WebSocketDisconnect: a client vanishing is the normal end of a
WebSocket's life, not an error.
Worked example: an authenticated chat room¶
Browsers can't set an Authorization header on a WebSocket, so tokens usually travel in
the query string (or a cookie, or the first message). Here, a query token, and a room
that broadcasts to everyone connected:
from typing import Annotated
from fastapi import Query, WebSocketException, status
TOKENS = {"tok-ada": "ada", "tok-bob": "bob"}
class Room:
def __init__(self):
self.members: dict[str, WebSocket] = {}
async def join(self, user: str, ws: WebSocket):
await ws.accept()
self.members[user] = ws
await self.broadcast({"event": "joined", "user": user})
def leave(self, user: str):
self.members.pop(user, None)
async def broadcast(self, msg: dict):
for name, ws in list(self.members.items()):
try:
await ws.send_json(msg)
except Exception:
self.leave(name)
room = Room()
@app.websocket("/ws/chat")
async def chat(ws: WebSocket, token: Annotated[str | None, Query()] = None):
user = TOKENS.get(token or "")
if user is None:
raise WebSocketException(code=status.WS_1008_POLICY_VIOLATION, reason="invalid token")
await room.join(user, ws)
try:
while True:
data = await ws.receive_json()
text = str(data.get("text", ""))[:500]
await room.broadcast({"event": "message", "user": user, "text": text})
except WebSocketDisconnect:
room.leave(user)
await room.broadcast({"event": "left", "user": user})
Two clients through the test client:
bad token -> WebSocketDisconnect 1008 invalid token
ada got: {'event': 'joined', 'user': 'ada'}
ada got: {'event': 'joined', 'user': 'bob'}
bob got: {'event': 'joined', 'user': 'bob'}
ada got: {'event': 'message', 'user': 'bob', 'text': 'hi ada'}
bob got: {'event': 'message', 'user': 'bob', 'text': 'hi ada'}
ada got: {'event': 'left', 'user': 'bob'}
And against a real Uvicorn server with the websockets client library:
ada: {"event":"joined","user":"ada"}
ada: {"event":"message","user":"ada","text":"anyone here?"}
bad token: InvalidStatus server rejected WebSocket connection: HTTP 403
with these server log lines:
INFO: 127.0.0.1:52265 - "WebSocket /ws/chat?token=tok-ada" [accepted]
INFO: 127.0.0.1:52266 - "WebSocket /ws/chat?token=nope" 403
INFO: connection rejected (403 Forbidden)
The test client reported the rejection as close code 1008, while the real client saw an
HTTP 403: the exception was raised before accept(), so the handshake itself was
refused and no WebSocket ever existed. That's the right place to reject — but test the
real behaviour with a real client at least once.
Note [:500] on the text and list(...) around the members while broadcasting: input
limits and safe iteration matter even more for long-lived connections.
How It Actually Works¶
Streaming: FastAPI wraps your generator in a streaming response. Starlette sends
http.response.start (status and headers) immediately, then one http.response.body
message with more_body=True per chunk, and a final empty one. Uvicorn writes each as
an HTTP/1.1 chunk. Because headers go first, the status code is fixed before your first
yield — an exception halfway through can't turn into a 500; the stream just breaks.
Send errors as an event in the stream when the client needs to know.
WebSockets: the client sends an HTTP GET with Upgrade: websocket. Uvicorn turns
that into an ASGI scope of type "websocket" and delivers a websocket.connect event.
Your endpoint's accept() sends websocket.accept, and the server completes the
handshake with 101 Switching Protocols. If the app closes or raises before accepting,
Uvicorn answers the handshake with HTTP 403 instead. After that, receive_* and
send_* exchange websocket.receive/websocket.send events, and a client close
arrives as websocket.disconnect, which Starlette raises as WebSocketDisconnect.
The Room here lives in one process's memory. With several workers or servers, Ada and
Bob may be connected to different processes that can't see each other's members. Real
deployments put a pub/sub layer (Redis, NATS, PostgreSQL LISTEN/NOTIFY) between
processes: each process subscribes, and broadcasts go through the broker.
Common mistakes¶
- No
WebSocketDisconnecthandling, leaving dead connections in your registry. - Blocking calls inside a WebSocket loop. One
time.sleepstalls every connection in the process (Level 2 lesson 3). - Authenticating after
accept(), or not at all. Reject during the handshake. - In-memory rooms with multiple workers.
- Proxies buffering streams — configure them (
x-accel-buffering, orproxy_buffering offfor nginx), and set proxy timeouts long enough for idle WebSockets or send periodic pings. - Using WebSockets for one-way updates that SSE would handle with less code.
- Unbounded message sizes from clients.
Exercise¶
- Add a
GET /jobs/{id}/eventsstream that sends aheartbeatcomment every 15 seconds while waiting for progress, and ends with adoneevent. - Support
Last-Event-ID: read the header and skip events the client has already seen. - Extend the chat room with named rooms (
/ws/rooms/{name}), per-room membership and a limit of 50 members per room. - Write a test using the test client that connects three users, disconnects one
abruptly, and asserts the others receive a
leftevent.