---
title: Publish-Subscribe
source: 'https://academia.sh/en/courses/in-memory-stores/publish-subscribe'
course: 'In-Memory Stores and Caching Systems'
language: en
updated: '2026-08-23T07:00:43+00:00'
license: 'CC BY-SA 4.0'
---

# Publish-Subscribe

The in-memory store's lightweight messaging path and its lack of a delivery guarantee: counting the loss of a message published with no subscriber present, the bytes a slow subscriber's output buffer holds and the threshold at which it gets dropped, and the memory the same event stream pays when built with a persistent structure instead.

The previous lesson ran collision control on a single key: whether the value a client read had
been changed by someone else was seen through the watch, and a colliding transaction was
retried. There, communication was established indirectly, between two clients touching the same
key. A client saying something **directly** to another client — the store used as a transmission
channel — was never covered.

In-memory stores carry a separate path for this: a client **publishes** to a **channel**, and
clients subscribed to that channel receive the message. Nothing gets written to the key space.
What this lesson measures is not that path's speed, but **what it does not guarantee**.

## The Store's Publish Path

In the library system, loan events land on branch boards instantly. The mechanism below sets up
three branch channels, one board per branch, and an audit subscriber that joins at step three
hundred. Board connectivity drops on a deterministic schedule; instead of a real network outage,
a rule based on the step number is used, so the numbers come out the same on every run (**CU10**).

```js
// publish-loss.mjs — a message published to a channel with no subscriber is lost; the loss is counted
class Store {                                       // publish-subscribe: no record, no rewind
  constructor() { this.channels = new Map(); }      // channel -> Set(subscriber)
  subscribe(channel, a) {
    if (!this.channels.has(channel)) this.channels.set(channel, new Set());
    this.channels.get(channel).add(a);
  }
  unsubscribe(channel, a) {
    this.channels.get(channel)?.delete(a);
  }
  publish(channel, message) {                        // return value: number of subscribers reached
    const set = this.channels.get(channel);
    if (!set || set.size === 0) return 0;             // no subscriber: message is written nowhere
    for (const a of set) a.receive(message);
    return set.size;
  }
}

const events = (count) => {                          // deterministic workload, seed visible
  let seed = 20250731;
  const rand = () => ((seed = (seed * 1103515245 + 12345) % 2147483648) / 2147483648);
  return Array.from({ length: count }, (_, i) => ({
    step: i + 1, branch: Math.floor(rand() * 3) + 1, bookId: Math.floor(rand() * 400) + 1,
  }));
};

const store = new Store();
const channel = (s) => `branch:${s}:loans`;
const boards = [1, 2, 3].map((s) => ({ name: `board${s}`, branch: s, received: 0, receive() { this.received += 1; } }));
const audit = { name: "audit", received: 0, receive() { this.received += 1; } };
const counters = [1, 2, 3].map(() => ({ published: 0, delivered: 0, unsubscribed: 0 }));

// Board connectivity: each board disconnects periodically (deterministic schedule, no randomness).
const boardOpen = (s, step) => (step + s * 37) % 100 >= 25;   // connected in 75% of steps
for (const e of events(600)) {
  for (const b of boards) {                          // connection state translates to subscription
    const open = boardOpen(b.branch, e.step);
    if (open) store.subscribe(channel(b.branch), b); else store.unsubscribe(channel(b.branch), b);
  }
  if (e.step === 301) for (const s of [1, 2, 3]) store.subscribe(channel(s), audit);  // late-joining subscriber

  const c = counters[e.branch - 1];
  const reached = store.publish(channel(e.branch), e);
  c.published += 1; c.delivered += reached;
  if (reached === 0) c.unsubscribed += 1;
}

console.log(["channel", "published", "delivered", "unsubscribed publish", "missed by board"]
  .map((h, i) => (i === 0 ? h.padEnd(16) : h.padStart(22))).join(""));
let total = [0, 0, 0, 0];
for (const [i, c] of counters.entries()) {
  const missed = c.published - boards[i].received;
  total = total.map((t, j) => t + [c.published, c.delivered, c.unsubscribed, missed][j]);
  console.log(channel(i + 1).padEnd(16) +
    [c.published, c.delivered, c.unsubscribed, missed].map((n) => String(n).padStart(22)).join(""));
}
console.log("total".padEnd(16) + total.map((n) => String(n).padStart(22)).join(""));
console.log(`received by the late subscriber: ${audit.received}  (published before it joined: 300)`);
```

```
channel                      published             delivered  unsubscribed publish       missed by board
branch:1:loans                     199                   266                    21                    46
branch:2:loans                     195                   227                    31                    55
branch:3:loans                     206                   260                    22                    46
total                              600                   753                    74                   147
received by the late subscriber: 300  (published before it joined: 300)
```

Three numbers need reading. The **delivered** column is larger than the published count
(753 > 600), because a message gets copied once for every subscriber present, and the publisher
does not know that count in advance. The **unsubscribed publish** column says that 74 messages
were written nowhere at all: no one was on the channel at that moment, and the message vanished
at the step it was produced. The **missed by board** column is larger still (147), because when
one board is disconnected but another subscriber is connected, the message is not lost — it just
never reaches that board.

