Kafka & NATS as SSE Event Sources Permalink to this section

Part of Backend Stream Generation & Connection Management.

Most production event streams do not originate in the process that serves them. Orders are written by one service, prices arrive from a feed handler, notifications are produced by a dozen others; the SSE tier’s job is to take events from a log or bus and deliver them to browsers. Redis pub/sub is the usual first bridge, and the Redis fan-out topic covers it. This guide covers the durable alternatives teams reach for next — Apache Kafka, NATS JetStream and PostgreSQL’s LISTEN/NOTIFY — and the design questions they raise that Redis pub/sub never did. Durable logs change what a reconnecting client can get back, what the SSE id should contain, and how consumers must be configured so that every SSE node sees every event. Getting the consumer model wrong is the most common failure: Kafka consumer groups are built to divide messages among consumers, while an SSE fan-out tier needs every node to receive all of them.

How It Works Permalink to this section

A bridge has three parts: a consumer on each SSE node that reads from the source, an in-process fan-out that delivers each message to the local connections interested in it, and a replay path that serves Last-Event-ID reconnects from the source’s retained history.

Every SSE node consumes the whole topic A Kafka topic or NATS stream is read independently by three SSE nodes, each with its own consumer, and each node fans messages out to its local browser connections. Every SSE node consumes the whole topic Producers orders, prices, jobs Topic / stream durable log append SSE node A own consumer SSE node B own consumer SSE node C own consumer read all
The consumers are independent readers, not members of one group. Dividing the topic among nodes would deliver each event to only one node's clients.

The three sources differ in what they retain and how you address a position:

Source Position Retention Replay by position Delivery to many readers
Kafka partition + offset time/size based, days yes, seek to offset each consumer group reads everything
NATS JetStream stream sequence limits or interest based yes, start at sequence ephemeral/ordered consumers per node
NATS core none none no all subscribers receive
Postgres LISTEN/NOTIFY none none (in-flight only) no — use a table all listeners receive

A durable position turns into an SSE id almost directly. For a single-partition source, the offset or sequence is the id. For Kafka topics with several partitions, a single number cannot represent a position, and the id becomes a compact encoding of per-partition offsets — or the events carry an application-level sequence assigned by the producer.

: NATS JetStream — stream sequence as the SSE id
id: 884213
event: order
data: {"orderId":"o_19f","status":"shipped"}

: Kafka, 3 partitions — encoded partition offsets as the id
id: 0-5812.1-5790.2-5833
event: order
data: {"orderId":"o_19g","status":"paid"}

Server-Side Implementation Permalink to this section

Kafka: one consumer group per node, or no group at all Permalink to this section

Each SSE node must read the entire topic. Two configurations achieve that: give each node a unique consumer group id (so the group has one member, which is assigned all partitions), or assign partitions manually without a group. Unique groups are simplest but leave orphaned groups behind after every deploy unless their ids are stable per node; manual assignment avoids group coordination entirely and suits a fan-out tier well.

// kafkajs — each node reads all partitions from the latest offset, no shared group.
import { Kafka } from 'kafkajs';

const kafka = new Kafka({ clientId: `sse-${process.env.HOSTNAME}`, brokers: BROKERS });
const consumer = kafka.consumer({ groupId: `sse-fanout-${process.env.HOSTNAME}` }); // unique per node
await consumer.connect();
await consumer.subscribe({ topic: 'orders', fromBeginning: false });

await consumer.run({
  autoCommit: false,                           // positions live in the browsers, not the group
  eachMessage: async ({ partition, message }) => {
    const evt = JSON.parse(message.value.toString());
    positions[partition] = message.offset;     // track the latest offset per partition
    hub.publish(evt.userId, {
      id: encodePositions(positions),
      event: 'order',
      data: message.value.toString(),
    });
  },
});

Disabling auto-commit is deliberate. The fan-out tier does not “process” messages in the at-least-once sense; it only needs the live tail, and the durable positions that matter belong to each browser, carried in Last-Event-ID. Bridging Kafka topics to SSE streams covers partition-offset ids and seeking for replay.

NATS JetStream: ordered consumers per node Permalink to this section

JetStream’s ordered consumer is designed for this job: an ephemeral, flow-controlled consumer that delivers the stream in order to one reader and recreates itself transparently if a message is missed.

// nats.js v2 — each node gets its own ordered consumer on the stream.
import { connect, DeliverPolicy } from 'nats';

const nc = await connect({ servers: NATS_URL });
const js = nc.jetstream();
const consumer = await js.consumers.get('ORDERS', {
  deliver_policy: DeliverPolicy.New,           // live tail only; replay is served separately
  filterSubjects: ['orders.>'],
});
for await (const m of await consumer.consume()) {
  const evt = m.json();
  hub.publish(evt.userId, { id: String(m.seq), event: 'order', data: m.string() });
}

For replay, a separate short-lived ordered consumer starts at the client’s sequence:

