The server speaks first

Pub/sub, streams, and a blocking XREAD.

  • Part 8
  • intermediate
  • about 80 minutes

You will build

pub/sub, streams, and a blocking XREAD

You will understand

what breaks when two threads write to one socket, why a subscribed connection is restricted, how stream IDs order themselves, and one trap that makes a blocking read never return

  • roughly 300 lines across four new classes
  • Stages 22, 23, and 27

Where we are going

A message arriving on a connection that asked for nothing:

# terminal A
$ redis-cli subscribe news
1) "subscribe"
2) "news"
3) (integer) 1

# terminal B
$ redis-cli publish news hello
(integer) 1

# terminal A prints, on its own
1) "message"
2) "news"
3) "hello"
flowchart LR
    A[Part 7<br/>blocking and transactions] --> B[Stage 22<br/>pub/sub]
    B --> C[Stage 23<br/>streams]
    C --> D[Stage 27<br/>blocking XREAD]
    D --> E([Part 9<br/>four more types])

Stage 22, publish and subscribe

Goal. A message published by one client reaches every subscriber.

The idea

Every reply so far answered a command on the same connection, written by the thread that read it. This breaks both halves of that.

The message arrives on A’s socket with no command from A. And it is written by B’s thread, reaching across into A’s connection.

That second point is the engineering problem. Two threads can now write to one socket: A’s own thread replying to A’s commands, and any publisher’s thread pushing messages. Interleave those mid-reply and the client receives corrupted RESP, half a +PONG with a *3 inside it, which desynchronises the connection permanently.

