Skip to content

Soak and failover

Half a million messages a second over a million partitions for hours on three nodes, with deduplication and retention on, through kill -9 and graceful restarts of the leader and a kill -9 of a follower: steady latency, what each fault looked like, memory and disk.

Updated View as Markdown

A minute at full speed shows a ceiling. A broker you trust has to hold its latency for hours, with retention deleting underneath it, while its nodes die. So on the night before 2.0.0-beta.1 we ran Queen at half a million messages a second, in three runs that lasted from 1 h 40 min to 3 h 23 min, and killed a node every ninety minutes.

The shape

1,000 queues of 1,000 partitions, one million partitions in all, created before the run. Nine load processes on three load machines offered 500,000 msg/s, three of them attached to each broker for the whole run. Messages were 256-byte JSON, pushed 100 at a time to one partition. 1,008 consumers popped with the client’s default sizing, where the leader picks how many partitions a pop claims, and acked asynchronously. Deduplication was on with a 60-second window, and completed messages were deleted after 300 seconds, so retention ran the whole time.

An agent on each broker injected a fault ten minutes in and every ninety minutes after that: a kill -9 of the leader, a SIGTERM of the leader, a kill -9 of a follower, in turn. It restarted the node 30 seconds later, 35 after a SIGTERM. The harness is soak.sh in benchmark-queen/2026-09-30-stage03-3node/, and the runs are runs/soak1 to runs/soak5.

Three runs

Run Build Length Steady windows e2e p50 / p99 Push p99 Faults
soak3 the night’s build after the reclaim fix 1 h 45 min 608 of 629 19 / 91 ms 29 ms kill -9 of the leader, SIGTERM of the leader
soak4 the next one, with the review fixes 3 h 23 min 1,187 of 1,219 20 / 86 ms 28 ms kill -9 of the leader, SIGTERM of the leader, kill -9 of a follower
soak5 7a2a1711, 2.0.0-beta.1 1 h 40 min 589 of 602 22 / 101 ms 30 ms kill -9 of the leader

All three ran between 2026-09-30 and 2026-10-01, and every build of the night is part of commit 7a2a1711. The latencies are the median over the steady 10-second windows: the median load process for p50, the worst for p99. A steady window has no failed request and an e2e p99 under one second; the others are the faults and the minute or two after each.

The night also showed what a soak is for. The run before these, soak2, held a steady p99 of about 350 ms with a push p99 of 175 ms, because every 5 seconds the retention pass rewrote log files while holding the log’s write lock, and appends queued behind it. Moving the rewrite off the lock, now part of 7a2a1711, took the steady p99 to about 90 ms and the push p99 to 28 ms.

What a fault looks like

Take soak4, the longest run. Each load process talks to one broker, so a fault has two audiences: the clients of the node that went down, and everyone else.

Messages pushed per second by all nine load processes, every 10 seconds of the 3 h 23 min soak4 run, offered 500,000 msg/s. The line is flat at 500,000 except at the three faults: after the kill -9 of the leader at 01:49:50, the SIGTERM of the leader at 03:19:50 and the kill -9 of a follower at 04:49:50 it drops to about 333,000, the two thirds of the load whose broker was still up, for about a minute, then overshoots to between 570,000 and 634,000 while the restarted node's clients catch up, and is back at 500,000 within two minutes of each fault.msg/s pushed, all nine load processes0200k400k600k01:4002:1002:4003:1003:4004:1004:40kill -9 leaderSIGTERM leaderkill -9 followerpushed
Pushes per second across the run. Each fault takes away the third of the load attached to the node that went down, and only for as long as that node is down. Source: benchmark-queen/2026-09-30-stage03-3node/runs/soak4 (g0.log to g8.log)
Push p99 latency of the worst of the nine load processes, every 10 seconds of the soak4 run, on a log scale. Outside the faults, half the windows are under 28 ms and 95% under 161 ms. After the kill -9 of the leader it reaches 3.8 seconds for the surviving nodes' clients while a new leader is elected, then about 32 seconds for the clients of the restarted node while it catches up. The SIGTERM and the follower kill show the same catch-up peak, 52 and 37 seconds, and no election spike.push p99, worst load process10 ms100 ms1 s10 s100 s01:4002:1002:4003:1003:4004:1004:40kill -9 leaderSIGTERM leaderkill -9 followerpush p99
The worst load process in each 10-second window. The tall peaks belong to the clients of the node that was down, waiting for it to catch up; the other clients saw at most the 3.8-second election after the leader kill. Source: benchmark-queen/2026-09-30-stage03-3node/runs/soak4 (g0.log to g8.log)

A kill -9 of the leader

At 01:49:50 the leader was killed. The clients of the two surviving nodes had one 10-second window with a push p99 of 3.2 to 3.8 seconds while the cluster elected a new leader. Each of their load processes saw about 220 pops time out and 300 to 900 acks fail, the ones in flight when the leader died. The load generator never retries an ack, so those messages came back when their 30-second leases ran out. The three processes attached to the killed node got errors until it returned. Restarted at 01:50:20, it took requests again while it caught up on what it had missed, and those requests waited for the catch-up (push p99 near 32 seconds). Each queue had about one consumer, so the third of the queues whose consumer sat on the killed node waited as well, and their messages arrived up to 75 seconds late. The redelivered and the delayed messages are the e2e p99 of 24 to 75 seconds in the windows of the next minute and a half. All nine processes were back to normal at 01:51:30, 100 seconds after the kill.

A SIGTERM of the leader

At 03:19:50 the leader got a SIGTERM and handed leadership over before it stopped. The surviving nodes’ clients saw a push p99 of 330 ms in that window and not one failed ack. The stopping node’s own clients had about 50,000 acks fail, for messages they had popped in its last moments; those came back after their leases, and the node’s restart went as above.

A kill -9 of a follower

At 04:49:50 a follower was killed. The other two nodes’ clients kept their full rate with no failed acks, and their push p99 stayed under 600 ms while the follower was away. The killed node’s clients went through the same restart and catch-up as above, and the messages they pushed during it reached the other consumers up to 36 seconds late.

Memory and disk

Over the 3 h 23 min of soak4, through every fault, the leader’s anonymous memory stayed between 7.3 and 8.3 GB, apart from one 9.2 GB sample during the follower kill, and the followers’ between 5.7 and 6.1 GB. Each node’s disk held between 47 and 50 GB the whole time: retention freed log files as fast as new messages filled them.

What this does not show

The longest run is 3 h 23 min at 500,000 msg/s; days, or a million a second for hours, are not shown. Each load process stayed on one broker, so the clients of a dead node waited for it. A client behind a load balancer, or one with more than one broker address, would have moved, and since 2.0.0-beta.1 a restarted node’s /health reports its catch-up lag, so a readiness probe keeps traffic off it until it has caught up; the soak used neither. And the load generator counts messages, not their identity, so loss and duplication are left to the Jepsen tests.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close