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.ReadTimeout after five seconds on a stream that is simply quiet.
  • requests with stream=True seems 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.

httpx timeouts and what they mean for a stream Matrix of the four httpx timeout kinds, what each measures, the default, and the value suited to an SSE stream. httpx timeouts and what they mean for a stream Timeout Measures Default For SSE connect TCP + TLS setup 5 s 5–10 s read gap between bytes 5 s 2–3 × heartbeat write sending the request 5 s default pool waiting for a connection 5 s default
Only the read timeout interacts with the stream's quiet periods. Set it above the heartbeat interval, not to None, so dead connections are still detected.

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.

Backpressure from a slow worker to the upstream Sequence diagram of the reader putting events into a bounded queue, the queue filling because workers are slow, the reader pausing, and the upstream's writes slowing through TCP flow control. Backpressure from a slow worker to the upstream Upstream Reader task Bounded queue Workers events put (queue fills) get, slowly put() waits, stops reading TCP window full, writes slow down
No buffer grows without bound. The slowest component sets the pace for the whole chain.

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
Why Python SSE consumers stop, by cause Bar chart of the share of consumer stoppages by cause in a fleet of Python consumers before fixes: default read timeout, no reconnect loop, slow handler blocking reads, and genuine server errors. Why Python SSE consumers stop, by cause Default 5 s read timeout 46 % No reconnect loop 28 % Handler blocking reads 17 % Genuine server errors 9 % share of unplanned consumer stoppages by cause
Illustrative breakdown from a fleet review. Three of the four causes are client configuration, fixed by the steps in this guide.

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.