The server speaks first
Pub/sub, streams, and a blocking XREAD.
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
- Subscribe in one terminal, then try
GET kin that same session. Read the error. Now runPINGin it, which is allowed. - Kill a subscriber with Ctrl-C, then publish. The count drops to 0 rather than throwing.
- 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
- Add three entries in the same millisecond using
5-*three times. The sequence climbs 0, 1, 2. - Run
XADD s 0-*on a fresh stream. You get0-1, because0-0is forbidden. - Compare
XRANGE s 1-1 +withXREAD STREAMS s 1-1. The first includes1-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
- 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. - Use
BLOCK 0and leave it running. Push an entry a minute later and it still wakes. - Compare
XREAD BLOCK 100 STREAMS s $againstXREAD 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
- What exactly breaks if two threads write to one socket without synchronising?
- Why must
onDisconnectrun even when the client’s thread threw? - Why does Redis restrict what a subscribed connection may do?
- Why is a stream ID two numbers rather than one timestamp?
- 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.