flowchart TD
    subgraph "one socket, two writers"
        T1[A's own thread<br/>replying to A] --> S[ClientSession.send<br/><i>synchronized</i>]
        T2[B's thread<br/>publishing] --> S
        S --> SOCK[A's socket]
    end

The fix is small and mandatory. All writes to a socket go through one synchronized method. ClientSession.send owns the stream, and the connection loop stops writing directly.

The code

ClientSession grows from a transaction queue into the connection’s identity:

private final Set<String> subscriptions = new LinkedHashSet<>();
private final OutputStream out;

public ClientSession(OutputStream out) {
    this.out = out;
}

// tests and anything that never pushes to the client
public ClientSession() {
    this(new ByteArrayOutputStream());
}

public Set<String> subscriptions() {
    return subscriptions;
}

public boolean isSubscribed() {
    return !subscriptions.isEmpty();
}

// publishers write here from their own threads, so one writer at a time
public synchronized void send(byte[] bytes) throws IOException {
    out.write(bytes);
    out.flush();
}

And a registry mapping channel to sessions:

package com.example.redis;

import java.io.IOException;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;

// Channel registry. Publishers write straight into subscribers' sockets from their own thread.
public class PubSub {

    private final Map<String, Set<ClientSession>> channels = new ConcurrentHashMap<>();

    // returns how many channels this session is now subscribed to
    public int subscribe(ClientSession session, String channel) {
        channels.computeIfAbsent(channel, k -> ConcurrentHashMap.newKeySet()).add(session);
        session.subscriptions().add(channel);
        return session.subscriptions().size();
    }

    public int publish(String channel, String message) {
        Set<ClientSession> subscribers = channels.get(channel);
        if (subscribers == null) {
            return 0;
        }

        byte[] payload = RespWriter.array(
                RespWriter.bulkString("message"),
                RespWriter.bulkString(channel),
                RespWriter.bulkString(message));

        int delivered = 0;
        for (ClientSession subscriber : subscribers) {
            try {
                subscriber.send(payload);
                delivered++;
            } catch (IOException e) {
                // subscriber went away mid-publish, drop it rather than failing the publish
                subscribers.remove(subscriber);
            }
        }
        return delivered;
    }

    // a disconnecting client must not stay in any channel
    public void removeAll(ClientSession session) {
        for (String channel : session.subscriptions()) {
            Set<ClientSession> subscribers = channels.get(channel);
            if (subscribers != null) {
                subscribers.remove(session);
            }
        }
        session.subscriptions().clear();
    }
}

ConcurrentHashMap.newKeySet() for the members, because publishers and subscribers touch the set from different threads.

Cleanup is not optional

A client that disconnects while subscribed leaves a dead session in every channel it joined. Publishing then writes to a closed socket forever, and the set grows without bound.

Two defences, and you need both. onDisconnect in the connection loop removes the session from all its channels, and it has to run even when the client crashed. publish drops a subscriber whose send throws, because the disconnect may not have been noticed yet and a write is how you find out.

The first one forced a change in handleClient. The session variable moved outside the try:

// declared out here so the disconnect cleanup can still see it after a failure
ClientSession session = null;

try (Socket socket = clientSocket) {
    ...
    session = new ClientSession(outputStream);
    ...
} catch (IOException e) {
    ...
}

if (session != null) {
    dispatcher.onDisconnect(session);
}

Cleanup that only runs on the happy path is not cleanup.

Subscribed mode

Once a connection subscribes, Redis restricts it:

allowed:  SUBSCRIBE  UNSUBSCRIBE  PING  QUIT  RESET
anything else → -ERR Can't execute 'get': only (P|S)SUBSCRIBE / ... are allowed in this context

The reason is protocol clarity. A subscriber’s stream is a sequence of pushed messages, and mixing ordinary replies into it makes a client unable to tell which is which. RESP3 fixed this with a separate push type. RESP2 solves it by forbidding the mix.

The restriction is per connection, not global. Other clients carry on normally, and the test otherClientsAreNotRestricted pins that, because it is an easy thing to accidentally make server-wide.

One command, several replies

SUBSCRIBE a b sends two confirmation arrays, not one array of two. Each carries the running count:

*3  subscribe  a  1
*3  subscribe  b  2

That does not fit RespWriter.array, which builds one value. Concatenating pre-encoded replies is the honest answer, and it keeps the writer a RESP encoder rather than growing a special case for one command’s quirk.

Run it

Two terminals.

# terminal A
$ redis-cli -t 30 subscribe news
# terminal B
$ redis-cli publish news hello
(integer) 1

Terminal A prints the message as it arrives. Then confirm a third terminal is unaffected:

$ redis-cli set k v
OK

Notice that

PUBLISH returned 1, the number of subscribers it reached. That number is the entire delivery guarantee. Redis pub/sub is fire and forget. A message published while a subscriber is disconnected is gone. No queue, no replay, no acknowledgement.

Say that out loud in a tutorial, because people reach for pub/sub expecting a message queue. Streams are the durable option, and they are next.

Try it yourself

  1. Subscribe in one terminal, then try GET k in that same session. Read the error. Now run PING in it, which is allowed.
  2. Kill a subscriber with Ctrl-C, then publish. The count drops to 0 rather than throwing.
  3. Subscribe two clients to the same channel and publish once. Both receive identical bytes.

What usually goes wrong

A subscriber receives corrupted RESP. Two threads wrote to the socket without synchronising.

Publishing to a channel nobody left grows forever. onDisconnect is inside the try and never runs on a crash.


Stage 23, streams

Goal. An append-only log with ordered IDs. XADD, XRANGE, XLEN, XREAD.

The idea

Lists are stored as List<String>. If a stream were List<StreamEntry>, instanceof List would be true for both, TYPE could not tell them apart, and LPUSH on a stream would type-check and corrupt it.

So a stream gets its own class rather than a bare collection. The lesson generalises. When two values share a Java type but not a domain meaning, the domain needs its own type. This is where Object in Entry starts paying rent.

IDs are the whole design

1757012345678-0
└── milliseconds ┘└ seq

Two parts, because many entries can land in the same millisecond. Ordering is by millisecond first and sequence second, which is exactly compareTo on the pair. That is why StreamId is a Comparable record rather than a string.

Never compare IDs as strings. "10-0" < "9-0" lexicographically, which is backwards.

package com.example.redis;

// A stream id is milliseconds plus a sequence number, so many entries can share a millisecond.
public record StreamId(long ms, long seq) implements Comparable<StreamId> {

    public static final StreamId MIN = new StreamId(0, 0);
    public static final StreamId MAX = new StreamId(Long.MAX_VALUE, Long.MAX_VALUE);

    // "5-3" exact, or "5" with the sequence defaulted, start of a range wants 0, end wants max
    public static StreamId parse(String raw, long defaultSeq) {
        int dash = raw.indexOf('-');
        if (dash < 0) {
            return new StreamId(Long.parseLong(raw), defaultSeq);
        }
        return new StreamId(
                Long.parseLong(raw.substring(0, dash)),
                Long.parseLong(raw.substring(dash + 1)));
    }

    @Override
    public int compareTo(StreamId other) {
        int byMs = Long.compare(ms, other.ms);
        return byMs != 0 ? byMs : Long.compare(seq, other.seq);
    }

    @Override
    public String toString() {
        return ms + "-" + seq;
    }
}

Three ways to specify one:

5-3      exact
5-*      millisecond fixed, sequence auto: last+1 if same ms, else 0
*        both auto: now, and seq 0, or last+1 within the same millisecond

The two rules XADD enforces

XADD s 0-0 …    →  -ERR The ID specified in XADD must be greater than 0-0
XADD s 4-0 …    →  -ERR The ID specified in XADD is equal or smaller than the target stream top item

IDs must strictly increase. That is what makes a stream an append-only log with a meaningful cursor. A consumer that has seen 5-5 can ask for everything after it and know nothing will be inserted behind its back.

One consequence worth noticing. Since 0-0 is forbidden, XADD s 0-* must produce 0-1, not 0-0. A special case, and a good test.

XRANGE and XREAD differ in one bit

XRANGE s 1-1 2-1     inclusive at both ends    browse history
XREAD  STREAMS s 1-1 exclusive at the start    follow a cursor

The same scan, one boolean apart. The distinction is the consumer’s mental model. A cursor you already consumed must not repeat.

XRANGE also takes shorthands. - is the smallest possible ID, + the largest, and a bare 5 means 5-0 as a start and 5-<max seq> as an end. That is why StreamId.parse takes a default sequence rather than assuming zero.

Reply shapes

Deeply nested, and the first place the writer’s design really pays off:

XRANGE →  *2
            *2  $3 "1-1"  *2 $1 "a" $1 "1"
            *2  $3 "2-1"  *2 $1 "a" $1 "2"

XREAD  →  *1
            *2  $1 "s"  <the XRANGE shape above>

Because array() takes already-encoded byte[] elements, this is composition rather than a recursive encoder. Streams with nothing new are omitted from XREAD entirely, and if no stream has anything, the reply is a null array.

Run it

$ redis-cli xadd s 1-1 temperature 36
"1-1"
$ redis-cli xadd s '*' temperature 37
"1788635051406-0"
$ redis-cli xrange s - +
1) 1) "1-1"
   2) 1) "temperature"
      2) "36"
