Skip to content

Streaming GraphQL Subscriptions over WebSockets

GraphQL subscriptions push results to clients over a long-lived connection — price ticks, chat messages, job progress. In Strawberry a subscription resolver is an async generator: each yield sends one result, and the generator lives as long as the client stays subscribed. That makes the lifecycle explicit, and it puts two obligations on your code: fan events out to every subscriber without letting one slow client hold up the others, and release per-subscriber resources when the client leaves. Measured with Strawberry 0.329 on FastAPI and uvicorn, Python 3.14, using the graphql-transport-ws protocol: 500 subscribers received a stream of 20 events published 10 ms apart, while half of them disconnected after the fifth event. The server delivered 6,250 events — exactly the 250 × 20 plus 250 × 5 expected — with a publish-to-receipt latency of 11.6 ms at p50 and 33.5 ms at p99, and when the clients had gone, the subscriber registry was back to 0 with all 500 generators' finally blocks having run. This guide builds subscriptions that fan out and clean up.

Prerequisites

1. Write the subscription as an async generator

A subscription field returns an AsyncGenerator; Strawberry iterates it and sends each yielded value to the client as a next message:

from collections.abc import AsyncGenerator
import strawberry


@strawberry.type
class PriceTick:
    symbol: str
    price: float
    sent_at: float


@strawberry.type
class Subscription:
    @strawberry.subscription
    async def prices(self, symbol: str = "ACME") -> AsyncGenerator[PriceTick, None]:
        queue: asyncio.Queue[PriceTick] = asyncio.Queue(maxsize=100)
        HUB.subscribers.add(queue)
        try:
            while True:
                tick = await queue.get()
                if tick.symbol == symbol:
                    yield tick
        finally:
            HUB.subscribers.discard(queue)            # runs when the client leaves


schema = strawberry.Schema(query=Query, subscription=Subscription)
app.include_router(GraphQLRouter(schema, subscription_protocols=[GRAPHQL_TRANSPORT_WS_PROTOCOL]),
                   prefix="/graphql")

Clients connect with the graphql-transport-ws subprotocol: connection_init, a subscribe message carrying the query, then a next message per result until either side sends complete. The finally block is the subscription's cleanup hook; Strawberry closes the generator when the client completes the subscription or the socket closes. Measured: after 500 clients connected and then all disconnected, the registry was empty and the cleanup counter read 500.

Verify: after a load test that ends with clients disconnecting, the server's subscriber registry returns to zero.

One subscription over graphql-transport-ws A sequence of 9 messages between 4 participants. One subscription over graphql-transport-ws client Strawberry subscription generator hub connection_init connection_ack subscribe {query} start generator register queue event -> queue yield -> next {data} socket closed finally: unregister The generator's lifetime is the subscription's lifetime.

2. Fan out through per-subscriber bounded queues

Publishing should never wait for subscribers. Give each subscriber its own bounded queue and have the publisher push to all of them without awaiting:

class Hub:
    def __init__(self) -> None:
        self.subscribers: set[asyncio.Queue] = set()

    def publish(self, event) -> None:
        for queue in self.subscribers:
            if queue.full():
                queue.get_nowait()                     # drop this subscriber's oldest event
                DROPPED.inc()
            queue.put_nowait(event)

publish is synchronous and constant-time per subscriber, so a stalled client cannot slow the publisher or other subscribers; its queue fills, and its oldest events are dropped. For price ticks, dropping stale values is the right policy; for chat or audit events, close the slow subscriber instead, so it reconnects and resynchronizes — the trade-offs are measured in bounding actor mailboxes under load. Measured: with 500 subscribers, all 20 events reached every subscriber still connected, with p50 latency 11.6 ms and p99 33.5 ms from publish to receipt on the same machine.

Verify: a deliberately stalled test client accumulates drops while other clients' latency stays unchanged.

3. Filter at the source, not in every generator

The generator above receives every event and filters by symbol — fine for a few subscribers, wasteful for thousands each wanting a different symbol. Key the registry by what subscribers asked for:

from collections import defaultdict


class TopicHub:
    def __init__(self) -> None:
        self.by_topic: dict[str, set[asyncio.Queue]] = defaultdict(set)

    def publish(self, topic: str, event) -> None:
        for queue in self.by_topic.get(topic, ()):
            if queue.full():
                queue.get_nowait()
            queue.put_nowait(event)

    def subscribe(self, topic: str, queue: asyncio.Queue) -> None:
        self.by_topic[topic].add(queue)

    def unsubscribe(self, topic: str, queue: asyncio.Queue) -> None:
        subs = self.by_topic.get(topic)
        if subs:
            subs.discard(queue)
            if not subs:
                del self.by_topic[topic]               # no empty sets left behind

Publishing cost then scales with the subscribers of that topic, not with all subscribers. Removing empty topic entries in unsubscribe keeps the registry from growing with every symbol anyone ever subscribed to. Across several server processes, the hub's input comes from a broker — Redis pub/sub or a stream — with each process fanning out to its own subscribers, the pattern in scaling WebSockets across processes with Redis pub/sub.

Verify: publishing to a topic with no subscribers costs a dictionary lookup, and the registry's size tracks active topics.

