Streaming NATS JetStream to SSE Clients Permalink to this section

Part of Kafka & NATS as SSE Event Sources, under Backend Stream Generation & Connection Management.

NATS JetStream is a particularly good fit behind a Server-Sent Events tier. Each stream has a single, monotonically increasing sequence — a ready-made SSE id — subjects can encode exactly who an event is for, and consumers can be created in milliseconds and filtered on the server. This guide wires a JetStream stream to browser clients: subject design, a live tail per node, replay by start sequence, and the handling of sequences that have been purged.

Symptom & Developer Intent Permalink to this section

  • The SSE tier receives every message on the stream and discards most of them after filtering in process.
  • Durable consumers pile up on the server, one per browser that ever connected.
  • Reconnecting clients either receive nothing they missed or receive the whole stream from the start.
  • After a stream limit purges old messages, reconnects with old ids fail.
  • Live delivery stalls occasionally and resumes with a burst.

The intent is a design where each node receives only the messages its clients need, reconnects replay exactly from the client’s sequence, and no consumer state leaks per browser.

Root Cause Analysis Permalink to this section

JetStream consumers are cheap, but durable consumers are persistent server state. Creating one per browser connection means cleaning them up forever. The right primitives are an ordered consumer per node for the live tail — ephemeral, flow-controlled, self-healing — and a short-lived consumer per reconnect for replay, started at last sequence + 1.

Subjects that encode the audience Layers of a subject hierarchy from the stream prefix through the audience type, the audience id, and the event type. Subjects that encode the audience events stream prefix user / room / public audience type 42 audience id order.shipped event type
When the audience is part of the subject, NATS filters for you. Nodes subscribe to the audiences they hold instead of receiving everything.

In-process filtering is wasteful because the subject can carry the audience: events.user.42.order.shipped. A consumer with filter subjects events.user.42.> and events.public.> receives only what matters. Stalls followed by bursts usually come from a push consumer without flow control, or from a slow SSE write loop blocking the consumer’s message handler.

Step-by-Step Resolution Permalink to this section

Step 1 — Define the stream with subjects and limits Permalink to this section

nats stream add EVENTS \
  --subjects 'events.>' \
  --storage file --retention limits \
  --max-age 72h --max-bytes 20GB \
  --discard old --dupe-window 2m

max-age is your replay window: a client can resume after up to 72 hours away. The duplicate window lets producers publish with a Nats-Msg-Id header so retries do not create duplicate events.

await js.publish(`events.user.${userId}.order.shipped`, JSON.stringify(evt), { msgID: evt.eventId });

Step 2 — One live consumer per node, filtered by held audiences Permalink to this section

// Each node follows the stream from "now", filtered to public events plus users it serves.
import { connect, DeliverPolicy } from 'nats';

const nc = await connect({ servers: NATS_URL });
const js = nc.jetstream();

async function startLive() {
  const c = await js.consumers.get('EVENTS', {
    deliver_policy: DeliverPolicy.New,
    filterSubjects: ['events.public.>', 'events.user.>'],
  });
  for await (const m of await c.consume({ max_messages: 1000 })) {
    const [, kind, id] = m.subject.split('.');
    const frame = `id: ${m.seq}\nevent: ${m.subject.split('.').slice(3).join('.')}\ndata: ${m.string()}\n\n`;
    if (kind === 'public') hub.broadcast(frame, m.seq);
    else hub.toUser(id, frame, m.seq);                           // only if this node holds user id
  }
}

For very large fleets, route users to nodes and give each node’s consumer filter subjects for just its users. JetStream supports many filter subjects per consumer, and recreating the ordered consumer with a new filter set when users arrive is cheap.

Step 3 — Replay from the client’s sequence on reconnect Permalink to this section

