Consuming SSE in Python Permalink to this section
Part of Non-Browser SSE Clients, under Frontend Consumption & Client Patterns.
Python code consumes Server-Sent Events in data pipelines, agents calling language-model APIs, monitoring scripts and integration tests. The ecosystem has good pieces — httpx for streaming HTTP and httpx-sse for parsing — but neither reconnects, and httpx’s defaults are tuned for short requests. This guide builds a consumer that stays connected, resumes after drops, and processes events without blocking the stream, in both synchronous and asyncio code.
Symptom & Developer Intent Permalink to this section
httpx.ReadTimeoutafter five seconds on a stream that is simply quiet.requestswithstream=Trueseems to work, then delivers events in bursts or not at all.- The consumer processes a slow event and the upstream connection times out while it does.
- After a network drop the script exits instead of reconnecting.
- A long-running asyncio consumer never notices that the connection died.
The intent is a small, reusable consumer that handles timeouts, reconnection and slow processing correctly.
Root Cause Analysis Permalink to this section
httpx applies a default timeout of five seconds to connect, read, write and pool operations. For a stream, a read timeout fires whenever no bytes arrive for that long — any quiet period longer than five seconds. requests has no default timeout at all, which avoids that problem but means a dead connection is never detected, and its line iteration helpers buffer internally, which can delay events.
Processing inside the read loop is the other trap: while the handler runs, nothing reads the socket. A long enough handler lets the server’s send buffer fill, or the read timeout fire. Handing events to a queue decouples the two.
Step-by-Step Resolution Permalink to this section
Step 1 — Configure timeouts for streaming Permalink to this section
import httpx
HEARTBEAT_S = 15
timeout = httpx.Timeout(connect=10.0, read=HEARTBEAT_S * 3, write=10.0, pool=10.0)
client = httpx.Client(timeout=timeout, headers={"User-Agent": "orders-consumer/1.4"})
A read timeout of three heartbeats both tolerates quiet periods and detects dead connections.
Step 2 — Wrap httpx-sse in a resuming loop (sync) Permalink to this section
import random, time
from httpx_sse import connect_sse
class Fatal(Exception): ...
def follow(url, handle, get_token, last_id=None, stop=lambda: False):
retry_ms, backoff = 3000, 1.0
while not stop():
headers = {"Authorization": f"Bearer {get_token()}"}
if last_id: headers["Last-Event-ID"] = last_id
try:
with connect_sse(client, "GET", url, headers=headers) as source:
status = source.response.status_code
if status in (401, 403, 404): raise Fatal(status)
source.response.raise_for_status()
for sse in source.iter_sse():
if sse.id: last_id = sse.id
if sse.retry: retry_ms = sse.retry
handle(sse)
backoff = 1.0
except Fatal:
raise
except (httpx.TransportError, httpx.HTTPStatusError):
pass # network or 5xx: retry
time.sleep(max(retry_ms / 1000, backoff) * (1 + random.random() * 0.3))
backoff = min(backoff * 2, 60)
return last_id
Treat 401 specially if tokens expire: refresh and retry once before giving up.
Step 3 — The asyncio version, with a processing queue Permalink to this section
import asyncio
from httpx_sse import aconnect_sse
async def follow_async(url, queue: asyncio.Queue, get_token, last_id=None):
async with httpx.AsyncClient(timeout=timeout) as aclient:
retry_s, backoff = 3.0, 1.0
while True:
headers = {"Authorization": f"Bearer {await get_token()}"}
if last_id: headers["Last-Event-ID"] = last_id
try:
async with aconnect_sse(aclient, "GET", url, headers=headers) as source:
async for sse in source.aiter_sse():
if sse.id: last_id = sse.id
if sse.retry: retry_s = sse.retry / 1000
await queue.put(sse) # blocks when full: reading pauses
backoff = 1.0
except httpx.TransportError:
pass
await asyncio.sleep(max(retry_s, backoff) * (1 + random.random() * 0.3))
backoff = min(backoff * 2, 60)
async def worker(queue):
while True:
sse = await queue.get()
await process(sse) # slow work happens here
queue.task_done()
async def main():
q = asyncio.Queue(maxsize=500) # bounded: backpressure to the socket
await asyncio.gather(follow_async(URL, q, token), *(worker(q) for _ in range(4)))
Because queue.put waits when the queue is full, the reader stops pulling from the socket, and TCP flow control slows the upstream instead of the consumer’s memory growing.
Step 4 — Persist the position Permalink to this section
For consumers that restart, store last_id after processing — in a file, Redis or a database row — and pass it back into follow on start-up. Handlers must tolerate the occasional repeat.
Step 5 — Stream from POST endpoints Permalink to this section
Language-model and search APIs expect a POST with a JSON body; connect_sse accepts any method and body:
with connect_sse(client, "POST", "https://llm.example.com/v1/responses",
json={"model": "m", "input": prompt, "stream": True},
headers={"Authorization": f"Bearer {key}"}) as source:
for sse in source.iter_sse():
if sse.event == "response.output_text.delta":
print(sse.json()["delta"], end="", flush=True)
Do not wrap one-shot generation streams in a blind reconnect loop — retrying a POST starts new work. Resume only through an API that supports it, as described in sending POST requests that return SSE.
Validation & Monitoring Permalink to this section
# Local test server that goes quiet for 30 s: the consumer must not time out (read timeout 45 s).
python -m tests.sse_server --pause 30 &
python consumer.py --url http://localhost:8000/events
Test the three configuration failures deliberately: a server that pauses for longer than the old default timeout, a server that drops the connection mid-stream, and a handler that sleeps longer than the heartbeat interval. With the fixes in place, the first produces no error, the second a quick reconnect with resume, and the third a pause in reading with no timeout.
Log each connection with the resume id, event counts and the exception that ended it. For long-running consumers, export time since last event, queue depth and reconnects by cause, and alert when time since last event exceeds several heartbeat intervals.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Can I use the requests library for SSE?
It can stream with stream=True and iter_lines, but it has no async support, no read-timeout defaults suited to streams and its line iteration can buffer. httpx with httpx-sse is the better-maintained path.
Why set a read timeout at all?
Without one, a connection that dies silently is never detected and the consumer waits forever. A read timeout of a few heartbeats detects it quickly without cutting healthy quiet periods.
Does httpx-sse reconnect automatically?
No. It parses the stream and exposes events. Reconnection, backoff and Last-Event-ID are the loop in this guide.
Can I follow several streams in one asyncio program?
Yes. Run one follow_async task per stream, each with its own queue or a shared bounded queue, and its own persisted cursor. The event loop holds idle streams cheaply, so dozens or hundreds are practical in one process.
How do I stop a consumer cleanly?
Cancel the reader task (asyncio) or set the stop flag (sync), let workers drain the queue, then persist the last processed id before exiting.