2) 1) "1788635051406-0"
   2) 1) "temperature"
      2) "37"
$ redis-cli xlen s
(integer) 2
$ redis-cli type s
stream
$ redis-cli xadd s 1-1 a 1
(error) ERR The ID specified in XADD is equal or smaller than the target stream top item

Notice that

redis-cli rendered the nesting as an indented tree. If it prints structure rather than an error, your encoding is right.

Try it yourself

  1. Add three entries in the same millisecond using 5-* three times. The sequence climbs 0, 1, 2.
  2. Run XADD s 0-* on a fresh stream. You get 0-1, because 0-0 is forbidden.
  3. Compare XRANGE s 1-1 + with XREAD STREAMS s 1-1. The first includes 1-1, the second does not.

Stage 27, blocking XREAD

Goal. XREAD BLOCK 5000 STREAMS s $ waits for the next entry.

The idea

The waiting machinery already exists, because BLPOP built it in Part 7. What changes is who signals:

flowchart LR
    RPUSH --> M[one monitor]
    XADD --> M
    M --> W1[wakes blocked BLPOP]
    M --> W2[wakes blocked XREAD]

One monitor for both, so listMonitor became dataMonitor and XADD gained the same notifyAll that push had.

The loop moved, though. BLPOP waits inside the store. Blocking XREAD waits in the dispatcher, because it has to re-poll several streams and assemble a reply. The store just exposes awaitWrite, and the caller decides what done means. That is the better split, since the store does not need to know what a client is waiting for.

