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 NUMSUB shows 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.

Three routing designs compared Matrix comparing per-user channels, a node-wide pattern subscription, and node-addressed channels with a presence registry across four criteria. Three routing designs compared Design Broker subscriptions Wasted delivery Restart cost Complexity Channel per user one per user none resubscribe all lowest Pattern user:* one per node every node gets all one call low Channel per node one per node none one call needs registry scales well acceptable problem at scale
Per-user channels are the simplest and scale with users. Node-addressed channels scale with nodes but need a registry of who is where.

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);
});
Node-addressed delivery A publisher looks up the node holding the user in the registry and publishes once to that node's channel; the node routes the frame to the user's two local streams. Node-addressed delivery notify(user 42) lookup conn:42 node:b channel one per node PUBLISH Tab on laptop stream 1 Phone app stream 2 route locally
The broker sees one subscription per node and one publish per node that holds the user — never a channel per user.

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.

Broker subscriptions as concurrent users grow Line chart of Redis subscription count against concurrent users for per-user channels and for node-addressed channels on a ten-node fleet. Broker subscriptions as concurrent users grow channel per user channel per node (10 nodes) 0 25 50 75 100 0 20 40 60 80 100 concurrent users (thousands) subscriptions (thousands)
Per-user channels grow linearly with users. Node-addressed channels stay at the node count.

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.