async function replay(res, userId, lastSeq) {
  const info = await js.streams.info('EVENTS');
  if (lastSeq + 1 < info.state.first_seq) {                      // purged by limits
    res.write('event: resync\ndata: {"reason":"expired"}\n\n');
    return info.state.last_seq;
  }
  const c = await js.consumers.get('EVENTS', {
    deliver_policy: DeliverPolicy.StartSequence,
    opt_start_seq: lastSeq + 1,
    filterSubjects: [`events.user.${userId}.>`, 'events.public.>'],
  });
  let seq = lastSeq;
  const batch = await c.fetch({ max_messages: 500, expires: 1500 });
  for await (const m of batch) {
    res.write(`id: ${m.seq}\nevent: ${m.subject.split('.').slice(3).join('.')}\ndata: ${m.string()}\n\n`);
    seq = m.seq;
  }
  return seq;                                                    // live frames must be > seq
}

Filtered replay skips sequences for other audiences, so ids the client receives have gaps; that is expected and harmless. Only the comparison “greater than the last id” matters.

Reconnect with a JetStream sequence Sequence diagram of a browser reconnecting with Last-Event-ID 884200, the SSE node registering for live delivery, creating a start-sequence consumer, replaying filtered messages, and switching to live. Reconnect with a JetStream sequence Browser SSE node JetStream GET Last-Event-ID: 884200 register user 42 for live consumer start_seq 884201, filter user.42 884203, 884210 … replayed frames, then live > 884210
The replay consumer lives for one request and is filtered to the user's subjects, so it never delivers another user's events.

Step 4 — Keep the consumer loop non-blocking Permalink to this section

The loop that reads JetStream messages must never wait on a browser socket. Hand frames to per-connection bounded buffers and return immediately; a slow client fills its own buffer and is disconnected, as in handling slow consumers with SSE backpressure. Otherwise one slow phone stalls delivery for every client on the node, which is the stall-then-burst symptom.

Step 5 — Authorise subjects, not just connections Permalink to this section

Because subjects carry the audience, authorisation becomes a question of which subjects a connection may be filtered to. Derive the filter list on the server from the authenticated identity — events.user.${req.user.id}.> plus the rooms the user is a member of — and never accept subject names from the client. For room membership that changes while a stream is open, publish a control message on events.user.<id>.control.revoke that the node handles by closing that user’s streams; the reconnect then computes a fresh filter list.

If NATS itself enforces permissions (per-account or per-user subject permissions), give the SSE tier its own NATS user with subscribe rights on events.> and no publish rights beyond its control subjects. The SSE tier is then a read-only projection of the stream, which limits the damage a compromised node can do.

Validation & Monitoring Permalink to this section

# Stream window and consumer count (should be ~one per node, not per browser).
nats stream info EVENTS | grep -E 'First Sequence|Last Sequence|Consumers'

# Publish a test event and watch it reach a stream.
nats pub events.user.42.order.shipped '{"orderId":"o_1"}' &
curl -sN -b s.txt https://app.example.com/api/stream | grep -m2 -E '^(id|event):'
Messages delivered to one node per second Bar chart comparing messages delivered to an SSE node from JetStream when consuming the entire stream, when filtering to public plus all user subjects, and when filtering to only the users held by that node. Messages delivered to one node per second Whole stream 8,000 / s public + all users 8,000 / s public + held users 450 / s messages per second received by one of 20 nodes, stream at 8,000 msg/s
Subject filtering moves the discard from the SSE node to the NATS server, where it costs almost nothing.

Monitor pending messages on each node’s live consumer (should stay near zero), replay requests and durations, and resyncs by reason. A growing pending count means the node’s fan-out cannot keep up and is blocking the consumer.

Production Checklist Permalink to this section

Frequently Asked Questions Permalink to this section

Should each browser get its own JetStream consumer?

Not a durable one. A short-lived consumer for a replay is fine, but live delivery should come from one consumer per node fanned out in memory, or consumer counts grow with every connection.

Why do the ids a client receives have gaps?

The sequence is per stream, and a filtered consumer skips messages for other subjects. Clients only need ids to increase; gaps are normal.

Is core NATS enough for live-only feeds?

Yes. For dashboards and presence where replay is not needed, core subjects with wildcard subscriptions give the lowest latency and no server state.

How do I handle a stream purge or a limit change?

Compare the client's sequence with the stream's first sequence on every reconnect. Anything older than the first sequence cannot be replayed; send resync so the client fetches current state.