500 subscribers, 20 events, half leave after 5 A grid of 4 rows by 2 columns. 500 subscribers, 20 events, half leave after 5 measure value subscribers connected 500 events delivered 6,250 (250 x 20 + 250 x 5) latency p50 / p99 11.6 ms / 33.5 ms registry after disconnects 0 subscribers, 500 cleanups Strawberry 0.329, FastAPI, uvicorn, graphql-transport-ws; client and server on one machine.

4. Authenticate and limit subscriptions at connect time

A subscription holds a connection and a generator for as long as the client wants. Check identity once, when the WebSocket connects, and bound how many subscriptions each user may hold:

from strawberry.fastapi import GraphQLRouter


class AuthedRouter(GraphQLRouter):
    async def on_ws_connect(self, context: dict) -> dict | None:
        token = (context.get("connection_params") or {}).get("token")
        user = await verify_token(token)
        if user is None:
            raise ConnectionRejectionError({"reason": "unauthorized"})
        context["user"] = user
        return {"user": user.id}                      # sent back in connection_ack


@strawberry.subscription
async def orders(self, info: strawberry.Info) -> AsyncGenerator[Order, None]:
    user = info.context["user"]
    async with SUBSCRIPTION_SLOTS[user.id]:           # e.g. Semaphore(5) per user
        ...

Tokens arrive in connection_init's payload because browsers cannot set arbitrary headers on WebSocket handshakes. A per-user cap stops one client from opening thousands of subscriptions on one socket; a server-wide cap and the connection-level limits in authenticating WebSocket connections bound total memory. Consider re-validating long-lived subscriptions periodically, since a token valid at connect time may be revoked an hour later.

Verify: a connection without a valid token is rejected before any subscription starts, and a user's sixth subscription fails.

5. Keep heartbeats and timeouts on long-lived sockets

Subscriptions can be idle for minutes between events, and idle connections are what proxies and load balancers close first. The graphql-transport-ws protocol has ping and pong messages, and the WebSocket layer has its own ping frames:

app.include_router(
    GraphQLRouter(schema,
                  subscription_protocols=[GRAPHQL_TRANSPORT_WS_PROTOCOL],
                  connection_init_wait_timeout=timedelta(seconds=10)),
    prefix="/graphql",
)
# uvicorn: --ws-ping-interval 20 --ws-ping-timeout 20

connection_init_wait_timeout closes sockets that connect but never initialize — cheap to open, they otherwise hold a slot indefinitely. Server pings at an interval shorter than the shortest idle timeout on the path keep healthy connections open and detect dead ones, after which the generator's finally runs. Tuning the interval is covered in tuning WebSocket ping/pong heartbeats.

Verify: a client that connects and stays silent is closed after the init timeout; a subscribed client idle for longer than the load balancer's idle timeout stays connected.

How should this subscription deliver events? A decision on What do the events represent with 4 outcomes. How should this subscription deliver events? What do the events represent? latest state, e.g. prices bounded queue, drop oldest stale values are worthless events that must all arrive close slow subscribers client resyncs many keys, few per subscriber topic-keyed hub publish cost per topic several server processes broker feeds each hub Redis pub/sub or streams Publish never awaits a subscriber; cleanup always runs in finally.

Verification

Subscriptions are production-ready when:

  • Each subscription is an async generator with cleanup in finally, and the registry returns to zero after clients leave.
  • Publishing never awaits subscribers; each has a bounded queue and a drop or close policy.
  • The hub is keyed by topic, and empty topics are removed.
  • Connections are authenticated at init, capped per user, and kept alive with pings.

Diagnostic Hook: export the number of active subscriptions, per-subscriber drops and publish-to-send latency. Active subscriptions that keep rising while connections stay flat point at generators that never finish their cleanup; rising drops for a few subscribers are slow clients; rising latency for everyone means the event loop, not the clients, is saturated.

Pitfalls & edge cases

  • Awaiting queue.put in the publisher. One stalled subscriber then stalls every event.
  • Cleanup outside finally. It does not run when the client disconnects mid-stream.
  • Filtering in every generator. Publishing cost grows with all subscribers instead of a topic's.
  • Authenticating per event. Do it once at connection init, and re-check on a schedule.

Frequently Asked Questions

How do GraphQL subscriptions work in Strawberry?

A field decorated with @strawberry.subscription returns an AsyncGenerator; each yield sends a result to the client over a WebSocket using the graphql-transport-ws protocol, and the generator is closed when the client completes or disconnects.

How do I clean up when a GraphQL subscriber disconnects?

Put the cleanup in a finally block in the subscription generator; Strawberry closes the generator on disconnect. In testing, all 500 subscribers' finally blocks ran and the registry returned to zero.

How do I broadcast one event to many GraphQL subscribers?

Give each subscriber its own bounded asyncio.Queue and have the publisher put_nowait into each, dropping the oldest when full. 500 subscribers received events at a p50 latency of 11.6 ms in testing.

How do I authenticate GraphQL subscriptions?

Send the token in the connection_init payload and verify it in the router's on_ws_connect hook, rejecting the connection if invalid; then cap subscriptions per user.