async function* replayFrom(seq) {
  const c = await js.consumers.get('ORDERS', { deliver_policy: DeliverPolicy.StartSequence, opt_start_seq: seq + 1 });
  const msgs = await c.fetch({ max_messages: 500, expires: 2000 });
  for await (const m of msgs) yield m;
}

Streaming NATS JetStream to SSE clients builds the full endpoint with per-user subject filters.

Joining replay and live without a gap Permalink to this section

Whatever the source, a reconnecting client needs history from its position followed by the live tail, with no gap and no duplicate at the seam. The node’s live consumer is already running, so the handler registers the connection with the in-process hub first (buffering live messages), then reads history from the source, then flushes the buffer, discarding anything at or below the last replayed position.

Replay from the log, then switch to the node's live tail Sequence diagram of a browser reconnecting with an event id, the SSE node registering for live messages, reading history from the log after that position, sending it, then sending buffered live messages above the last replayed position. Replay from the log, then switch to the node's live tail Browser SSE node Hub Log GET Last-Event-ID: 884200 register, buffer live read from 884201 live 884215 (buffered) 884201 … 884215 replay, then live from 884216
The seam is handled per connection with a position comparison. The node's shared live consumer never rewinds for anyone.

The same shape in Python with confluent-kafka, for a single-partition topic where the offset is the id:

# replay.py — seek a short-lived consumer to the client's offset, read up to the high watermark.
from confluent_kafka import Consumer, TopicPartition

def replay(topic: str, after: int, limit: int = 500):
    c = Consumer({"bootstrap.servers": BROKERS, "group.id": "sse-replay",
                  "enable.auto.commit": False})
    tp = TopicPartition(topic, 0, after + 1)
    low, high = c.get_watermark_offsets(tp, timeout=2)
    if after + 1 < low:
        c.close()
        return None                                  # aged out: caller sends resync
    c.assign([tp])
    out = []
    while len(out) < limit and tp.offset < high:
        msg = c.poll(1.0)
        if msg is None or msg.error():
            break
        out.append((msg.offset(), msg.value()))
        tp.offset = msg.offset() + 1
    c.close()
    return out

Replay consumers are assigned partitions directly and never join the live group, so they cannot trigger rebalances. Cap how much one reconnect may replay; beyond the cap, send resync and let the client fetch state.

Postgres LISTEN/NOTIFY: a signal, not a log Permalink to this section

NOTIFY delivers a payload of up to 8,000 bytes to every session currently listening on a channel, at transaction commit. Nothing is retained: a node that was reconnecting when the notification fired never receives it. That makes it an excellent wake-up signal and a poor event log. The robust design writes events to a table and uses NOTIFY only to tell listeners “there is something new after sequence N”:

CREATE TABLE events (seq bigserial PRIMARY KEY, topic text NOT NULL, body jsonb NOT NULL, at timestamptz DEFAULT now());

CREATE FUNCTION notify_event() RETURNS trigger AS $$
BEGIN
  PERFORM pg_notify('events', NEW.seq::text);    -- tiny payload: just the new sequence
  RETURN NEW;
END $$ LANGUAGE plpgsql;

CREATE TRIGGER events_notify AFTER INSERT ON events FOR EACH ROW EXECUTE FUNCTION notify_event();

Each node holds one listening connection, and on every notification reads rows after its last seen sequence. Missed notifications cost nothing, because the next one — or a periodic check — reads everything since. Using Postgres LISTEN/NOTIFY for SSE covers connection pooling pitfalls and payload limits.

Choosing a source for an SSE fan-out tier Matrix comparing Redis pub/sub, Postgres LISTEN/NOTIFY, NATS JetStream and Kafka on durability, replay, throughput and operational weight. Choosing a source for an SSE fan-out tier Source Durable Replay by id Throughput Ops weight Redis pub/sub no no high light Postgres NOTIFY + table table by seq moderate already there NATS JetStream yes by seq high moderate Kafka yes per partition very high heavy strong adequate weak
Replay capability is the dividing line. Without it, reconnect replay must come from a separate store; with it, the source itself serves Last-Event-ID.

Client-Side Consumption Permalink to this section

The browser side does not change: EventSource stores whatever string the server put in id: and sends it back as Last-Event-ID. The only rule is that the server must be able to interpret every id it has ever issued, for as long as the source retains the data. Encoded Kafka positions are opaque to the client but must stay parseable by future server versions — version the encoding (v1:0-5812.1-5790) so it can evolve.

Clients do need to handle the case where their id is older than the source’s retention: the server answers with a resync event, and the client refetches current state over HTTP. For notification-style feeds, that is the notification feed inbox fetch; for state feeds, a snapshot.

es.addEventListener('resync', async () => {
  const state = await (await fetch('/api/orders/recent', { credentials: 'include' })).json();
  store.replaceAll(state);             // the stream continues from the server's current position
});

