Bridging Kafka Topics to SSE Streams Permalink to this section
Part of Kafka & NATS as SSE Event Sources, under Backend Stream Generation & Connection Management.
Kafka is already the system of record for events in many organisations, so serving those events to browsers over Server-Sent Events is an obvious win: no extra broker, and a durable log that can replay what a disconnected client missed. The bridge is also easy to get subtly wrong, because Kafka’s defaults are built for work queues, not fan-out. This guide builds a bridge for a multi-partition topic that delivers every event to every interested browser, encodes a resumable position in the SSE id, and replays from Kafka itself on reconnect.
Symptom & Developer Intent Permalink to this section
- Clients connected to one SSE node receive only some events; clients on another node receive the rest.
- Every deploy pauses all streams for several seconds while “rebalancing” appears in the logs.
- After a reconnect, clients miss events that were published during the gap, even though Kafka retains them.
- A client that reconnects after a long weekend receives an error or an unexpected flood of old events.
- Events for the same order occasionally arrive out of order in the browser.
The intent is a bridge where every node sees the whole topic, reconnects resume at exactly the right position across all partitions, and positions that have aged out of retention are handled gracefully.
Root Cause Analysis Permalink to this section
Kafka’s consumer group protocol assigns each partition to exactly one member of a group. That is right for processing orders once; for fan-out it means each SSE node only sees its assigned share. Every membership change — a node starting or stopping — triggers a rebalance that pauses consumption.
The missing-events-after-reconnect symptom comes from not using Kafka for replay at all: the bridge only tails the live topic, so anything published while a client was away is gone from that client’s point of view. The fix is to make the SSE id a Kafka position and seek to it.
Out-of-order delivery comes from partitioning: Kafka orders messages within a partition, not across them. If events for one order are spread over partitions, the browser can see them in any order.
Step-by-Step Resolution Permalink to this section
Step 1 — Key messages by the entity whose order matters Permalink to this section
// Producer: all events for one order go to one partition, in order.
await producer.send({
topic: 'orders',
messages: [{ key: order.id, value: JSON.stringify(event) }],
});
Step 2 — Consume every partition on every node, from latest Permalink to this section
// Manual assignment: no group membership, no rebalances.
import { Kafka } from 'kafkajs';
const kafka = new Kafka({ brokers: BROKERS });
const admin = kafka.admin();
const offsets = await admin.fetchTopicOffsets('orders'); // [{ partition, high, low }]
const consumer = kafka.consumer({ groupId: `sse-${process.env.HOSTNAME}` });
await consumer.connect();
await consumer.subscribe({ topic: 'orders', fromBeginning: false });
const head = new Map(offsets.map((o) => [o.partition, Number(o.high) - 1]));
await consumer.run({
autoCommit: false,
eachMessage: async ({ partition, message }) => {
head.set(partition, Number(message.offset));
const evt = JSON.parse(message.value);
hub.deliver(evt.userId, frame(partition, Number(message.offset), message.value));
},
});
A stable per-node group id (the pod name in a StatefulSet, for example) avoids leaving thousands of abandoned groups on the brokers; with ephemeral hostnames, clean up old groups with an admin job.
Step 3 — Encode per-partition positions in the SSE id Permalink to this section
A browser’s position is the last offset it has seen in each partition. Encode it compactly and version the format:
// v1:<p>-<offset>.<p>-<offset>… e.g. "v1:0-5812.1-5790.2-5833"
function encodeId(pos) {
return 'v1:' + [...pos].sort((a, b) => a[0] - b[0]).map(([p, o]) => `${p}-${o}`).join('.');
}
function decodeId(id) {
if (!id?.startsWith('v1:')) return null;
return new Map(id.slice(3).split('.').map((s) => s.split('-').map(Number)));
}
Each connection tracks its own position map, starting from the replay position and advancing as it writes events. Every frame’s id: is the encoded map, so the browser always holds a complete resume point. For topics with many partitions, encode only the partitions relevant to the user (those their keys hash to) to keep ids short.
Step 4 — Replay by seeking a short-lived consumer Permalink to this section
async function* replay(positions, userId, limit = 500) {
const c = kafka.consumer({ groupId: `sse-replay-${crypto.randomUUID()}` });
await c.connect();
await c.subscribe({ topic: 'orders' });
const bounds = new Map((await admin.fetchTopicOffsets('orders'))
.map((o) => [o.partition, { low: Number(o.low), high: Number(o.high) }]));
for (const [p, off] of positions) {
if (off + 1 < bounds.get(p).low) throw new AgedOut(); // retention passed this position
}
const out = [];
await c.run({ eachMessage: async ({ partition, message }) => {
const evt = JSON.parse(message.value);
if (evt.userId === userId) out.push({ partition, offset: Number(message.offset), value: message.value });
}});
for (const [p, off] of positions) c.seek({ topic: 'orders', partition: p, offset: String(off + 1) });
await waitUntilCaughtUp(c, bounds); // stop at the high watermarks
await c.disconnect();
yield* out.slice(0, limit);
}
The handler registers the connection for live delivery first, runs the replay, sends it, then sends buffered live frames whose position is past the replay — the same seam discussed in the Kafka and NATS topic.
Step 5 — Answer aged-out positions with resync Permalink to this section
try {
for await (const m of replay(decodeId(req.get('Last-Event-ID')), user.id)) write(m);
} catch (e) {
if (!(e instanceof AgedOut)) throw e;
res.write('event: resync\ndata: {"reason":"retention"}\n\n'); // client refetches state
}
Replaying from Kafka is not free — each replay creates a consumer and reads partitions. For high-churn clients, keep a small in-memory ring of recent frames per node and serve short gaps from it, falling back to Kafka only for longer ones.
Validation & Monitoring Permalink to this section
# Every node should read all partitions: check assignments for each node's group.
kafka-consumer-groups.sh --bootstrap-server $B --describe --group sse-node-a | awk '{print $3}' | sort -u
# Replay: reconnect with a known position and compare against kcat output from that offset.
kcat -b $B -t orders -p 1 -o 5791 -c 5 -f '%o %s\n'
curl -sN -H 'Last-Event-ID: v1:0-5812.1-5790.2-5833' https://app.example.com/api/orders/stream | head
Track per node: consumer lag (should be near zero), replays per minute, replay duration, and resyncs by reason. A spike in resyncs after a retention change is expected; a steady rate means retention is shorter than typical client absences.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
Can I store positions in Kafka's committed offsets instead of the SSE id?
Committed offsets belong to a consumer group, not to a browser. Thousands of clients at different positions cannot share one group's offsets, so the position must travel with the client in Last-Event-ID.
How long can the encoded id be?
Browsers accept long ids, but they are resent on every reconnect and logged by proxies. Keep them to the partitions a user's events can land in, and use a compact encoding.
What about topics with hundreds of partitions?
Have producers stamp a per-user or global sequence and keep a small index from sequence to partition offset, or serve replay from a secondary store keyed by user. Per-partition ids stop being practical beyond a few dozen partitions.
Does this work with Kafka-compatible brokers?
Yes. Redpanda and other Kafka API implementations support the same consumer, seek and watermark operations the bridge relies on.