Waiting and pretending
BLPOP, MULTI with EXEC and DISCARD, and WATCH.
You will build
BLPOP, MULTI with EXEC and DISCARD, and WATCH
You will understand
how to block a thread without burning CPU, why per-connection state needs its own object, why Redis transactions have no rollback, and optimistic locking
- roughly 180 lines across two new classes
- Stages 20, 21, and 26
Where we are going
A client that waits for data that has not arrived yet:
# terminal A, blocks
$ redis-cli blpop waiting 5
# terminal B, a second later
$ redis-cli rpush waiting hello
Terminal A returns immediately with waiting and hello.
And a transaction that refuses to run if someone got there first:
127.0.0.1:6379> WATCH balance
OK
127.0.0.1:6379> MULTI
OK
127.0.0.1:6379(TX)> DECRBY balance 100
QUEUED
127.0.0.1:6379(TX)> EXEC
(nil) ← someone else changed balance, try again
flowchart LR
A[Part 6<br/>atomic commands] --> B[Stage 20<br/>BLPOP]
B --> C[Stage 21<br/>MULTI and EXEC]
C --> D[Stage 26<br/>WATCH]
D --> E([Part 8<br/>the server speaks first])
Stage 20, BLPOP, the command that waits
Goal. Block until a list has an element, or a timeout passes.
The idea
Every command until now answered immediately. This one holds the connection open and answers later, or never. That breaks an assumption baked in since Part 1.
Three problems appear together. How does a thread wait without burning CPU. How does a pushing thread tell it to stop waiting. What happens when two clients wait for the same key.
Waiting without spinning
The wrong version:
while (true) {
value = lpop(key);
if (value != null) return value;
}
Correct, and it pins a CPU core per waiting client. Five idle BLPOPs put the machine at 500% CPU
doing nothing.
The right primitive is Object.wait() with notifyAll(). A waiting thread parks, consuming
nothing, until another thread signals it. RPUSH does the signalling.
sequenceDiagram
participant A as Client A (BLPOP)
participant M as monitor
participant B as Client B (RPUSH)
A->>M: synchronized, list empty
A->>M: wait() releases the monitor, parks
B->>M: synchronized, push element
B->>M: notifyAll()
M-->>A: wakes
A->>A: re-check, element is there
Two details are easy to get wrong.
wait() must be inside synchronized, and it releases the monitor while parked. That release is
the whole point, because the pusher needs the monitor to notify and could never get it if the waiter
held it.
wait() must be in a loop, never an if. notifyAll wakes every waiter, and only one can have the
element, so the rest must re-check and sleep again. Java also permits spurious wakeups, where a
thread wakes with no notification at all, so the loop is required even with one waiter.
The code
// ponytail: one monitor for every list, so a push wakes all waiters and most go back to
// sleep. per-key monitors if a real workload ever makes the herd measurable
private final Object listMonitor = new Object();
// blocks until one of the keys has an element or the timeout passes.
// returns [key, value], or null on timeout. timeoutMillis of 0 waits forever
public String[] blockingPop(String[] keys, long timeoutMillis) throws InterruptedException {
long deadline = timeoutMillis == 0
? Long.MAX_VALUE
: System.currentTimeMillis() + timeoutMillis;
synchronized (listMonitor) {
while (true) {
// keys are checked in order, so an earlier key wins when both have data
for (String key : keys) {
String value = lpop(key);
if (value != null) {
return new String[]{key, value};
}
}
long remaining = deadline - System.currentTimeMillis();
if (remaining <= 0) {
return null;
}
// wait() releases the monitor, so a pushing thread can get in and notify us.
// the surrounding loop re-checks because notifyAll wakes every waiter
listMonitor.wait(Math.min(remaining, 100));
}
}
}
push gains three lines:
// a blocked BLPOP is waiting for exactly this
synchronized (listMonitor) {
listMonitor.notifyAll();
}
Wall-clock time is right here, unlike expiry. A timeout is a promise to a human, so a clock
correction should affect it. wait(millis) uses the system clock anyway.
Only one waiter can win. The element is taken by lpop inside the monitor, and lpop is itself
atomic, so the loser finds nothing and resumes waiting.
Timeouts and the reply
BLPOP key 0 block forever
BLPOP key 0.5 half a second, fractional seconds allowed
On timeout the reply is a null array, *-1\r\n, not an empty array *0\r\n. Clients distinguish
them. Null means nothing arrived. Empty would mean here is a list of nothing.
Run it
$ redis-cli rpush fruits apple
$ redis-cli blpop fruits 1
1) "fruits"
2) "apple"
$ redis-cli -t 5 blpop empty 0.5
(nil) # after half a second
Two terminals, the real test:
# terminal A
$ redis-cli -t 10 blpop waiting 5
# terminal B, a second later
$ redis-cli rpush waiting hello
A returns at once with waiting and hello.
Notice that
Terminal A used no CPU while waiting. Check with top if you like. That is wait() rather than a
spin loop, and it is the difference between a server and a space heater.
Try it yourself
- Replace
wait()with a busy loop and watch CPU usage while one client blocks. Undo it. - Change the
whilearoundwait()to anif. Start two blocked clients, push one element, and watch both wake. One returns a value. The other returnsnullimmediately instead of waiting out its timeout. - Run
BLPOP a b 0with data inbonly. It returnsb. Now put data in both andawins, because keys are checked in order.
What usually goes wrong
IllegalMonitorStateException. You called wait() outside synchronized.
The waiter never wakes. The pusher notifies a different monitor object, or does not notify.
Both waiters get the same element. The pop is outside the monitor.
Go deeper: how single-threaded Redis does this
Redis cannot block a thread on a client, because it has one thread for everything. Instead it puts the client in a per-key list of blocked clients and returns to its event loop. When a push happens, it checks that list and serves the longest-waiting client.
Our design blocks a thread and stays simple. Redis keeps the thread and takes on bookkeeping.
- Redis BLPOP, the blocking behaviour section
Object.waitJavadoc, and read the spurious wakeup paragraph
Stage 21, transactions
Goal. MULTI queues commands, EXEC runs them all, DISCARD throws them away.
The idea
Everything built so far is shared. One store, one dispatcher, every thread using both. A transaction belongs to one connection, and no other client may see it.
That forces a new object.
flowchart TD
S[(RedisStore)]
D[CommandDispatcher]
C1[ClientSession A]
C2[ClientSession B]
S ---|one per server, shared| D
D --- C1
D --- C2
C1 -.->|"MULTI queue<br/>never shared"| C1
C2 -.-> C2
If the queue lived on the dispatcher, two clients in MULTI at once would append to the same list
and EXEC would run each other’s commands.
The code
package com.example.redis;
import java.util.ArrayList;
import java.util.List;
// Per-connection state. Everything else in the server is shared, this is not.
public class ClientSession {
private boolean inTransaction;
private boolean aborted;
private final List<String[]> queued = new ArrayList<>();
public boolean inTransaction() {
return inTransaction;
}
public void begin() {
inTransaction = true;
aborted = false;
queued.clear();
}
public void queue(String[] command) {
queued.add(command);
}
// a command that failed to queue poisons the whole transaction
public void abort() {
aborted = true;
}
public boolean aborted() {
return aborted;
}
public List<String[]> drain() {
List<String[]> commands = new ArrayList<>(queued);
reset();
return commands;
}
public void reset() {
inTransaction = false;
aborted = false;
queued.clear();
}
}
Queuing intercepts dispatch before the switch, so it applies to every command including ones added later:
if (session.inTransaction() && !name.equals("EXEC") && !name.equals("DISCARD")) {
return queueForLater(name, command, session);
}
// inside MULTI nothing runs, it just piles up until EXEC
private byte[] queueForLater(String name, String[] command, ClientSession session) {
if (name.equals("MULTI")) {
return RespWriter.error("ERR MULTI calls can not be nested");
}
if (!KNOWN_COMMANDS.contains(name)) {
// redis refuses the whole transaction rather than skipping the bad command
session.abort();
return RespWriter.error("ERR unknown command '" + command[0] + "'");
}
session.queue(command);
return RespWriter.simpleString("QUEUED");
}
main creates one session per connection and passes it into execute.
Redis transactions are not database transactions
This surprises people, and it is the most important thing in the stage.
There is no rollback. If one command inside EXEC fails, the others still run:
SET word hello
MULTI
INCR word → +QUEUED
SET after ran → +QUEUED
EXEC → *2
-ERR value is not an integer or out of range
+OK
after is now set. Nothing was undone. Redis offers atomicity in the sense of isolation, meaning no
other client interleaves between the queued commands. It does not offer all-or-nothing.
Errors caught at queue time behave differently. An unknown command is rejected as it is queued and poisons the whole transaction:
MULTI
SET k v → +QUEUED
nope → -ERR unknown command 'nope'
EXEC → -EXECABORT Transaction discarded because of previous errors.
k is never set.
Errors Redis can see before running abort everything. Errors that only appear at runtime abort nothing.
Run it
127.0.0.1:6379> MULTI
OK
127.0.0.1:6379(TX)> SET k v
QUEUED
127.0.0.1:6379(TX)> INCR n
QUEUED
127.0.0.1:6379(TX)> EXEC
1) OK
2) (integer) 1
Notice that
redis-cli switched its prompt to (TX). It saw +OK from MULTI and is tracking the transaction
itself, which is a small sign your replies are shaped correctly.
Try it yourself
- Open two
redis-clisessions. StartMULTIin one and queueSET k v. In the other, runGET k. Still(nil), because nothing has run. - Queue a command, then
DISCARD, thenEXEC. You getERR EXEC without MULTI. - Reproduce the no-rollback behaviour with the
INCR wordexample. Convince yourself Redis is right to do this. What would rolling back a list push even mean?
What usually goes wrong
Two clients’ transactions mix. The queue is on the dispatcher rather than the session.
MULTI inside MULTI gets queued. The nesting check sits inside the switch, which the
interceptor never reaches.
Stage 26, WATCH
Goal. A transaction that gives up if someone else touched the key first.
This is stage 26 in the build order, arriving after the data types. It belongs here in the reading order, because it completes the transaction story.
The idea
Stage 21 gave you isolation. What it could not give you is a decision:
GET balance → 500 read it
← another client spends 400 here
MULTI
DECRBY balance 100
EXEC → runs anyway, balance is now negative
The read happened before the transaction, so the transaction cannot know it went stale. That is
check-then-act across a network, and no amount of isolation inside EXEC fixes it.
WATCH is Redis’s answer. Optimistic locking. Do not lock the key, just notice if it moved.
| Pessimistic lock | Optimistic with WATCH | |
|---|---|---|
| Others blocked | yes | no |
| Needs timeouts for dead clients | yes | no |
| Loser finds out | by waiting | by retrying |
| Cost | contention | a retry loop |
sequenceDiagram
participant A as Client A
participant S as Server
participant B as Client B
A->>S: WATCH balance
Note over S: remembers balance is at version 7
B->>S: SET balance 999
Note over S: balance is now version 8
A->>S: MULTI, DECRBY, EXEC
S-->>A: *-1 (nil), version moved, nothing ran
Detecting the change
The store must answer whether a key changed since you looked. Storing the old value would be wrong.
A key set to 5, changed to 9, then back to 5 did change, and a value comparison misses it. That
is the ABA problem, and the intermediate write may have been exactly what the client cared about.
So the store keeps a version counter per key, bumped on every write:
// bumped on every write, so WATCH can tell whether a key changed under it
private final Map<String, Long> versions = new ConcurrentHashMap<>();
public long version(String key) {
return versions.getOrDefault(key, 0L);
}
private void touch(String key) {
versions.merge(key, 1L, Long::sum);
}
Deletion bumps it. Creation bumps it too, because watching a key that does not exist yet and having someone create it must abort, which a null-to-null value comparison would also miss.
EXEC checks before running anything:
// optimistic locking: if anything we watched moved, the whole transaction is off
for (Map.Entry<String, Long> watched : session.watched().entrySet()) {
if (store.version(watched.getKey()) != watched.getValue()) {
session.reset();
return RespWriter.nullArray();
}
}
The maintenance hazard
Every mutating path must bump the version. Miss one and WATCH silently stops protecting against
that command. No error, no crash, just a transaction that runs when it should not. Silent, rare, and
it costs money in exactly the use case WATCH exists for.
The defence is a test that walks the write commands one by one:
@Test
void everyMutatingCommandInvalidatesAWatch() {
// if a command mutates without bumping the version, WATCH silently stops working
assertAborts("INCR", "k");
assertAborts("RPUSH", "k", "x");
assertAborts("APPEND", "k", "x");
assertAborts("SETNX", "k", "x");
assertAborts("XADD", "k", "1-1", "f", "v");
assertAborts("MSET", "k", "x");
}
Add a command that writes, add it to that list. It is the cheapest possible guard against a whole class of quiet bug. When hashes, sets, and sorted sets arrive in Part 9, extending this test catches nothing, which is the point.
Details that are easy to get wrong
Your own writes inside the transaction must not abort it. The version check happens once, at EXEC,
before any queued command runs. Otherwise the first SET in your own transaction would invalidate
your own watch.
EXEC and DISCARD both clear the watch, whether or not anything ran.
WATCH inside MULTI is an error rather than a queued command. It has to bypass the queuing
interceptor to say so. The natural implementation queues it and replies +QUEUED, which is wrong
and looks fine.
Run it
Two redis-cli sessions.
# session A
127.0.0.1:6379> SET k 1
127.0.0.1:6379> WATCH k
127.0.0.1:6379> MULTI
127.0.0.1:6379(TX)> SET k 2
QUEUED
# session B, before A runs EXEC
127.0.0.1:6379> SET k 999
OK
# back in session A
127.0.0.1:6379(TX)> EXEC
(nil)
Without B’s write, EXEC returns the results normally.
Notice that
(nil) is redis-cli rendering the null array. That is deliberately not an error. It means try
again, rather than you did something wrong. A client’s retry loop switches on exactly this.
The retry loop this enables
while (true) {
WATCH balance
current = GET balance
if (current < 100) { UNWATCH; fail }
MULTI
DECRBY balance 100
if (EXEC != nil) break // someone else won, go round again
}
That is compare-and-swap over a network, and it is the whole reason WATCH exists.
Try it yourself
- Run the two-session experiment above. Then run it again without B’s write and see
EXECsucceed. - Watch a key that does not exist, have another client create it, then
EXEC. It aborts, because creation counts as a change. - Remove the
touch(key)call fromincrementand runeveryMutatingCommandInvalidatesAWatch. One assertion fails, and it names the command.
What usually goes wrong
EXEC always aborts. You are checking the version after the queued commands run, so your own
writes invalidate your own watch.
WATCH inside MULTI returns +QUEUED. It is not excluded from the queuing interceptor.
What you built
A command that waits, transactions that queue, and optimistic locking that turns a queue into a conditional one.
Checkpoint
- Why must
wait()release the monitor? - Why must
wait()sit in a loop rather than anif? - Why must the transaction queue live on the session rather than the dispatcher?
- In what sense is a Redis transaction atomic, and in what sense is it not?
- Why a version counter rather than remembering the value?
Resources
- Redis transactions, the no-rollback rationale from the source
- Redis BLPOP and Redis WATCH
Object.waitJavadoc- ABA problem
Next
The server sends you something you did not ask for.