The last line is the sharpest one. The audit subscriber joined at step three hundred and
received exactly 300 messages; it has no way to see the 300 that came before it. No trace of
those messages remains anywhere in the store's key space. In the publish-subscribe path,
delivery reduces to whether the moment a message is published and the moment a subscriber is
connected happen to overlap.

## The Slow Subscriber's Buffer

Loss does not come only from a broken connection. When a subscriber is connected but messages
are produced faster than it reads them, the unread messages pile up in an **output buffer** the
store holds. That buffer eats memory, and it is invisible in the publisher's accounting. The run
below drives the subscriber at different read speeds and counts the bytes the buffer holds. Once
the buffer exceeds 32,768 bytes, the store drops the subscriber; for a dropped subscriber, every
publish after that is a loss.

```js
// slow-subscriber.mjs — a slow subscriber's output buffer: bytes held and the step it hits the limit
const MESSAGE = (step, branch, bookId) =>          // body carried verbatim; its bytes are measured
  JSON.stringify({ event: "loan", step, branch, bookId });

function run(readInterval, byteLimit, count = 2000) {
  let seed = 20250731;
  const rand = () => ((seed = (seed * 1103515245 + 12345) % 2147483648) / 2147483648);
  const buffer = [];                               // messages not yet written to the subscriber
  let bytes = 0, peak = 0, dropStep = 0, afterDrop = 0, read = 0;
  for (let step = 1; step <= count; step += 1) {
    const g = MESSAGE(step, Math.floor(rand() * 3) + 1, Math.floor(rand() * 400) + 1);
    if (dropStep) { afterDrop += 1; continue; }      // dropped subscriber: publish is a loss
    buffer.push(g); bytes += Buffer.byteLength(g);   // the store writes to the buffer first
    peak = Math.max(peak, bytes);
    if (bytes > byteLimit) { dropStep = step; continue; }   // limit exceeded: subscriber is dropped
    if (step % readInterval === 0) {                 // subscriber reads once every `readInterval` steps
      const c = buffer.shift();
      if (c !== undefined) { bytes -= Buffer.byteLength(c); read += 1; }
    }
  }
  return { readInterval, read, peak, pending: buffer.length, dropStep, afterDrop };
}

console.log(`message body: ${MESSAGE(1, 1, 1)}  (${Buffer.byteLength(MESSAGE(1, 1, 1))} bytes)`);
console.log(["read interval", "read", "peak bytes", "pending in buffer", "drop step", "after drop"]
  .map((h, i) => (i === 0 ? h.padEnd(15) : h.padStart(19))).join(""));
for (const interval of [1, 2, 4, 10]) {
  const r = run(interval, 32768);
  console.log(`every ${String(interval).padStart(2)} steps`.padEnd(15) +
    [r.read, r.peak, r.pending, r.dropStep, r.afterDrop]
      .map((n) => String(n).padStart(19)).join(""));
}
```

```
message body: {"event":"loan","step":1,"branch":1,"bookId":1}  (47 bytes)
read interval                 read         peak bytes  pending in buffer          drop step         after drop
every  1 steps                2000                 52                  0                  0                  0
every  2 steps                 639              32800                641               1280                720
every  4 steps                 214              32777                646                860               1140
every 10 steps                  71              32801                647                718               1282
```

A subscriber that keeps up adds 52 bytes to memory: the buffer never holds more than one
message. When the read interval grows to two, the same subscriber reaches 32,800 bytes and gets
dropped at step 1,280; the remaining 720 publishes do not exist for it. When the interval grows
to ten, the limit is hit at step 718, and the loss climbs to 1,282 messages.

The budget relationship here runs in reverse. A subscriber's **slowness** is paid out of the
store's **memory**, and that cost is per subscriber: a hundred slow subscribers with the same
limit means a hundred separate buffers, up to 3.2 MB in the worst case. The limit is the store's
way of protecting itself; its cost is that a subscriber that hits the limit gets dropped
silently. The publisher sees none of this, because the publish call does not wait for the
message to be read (**CU11**).

## Persistent Versus Not

The same event stream can also be built with the store's **stream** structure. The difference is
a single point: a stream stores every entry with an id, so a subscriber that arrives late or
reconnects after dropping can read from where it left off. The price of that is memory. The run
below passes 2,000 events through three paths; the subscriber drops at step 800 and returns at
step 1,500. A 16-byte pointer/length overhead per entry is assumed (**CU12**).

