Per-User Channels for Notification Delivery Permalink to this section
Part of Notification & Activity Feeds, under Real-Time Application Patterns.
“Publish to user:42” is the first design everyone writes, and it works until the node holding user 42’s connection has 40,000 other users on it too. This guide compares the three ways to route a per-user message to the right SSE connection across a fleet, shows where each one breaks, and builds the design that holds up: a local routing table fed by a bounded number of subscriptions per node.
Symptom & Developer Intent Permalink to this section
Routing problems surface as infrastructure symptoms rather than application bugs:
- Redis CPU climbs steadily as concurrent users grow, even though notification volume is flat.
- Deploys take longer each month, because every restarting node re-subscribes to tens of thousands of channels.
PUBSUB NUMSUBshows channels with zero subscribers receiving publishes, or users with two subscriptions after a reconnect.- Some notifications are delayed by seconds when a node is busy re-subscribing.
- The broker’s client output buffer limit disconnects a node, and every user on it misses notifications until they reconnect.
The intent is routing whose cost on the broker is proportional to the number of nodes, not users, with in-memory delivery to local connections.
Root Cause Analysis Permalink to this section
With one channel per user, each connected user adds one SUBSCRIBE on the node that holds them. Redis maintains a hash table from channel to subscribers; publish is cheap, but subscribe and unsubscribe churn is not free, and connection-heavy nodes spend real time managing tens of thousands of subscriptions. Worse, the subscription set is rebuilt from scratch when a node restarts.
A pattern subscription (PSUBSCRIBE user:*) fixes the subscription count but inverts the cost: every node receives every notification and discards all but its own users’ messages. That is fine at low notification volume and wasteful at high volume, because the broker’s outbound traffic becomes notifications × nodes.
Node-addressed channels get both benefits. A small registry records which node holds each user’s connection; publishers look up the node and publish to node:<id>; each node has one subscription and routes locally. The price is keeping the registry correct.
Step-by-Step Resolution Permalink to this section
Step 1 — Keep a local routing table on every node Permalink to this section
Whatever the broker design, each node needs a map from user id to the set of open streams for that user (a user may have several tabs or devices on the same node).
// routing.js — per-node, in memory.
const streams = new Map(); // userId → Set<res>
export function attach(userId, res) {
let set = streams.get(userId);
if (!set) streams.set(userId, (set = new Set()));
set.add(res);
return () => {
set.delete(res);
if (!set.size) streams.delete(userId);
};
}
export function deliverLocal(userId, frame) {
const set = streams.get(userId);
if (!set) return 0;
for (const res of set) if (!res.writableNeedDrain) res.write(frame);
return set.size;
}
Step 2 — Register which node holds each user, with a lease Permalink to this section
// registry.js — Redis hash per user with expiring node entries.
const NODE_ID = process.env.NODE_ID;
const LEASE_S = 60;
export async function register(userId) {
await redis.multi()
.hset(`conn:${userId}`, NODE_ID, Date.now())
.expire(`conn:${userId}`, LEASE_S)
.exec();
}
// Refresh leases for every locally connected user, in batches, every 20 s.
setInterval(async () => {
const pipe = redis.pipeline();
for (const userId of streams.keys()) {
pipe.hset(`conn:${userId}`, NODE_ID, Date.now()).expire(`conn:${userId}`, LEASE_S);
}
await pipe.exec();
}, 20_000);
The lease makes the registry self-healing: a node that crashes stops refreshing, and its entries expire within a minute. During that minute, publishes to the dead node are lost from the live path — which is acceptable because notifications are durable and the user’s reconnect replays them.
Step 3 — Publish to nodes, not users Permalink to this section
export async function publishToUser(userId, frameText) {
const nodes = await redis.hkeys(`conn:${userId}`);
if (!nodes.length) return; // offline: the table is enough
const payload = JSON.stringify({ userId, frame: frameText });
await Promise.all(nodes.map((n) => redis.publish(`node:${n}`, payload)));
}
Step 4 — One subscription per node, routed locally Permalink to this section
const sub = redis.duplicate();
await sub.subscribe(`node:${NODE_ID}`);
sub.on('message', (_channel, raw) => {
const { userId, frame } = JSON.parse(raw);
deliverLocal(userId, frame);
});
The same node-side routing in Go is a map guarded by a mutex and a non-blocking send into each stream’s buffered channel. A slow client’s full channel drops the frame for that client only — its reconnect replay repairs the gap — and never blocks delivery to the other users on the node:
// router.go — one Redis subscription per node, local fan-out by user id.
type Router struct {
mu sync.RWMutex
streams map[int64]map[chan []byte]struct{}
}
func (r *Router) Attach(user int64, ch chan []byte) func() {
r.mu.Lock()
if r.streams[user] == nil {
r.streams[user] = make(map[chan []byte]struct{})
}
r.streams[user][ch] = struct{}{}
r.mu.Unlock()
return func() {
r.mu.Lock()
delete(r.streams[user], ch)
if len(r.streams[user]) == 0 {
delete(r.streams, user)
}
r.mu.Unlock()
}
}
func (r *Router) Run(ctx context.Context, rdb *redis.Client, nodeID string) {
sub := rdb.Subscribe(ctx, "node:"+nodeID)
defer sub.Close()
for msg := range sub.Channel() {
var m struct {
UserID int64 `json:"userId"`
Frame string `json:"frame"`
}
if json.Unmarshal([]byte(msg.Payload), &m) != nil {
continue
}
r.mu.RLock()
for ch := range r.streams[m.UserID] {
select {
case ch <- []byte(m.Frame): // delivered to this stream's writer goroutine
default: // buffer full: skip, the client's replay will cover it
}
}
r.mu.RUnlock()
}
}
Each HTTP handler owns its channel, writes frames to the http.ResponseWriter and flushes, exactly as in implementing SSE with Go channels and Flusher. The router never touches a socket, which keeps a slow network write from holding the read lock.
Step 5 — Handle the transitions Permalink to this section
Two moments need care. When a user connects, register before replaying from the database, so any notification committed after the replay query finds the node. When the last stream for a user closes, remove the node from the hash so publishers stop sending to it:
req.on('close', async () => {
detach();
if (!streams.has(userId)) await redis.hdel(`conn:${userId}`, NODE_ID);
});
If a user’s stream moves between nodes during a deploy, they may briefly be registered on both. The duplicate delivery is harmless because the client store is keyed by notification id.
Validation & Monitoring Permalink to this section
# One subscription per node, however many users are connected.
redis-cli PUBSUB NUMSUB node:a node:b node:c
redis-cli PUBSUB CHANNELS 'user:*' | wc -l # expect 0 after the migration
# The registry agrees with reality for a test user.
redis-cli HGETALL conn:42
Track three numbers per node: local users, local streams, and routed frames per second. And one on the publisher: publishes that found no registered node. That last one should roughly equal notifications sent to offline users; a sudden rise means leases are expiring for connected users, usually because the refresh loop is blocked.
Production Checklist Permalink to this section
Frequently Asked Questions Permalink to this section
When is a channel per user good enough?
Up to a few thousand concurrent users per node, per-user channels are simple and perform well. Move to node-addressed channels when subscription churn shows up in broker CPU or deploys get slow.
Why not use a pattern subscription on every node?
It keeps subscription count low but makes every node receive every notification. At low volume that is fine; at high volume the broker's outbound traffic multiplies by the node count.
What happens to notifications published while a node is restarting?
They are lost from the live path but not from the database. The user's client reconnects to another node, which replays from the cursor, so the notification still arrives — a few seconds late.
Does this work with Redis Cluster?
Yes. Node channels are few, so they spread cleanly with sharded pub/sub, and the registry hashes are ordinary keys distributed across slots.