Edge Cases & Network Interference Permalink to this section

  • Rebalancing storms. Putting all SSE nodes in one Kafka consumer group means every deploy triggers a rebalance, during which consumption pauses, and each node only sees a share of partitions. Use per-node groups or manual assignment.
  • Consumer lag at startup. A node that starts with fromBeginning: true replays days of history into the fan-out before reaching the live tail. The live consumer should start at the latest position; replay is per client, on demand.
  • Retention shorter than disconnection. A client offline for a weekend may present an id whose offset has been deleted. Seek fails or returns a later offset; detect it and send resync.
  • Message size. Kafka and JetStream accept large messages, but browsers render SSE frames on the main thread. Keep event payloads small; send identifiers and let clients fetch details.
  • Ordering across partitions. Kafka orders within a partition only. Key messages by the entity whose order matters (user, order id) so all of its events land in one partition.
  • Pooled Postgres connections. LISTEN is session state. Through PgBouncer in transaction mode, the listening session disappears between transactions; the listener needs a direct, dedicated connection.

Mitigation checklist:

Performance & Scale Considerations Permalink to this section

The fan-out tier multiplies traffic twice: every node reads every message, and every node writes each message to many connections. With N nodes and M messages per second, the broker serves N × M reads per second — cheap for Kafka and JetStream, which are built for many readers of the same data, but worth checking for broker egress bandwidth.

Broker egress as the SSE tier scales out Line chart of broker egress bandwidth as the number of SSE nodes grows, for a 2 MB/s topic, compared with the egress needed if nodes only read the partitions their users need. Broker egress as the SSE tier scales out every node reads all key-routed nodes 0 20 40 60 80 0 8 16 24 32 40 SSE nodes broker egress (MB/s)
Whole-topic consumption grows linearly with nodes. At large node counts, route users to nodes by key so each node reads only its share of partitions.

For very large deployments, invert the design: route each user to an SSE node by a hash of the same key used to partition the topic, and have each node consume only the partitions its users hash to. That caps broker egress at one copy of the topic, at the cost of a routing layer (consistent hashing at the load balancer, or a redirect on connect). Below a few dozen nodes, whole-topic consumption is simpler and entirely adequate.

Filtering belongs as close to the source as possible. NATS subjects can encode the audience (orders.user.42), so a node’s consumer can filter by subject and never receive messages for users it does not hold; JetStream supports multiple filter subjects per consumer. Kafka has no server-side filtering, so everything reaches the node and is filtered in process — another reason to keep payloads small and to partition by the key that decides the audience. With Postgres, the notification carries only a sequence number and each node’s follow-up query can include a WHERE topic = ANY($1) clause listing the topics its connected users care about.

Per node, the hot path is deserialising each message once, looking up interested local connections, and writing a pre-serialised frame to each. Avoid re-serialising per connection, and avoid per-connection consumers: a consumer per browser connection turns 50,000 viewers into 50,000 broker sessions.

Validation & Debugging Permalink to this section

# Kafka: each SSE node should appear as its own group with lag near zero.
kafka-consumer-groups.sh --bootstrap-server $B --list | grep sse-fanout
kafka-consumer-groups.sh --bootstrap-server $B --describe --group sse-fanout-node-a

# NATS: list consumers on the stream; each node's ordered consumer shows its delivered sequence.
nats consumer ls ORDERS
nats stream info ORDERS    # first/last sequence = the replayable window

# Replay: reconnect with an old id and check that events resume after it.
curl -sN -H 'Last-Event-ID: 884200' https://app.example.com/api/orders/stream | grep -m3 '^id:'

Monitor end-to-end latency from produce time to SSE write (stamp the producer’s timestamp in the message and observe the difference at write), consumer lag per node, and the rate of resync events. Rising consumer lag on one node while others are current means that node’s fan-out is too slow for the topic — usually slow clients blocking the consumer loop, which the fan-out must never allow.

Production Checklist Permalink to this section

Frequently Asked Questions Permalink to this section

Why not put all SSE nodes in one Kafka consumer group?

A consumer group divides partitions among its members, so each node would receive only part of the topic and its clients would miss events published to other partitions. Fan-out needs every node to read everything.

Can the SSE id be a Kafka offset?

Only for single-partition topics. With several partitions, encode an offset per partition in the id, or have producers assign an application-level sequence that the server can map back to positions.

Is NATS core enough, without JetStream?

For a live-only feed, yes: core NATS delivers to all subscribers with very low latency. It retains nothing, so resumable streams need JetStream or a separate replay store.

Is LISTEN/NOTIFY reliable enough for SSE?

As a signal, yes. Notifications are lost when no session is listening, so store events in a table and use notifications only to trigger reads after the last seen sequence.

How much history should the source retain?

At least as long as the disconnections you want to repair without a resync: minutes for dashboards, hours to days for notification feeds. Anything older is better served by an HTTP fetch of current state.

Deep Dives