The $ trap

$ means entries added after this call. The obvious implementation re-resolves it on each retry:

loop:
    after = lastId(stream)      ← re-read every time
    entries = read after `after`
    if empty: wait, retry

That never returns. Each retry moves the cursor to whatever just arrived, so the new entry is always before the cursor by the time you look. A blocking read that can never succeed, with no error to explain it.

$ has to be pinned once, before the loop, to the last ID at call time. Then it is an ordinary exclusive cursor and the normal path works.

This is the kind of bug that only shows up in the blocking case. Non-blocking XREAD $ returns empty immediately and looks correct.

Why the internal wait is capped

awaitWrite never parks for longer than a tenth of a second, even when the caller asked for five. The outer loop re-checks and waits again.

That is deliberate insurance. A missed notification, a write that lands in the gap between the poll and the wait, would otherwise hang the client for the full timeout. Capping the wait turns a lost wakeup into at most 100ms of extra latency instead of a five-second stall.

Correct code would not need it. Code that might have a subtle race, in a language where that race is invisible until production, benefits from the belt.

Run it

Two terminals.

# terminal A, blocks
$ redis-cli -t 10 xread block 5000 streams s '$'
# terminal B, a second later
$ redis-cli xadd s '*' temperature 37

A returns immediately with the new entry. Run A alone and it returns (nil) after five seconds.

Notice that

A saw only the entry added after it started blocking. The existing entries in the stream were not returned, which is what $ means and what the pinning makes work.

Try it yourself

  1. Re-resolve $ inside the loop instead of before it. The blocking read now never returns, with no error. Undo it. This is the bug the stage exists to teach.
  2. Use BLOCK 0 and leave it running. Push an entry a minute later and it still wakes.
  3. Compare XREAD BLOCK 100 STREAMS s $ against XREAD STREAMS s $. The non-blocking one returns (nil) instantly and looks identical to a timeout.

What usually goes wrong

The blocking read never returns. $ is being re-resolved on each retry.

It returns instantly with old entries. You pinned $ to 0 rather than to the current last ID.

Go deeper: what streams are for, and what we skipped

Streams are Redis’s durable log. Unlike pub/sub, entries persist, and a consumer can replay history from any ID.

Skipped here: consumer groups (XGROUP, XREADGROUP, XACK), which turn a stream into a work queue with per-consumer acknowledgement and pending-entry tracking. Also XDEL, XTRIM, and MAXLEN.


What you built

A server that pushes without being asked, a durable log with ordered IDs, and a read that follows that log as it grows.

Checkpoint

  1. What exactly breaks if two threads write to one socket without synchronising?
  2. Why must onDisconnect run even when the client’s thread threw?
  3. Why does Redis restrict what a subscribed connection may do?
  4. Why is a stream ID two numbers rather than one timestamp?
  5. Why does re-resolving $ on each retry make a blocking read never return?

Resources

Next

Hashes, sets, sorted sets, and the command that replaces the dangerous one you wrote in Part 6.

Part 9: Four More Types →