```js
// stream-comparison.mjs — same event stream: bytes for a non-persistent publish vs. a trimmed stream
const OVERHEAD = 16;                                // per-entry id/length overhead (my own accounting)
const MESSAGE = (step, branch, bookId) => JSON.stringify({ event: "loan", step, branch, bookId });
const DROP = 800, RETURN = 1500, COUNT = 2000;      // subscriber drops at 800, returns at 1500

function events(count) {
  let seed = 20250731;
  const rand = () => ((seed = (seed * 1103515245 + 12345) % 2147483648) / 2147483648);
  return Array.from({ length: count }, (_, i) =>
    MESSAGE(i + 1, Math.floor(rand() * 3) + 1, Math.floor(rand() * 400) + 1));
}

function pubsubPath() {                             // non-persistent: nothing but the channel record
  let bytes = 0;
  for (const name of ["branch:activity", "board"]) bytes += Buffer.byteLength(name) + OVERHEAD;
  return { name: "publish-subscribe", returnBytes: bytes, endBytes: bytes, entries: 0, recovered: 0 };
}

function streamPath(bodies, trim) {                  // persistent: entry id and body sit in memory
  const stream = [];
  let bytes = 0, returnBytes = 0, recovered = 0;
  const size = (g) => Buffer.byteLength(g.id) + Buffer.byteLength(g.body) + OVERHEAD;
  bodies.forEach((body, i) => {
    const entry = { id: `${i + 1}-0`, body };
    stream.push(entry); bytes += size(entry);
    while (trim && stream.length > trim) bytes -= size(stream.shift());   // oldest entry drops
    if (i + 1 === RETURN) {                          // subscriber returns at exactly this step
      returnBytes = bytes;
      recovered = stream.filter((g) => {              // only entries still in the window are recoverable
        const step = Number(g.id.split("-")[0]);
        return step >= DROP && step < RETURN;
      }).length;
    }
  });
  return { name: trim ? `stream (trim=${trim})` : "stream (no trim)",
           returnBytes, endBytes: bytes, entries: stream.length, recovered };
}

const bodies = events(COUNT);
const gap = RETURN - DROP;                           // events the subscriber was disconnected for
console.log(`events=${COUNT}  gap=${gap} events (steps ${DROP}-${RETURN - 1})`);
console.log(["path", "bytes@1500", "bytes@2000", "entries", "recovered", "lost"]
  .map((h, i) => (i === 0 ? h.padEnd(20) : h.padStart(13))).join(""));
for (const r of [pubsubPath(), streamPath(bodies, 500), streamPath(bodies, 0)]) {
  console.log(r.name.padEnd(20) +
    [r.returnBytes, r.endBytes, r.entries, r.recovered, gap - r.recovered]
      .map((n) => String(n).padStart(13)).join(""));
}
```

```
events=2000  gap=700 events (steps 800-1499)
path                   bytes@1500   bytes@2000      entries    recovered         lost
publish-subscribe              52           52            0            0          700
stream (trim=500)           36865        36862          500          499          201
stream (no trim)           108381       145243         2000          700            0
```

The three rows show this course's rule exactly as it is. The publish-subscribe path holds 52
bytes; that is nothing more than the channel name and the subscription record, and it never
grows with event count. In exchange, all 700 of the dropped subscriber's events are lost.

The trimmed stream holds 36,865 bytes and returns 499 of the gap; the remaining 201 events are
lost because they were trimmed away before the subscriber returned. The trim length is directly
a **recovery window** here: 500 entries means 500 events of lookback.

The untrimmed stream loses nothing, but holds 108,381 bytes at step 1,500 and 145,243 at step
2,000. The growth between those two steps is the real warning: an untrimmed stream's memory
grows without bound as events accumulate, and once a memory limit kicks in, that growth forces
the eviction of other data.

## Where Publishing Fits

These measurements do not make the publish-subscribe path worthless; they set its place. Work
where loss is acceptable is work where missing a message is **correctable**: a board's live
counter, a cache-invalidation notice, an announcement of a configuration change. In every one of
these, a subscriber makes up what it missed on the next full read.

For work where loss is not acceptable, the measured cost gets paid: a loan record posting to
accounting, a late fee being processed, the record of a book returning to the shelf. This work
either gets written to a persistent stream or handed off outside the store, to a structure that
gives a delivery guarantee. The distinction is not made by asking "which is faster" but by
asking "can a missed message be recovered."

## Summary

- The publish-subscribe path does not write the message to the key space: 74 of 600 publishes
  vanished without reaching anywhere because no subscriber was on the channel at that moment,
  and boards missed 147 messages besides.
- The delivered count is independent of the published count; 600 publishes produced 753
  delivered copies, and the publisher does not know that number in advance.
- A late-joining subscriber cannot see any of the 300 messages published before it joined; the
  history sits nowhere in the store.
- The cost of a slow subscriber is the store's memory: when the read interval grew to two, the
  output buffer reached 32,800 bytes, the subscriber was dropped at step 1,280, and the next 720
  publishes were lost for it.
- When the same stream is built with a persistent structure, it holds 36,865 bytes (trim 500) or
  145,243 bytes (no trim) instead of 52; what that buys is 499 and 700 recoverable events
  respectively.

## Next Step

The lessons so far measured the store's abilities one at a time: replication, failover,
partitioning, collision control, and now publishing. Each was defined by its own cost, but none
of them answered end to end the question "which structure does this job get built with in this
store, and how many bytes does it hold." The next lesson takes on four typical jobs from the
library system — the catalog cache, session records, a rate limit, and a popular-book
leaderboard — one at a time; it fits each one to a structure and gathers each one's memory
accounting into the same table.
