Skip to content

Run a cluster

Three or five Queen nodes as one Raft cluster: a Compose file to try it on a laptop, how the leader and the followers share the work, what each failure costs, rolling upgrades, and what to set in production.

Updated View as Markdown

Three copies of the same binary, each told where the other two are, make a Queen cluster. They keep one replicated log between them, so a write is on a majority of disks before anyone hears yes, and any of them can serve any client. One node can be down for a crash or an upgrade while the other two go on committing. The quickest way to see it is the Compose file in the repository, which runs the whole cluster on a laptop.

Try it on a laptop

deploy/compose/three-node/ starts three nodes on a private network, each with its own volume. Node N answers on localhost:N6632 (16632, 26632 and 36632), API and dashboard alike, which keeps the cluster clear of a single node on 6632. Give it a token, start it, and ask each node what it is:

cd deploy/compose/three-node        # in a clone of github.com/queen-mq/queen
echo "QUEEN_RAFT_TOKEN=$(openssl rand -hex 32)" > .env
docker compose up -d --wait
for n in 1 2 3; do echo "queen-$n $(curl -s localhost:${n}6632/health | jq -r .raft.role)"; done
queen-1 follower
queen-2 follower
queen-3 leader

The three nodes read the same peer list on empty volumes, voted, and node 3 won. --wait returns once every node’s healthcheck passes, which on an Apple-silicon laptop, running the amd64 image under emulation, took eight seconds. Now push through node 1, a follower, pop through node 2 and ack through node 3:

curl -s -X POST localhost:16632/api/v1/push -H 'content-type: application/json' -d '{
  "items": [{ "queue": "orders", "partition": "customer-42",
              "transactionId": "order-1001-created",
              "payload": { "orderId": 1001, "total": 42 } }]
}'
POP=$(curl -s "localhost:26632/api/v1/pop/queue/orders?batch=1&wait=true&timeout=5000")
echo "$POP" | jq -c '.messages[0] | {transactionId, partition, data}'
curl -s -X POST localhost:36632/api/v1/ack -H 'content-type: application/json' -d "$(echo "$POP" | jq '{
  transactionId: .messages[0].transactionId, partitionId: .messages[0].partitionId,
  leaseId: .leaseId, status: "completed" }')"
[{"index":0,"message_id":"01a0fcd6-2243-7000-a59b-d7da302543ea","transaction_id":"order-1001-created","queueName":"orders","status":"queued","offset":0}]
{"transactionId":"order-1001-created","partition":"customer-42","data":{"orderId":1001,"total":42}}
[{"index":0,"transactionId":"order-1001-created","success":true,"error":null,"leaseReleased":true,"dlq":false,"noop":false}]

Each node answered as the leader would have. Now stop the leader:

docker compose stop queen-3
for n in 1 2; do echo "queen-$n $(curl -s localhost:${n}6632/health | jq -r .raft.role)"; done
queen-1 follower
queen-2 leader

docker compose stop sends SIGTERM, and a leader that receives one hands its leadership to the most caught-up follower before it exits. Node 2 was leading within a second of the stop, and node 3’s log says how long the hand-off itself took:

INFO rsm: raft: leadership handed off to=2 ms=121

Two nodes of three are a majority, so pushes, pops and acks go on through nodes 1 and 2 as before. Start node 3 again and it catches up from the new leader:

docker compose up -d --wait queen-3
curl -s localhost:36632/health | jq -c .raft
{"role":"follower","leader":true,"term":3,"applied":13,"commit":13,"lag":0,"storageReady":true,"clusterVersion":3,"kinds":3,"apply":{"failure":null,"skipped":[]}}

applied equal to commit, the same 13 the leader reports, means node 3 has applied everything the cluster committed, the writes it missed included. docker compose down -v removes the containers and the three volumes when you are done.

This is the whole file. The comments give the reason for each line that is not obvious:

# Three Queen nodes on one machine: a real Raft cluster for trying replication,
# leader failover and rolling restarts. Walkthrough in README.md next to this
# file and at https://queenmq.com/operate/cluster/.
#
#   echo "QUEEN_RAFT_TOKEN=$(openssl rand -hex 32)" > .env
#   docker compose up -d --wait
#   curl -s localhost:16632/health
#
# Node N answers on localhost:N6632 (16632, 26632, 36632), so the cluster does
# not collide with a single node on 6632. Each node keeps its state in its own
# volume; `docker compose down -v` deletes them.

name: queen-cluster

x-node: &node
  # QUEEN_VERSION in .env picks another release (README.md: rolling upgrade).
  image: ghcr.io/queen-mq/queen:${QUEEN_VERSION:-latest}
  # The images are published for linux/amd64 only. On an arm64 machine
  # (Apple silicon) Docker runs them under emulation: slower, fine for a trial.
  platform: linux/amd64
  # A node that was away long enough to need a snapshot exits with code 75 to
  # load it at start; the restart policy brings it back.
  restart: unless-stopped
  # On SIGTERM a node hands its ephemeral partitions and, if it leads, its
  # leadership to a peer, then lets the long polls in flight finish (30 s by
  # default). Docker's 10 s default would kill it halfway.
  stop_grace_period: 60s
  healthcheck:
    # 200 once the node knows a leader and has caught up, 503 "settling" before.
    test: ["CMD", "curl", "-fsS", "-o", "/dev/null", "http://localhost:6632/health"]
    interval: 5s
    timeout: 3s
    retries: 3
    start_period: 60s

x-env: &env
  QUEEN_RAFT_REPLICATOR: openraft
  # Every voter as id=raft_address/http_address, the same list on every node.
  QUEEN_RAFT_PEERS: "1=queen-1:7400/queen-1:6632,2=queen-2:7400/queen-2:6632,3=queen-3:7400/queen-3:6632"
  # The secret every Raft RPC carries. The raft port (7400) is not published, so
  # only these three containers reach it.
  QUEEN_RAFT_TOKEN: ${QUEEN_RAFT_TOKEN:?put QUEEN_RAFT_TOKEN in .env next to compose.yaml, see the header}
  QUEEN_RAFT_DIR: /var/lib/queen/raft

services:
  queen-1:
    <<: *node
    hostname: queen-1
    environment:
      <<: *env
      QUEEN_RAFT_NODE_ID: "1"
    # The broker port has no authentication of its own: published on loopback.
    ports: ["127.0.0.1:16632:6632"]
    volumes: ["queen-1:/var/lib/queen/raft"]

  queen-2:
    <<: *node
    hostname: queen-2
    environment:
      <<: *env
      QUEEN_RAFT_NODE_ID: "2"
    ports: ["127.0.0.1:26632:6632"]
    volumes: ["queen-2:/var/lib/queen/raft"]

  queen-3:
    <<: *node
    hostname: queen-3
    environment:
      <<: *env
      QUEEN_RAFT_NODE_ID: "3"
    ports: ["127.0.0.1:36632:6632"]
    volumes: ["queen-3:/var/lib/queen/raft"]

volumes:
  queen-1:
  queen-2:
  queen-3:

How it works

A cluster runs one Raft group over one log, with the openraft crate doing the protocol. The leader puts every write into that log in one order and replicates it, and an entry commits once a majority of the voters have written it and fsynced it. Each node then applies the same entries in the same order to its own copy of the state: the queue logs, the cursors and leases, the KV store, the timers. There is no second store beside the log, which is why a node that was away needs only the entries it missed (or a snapshot, if it was away long) to be whole again.

The three-node Compose cluster. A client holds the URLs of all three nodes. queen-1, queen-2 and queen-3 each publish their HTTP port 6632 on the host as 127.0.0.1:16632, 26632 and 36632, talk raft to each other on port 7400 with one shared QUEEN_RAFT_TOKEN, and keep their data directory in a volume of their own.your appurls: all three nodesqueen-1127.0.0.1:16632queen-2127.0.0.1:26632queen-3127.0.0.1:36632volume queen-1/var/lib/queen/raftvolume queen-2/var/lib/queen/raftvolume queen-3/var/lib/queen/raftraft on :7400, one QUEEN_RAFT_TOKENHTTP
Any node takes any request: followers forward writes to the leader, so a client can hold all three URLs without knowing which one leads. Source: deploy/compose/three-node/compose.yaml

We keep a tenant inside one Raft group because then a transaction is a single entry. An ack, three pushes and a KV write go into the log together and commit together, or none of them does, and no protocol between logs has to agree about it.

Followers carry their share of the work. A follower parses the requests of its own clients, prepares each one as a command, and sends the commands to the leader over a few long-lived streams on the raft port; it answers its client once the entry has committed. A pop is answered once this node has applied the claim, because the messages it hands out come from its own copy of the queue log. A read first waits until the node has applied everything the cluster had committed when the read began, so no client reads behind a write it was already told about, whichever node it asks. That is what made the walkthrough above work, and it is why a client can hold every node’s URL without caring which one leads:

const queen = new Queen({ urls: ['http://queen-1:6632', 'http://queen-2:6632', 'http://queen-3:6632'] })

With several URLs the JavaScript client spreads its requests over the nodes and moves to the next one when a node fails with a network error or a 5xx.

A stop is planned work, so we made it cheap. On SIGTERM a node first hands its ephemeral partitions to their next owners. If it leads, it then lets the acks already in flight commit (at most 500 ms) and transfers leadership to the most caught-up voter, waiting up to 3 s for it to take over. Only then does it close its listener and finish the requests in flight, as a follower. A rolling restart therefore costs the cluster one transfer and a vote, about a tenth of a second in the run above.

A rolling restart of the leader. The operator sends SIGTERM to node 1, the leader. It hands its ephemeral partitions to their next owners, lets the acks already in flight commit for at most 500 ms, and transfers leadership to node 2, the most caught-up voter, waiting up to 3 seconds for it to take over. Node 1 then closes its listener, finishes its requests in flight as a follower and exits. Restarted on the same data directory, it catches up on the entries it missed from node 2, and its /health answers 200 once it is within 1,000 entries or 2 seconds of the leader. Then the operator moves to the next node.operatornode 1leader, stoppingnode 2most caught upSIGTERMephemeral partitionsin-flight acks commit500 ms at mosttransfer leadershipleads within 3 sclose the listener, finishrequests as a follower, exitrestarted on the same data directorythe entries it missed/health 200, caught upnext node
A planned stop costs one transfer and a vote, about a tenth of a second in the Compose rehearsal. /health turning 200 is the signal to touch the next node.

What each failure costs

Event What the cluster does
A follower stops or crashes The other two keep committing, and its clients move to another URL. When it starts again it catches up from the leader.
The leader gets SIGTERM It hands leadership over first: under a second to a new leader in the run above.
The leader crashes (kill -9, OOM, a lost host) The followers wait out the election timeout (1 to 2 s, QUEEN_RAFT_ELECTION_MS) and elect a new leader: 3.7 s in the Compose cluster under emulation, at most 4 s in the recovery rehearsals.
A leader hears from no quorum for 3 s It hands leadership to the most caught-up voter, so a leader that can send but not receive does not hold every write hostage.
A majority is down Nothing commits. Requests wait out their deadline and fail with 503, code no_leader or timeout, and Retry-After: 1, until a majority is back.
A node was away more than 10 minutes (QUEEN_RAFT_PURGE_HOLD_S, 600) The leader stops holding its log for it. If the entries it missed are purged, the node receives a snapshot, exits with code 75 and loads it on its next start, and a restart policy brings it back.
A majority of the disks is lost for good Only a forced recovery brings the cluster back, and it can lose acknowledged writes (recovery).

In the recovery rehearsals of 2026-10-01 (three nodes in Docker under load, test/recovery/), kill -9 of the leader or of a follower and SIGTERM of the leader lost none of 3,300 acknowledged writes, and no leader change took more than 4 s. The guarantees behind those numbers are Jepsen-tested.

Losing the majority deserves a look, because the node left standing does not notice. With nodes 1 and 3 of the Compose cluster stopped, node 2 went on answering /health with 200, "status":"healthy" and "role":"leader", while a push to it waited out its 30 s deadline and got this:

HTTP/1.1 503 Service Unavailable
retry-after: 1

{"error":"the request deadline elapsed","code":"timeout"}

A write that fails like that, or with 503 outcome_unknown while the leader changes, may or may not have applied, so send it again. A repeated ack of a lease it already released is answered as the first one was, and a push repeated with the same transactionId is deduplicated (dedup). The push above, sent again once the majority was back, answered "status":"queued": it had never committed, and the retry wrote it once.

Three nodes or five

Run three voters, or five. A majority of three is two, so three nodes survive one failure; five survive two, at the price of a larger majority on every commit and two more followers for the leader to feed. Four nodes survive one failure like three while every write waits for three of them, so an even count only costs. Put each voter on its own machine, and in its own zone where you have zones: two voters on one host turn that host’s failure into a lost majority, which is also why the Compose cluster is a place to try things and nothing more.

More nodes make a cluster harder to kill, not faster, because one leader orders every write of a group. To spread write load, run several Raft groups per process (below) or more clusters. What one three-node cluster carries is on the benchmarks page.

Give every voter the same disk size. The disk gate (507 above 85% full, QUEEN_RAFT_DISK_HIGH_PCT) only turns away a node’s own clients, and a smaller disk fills from replication anyway.

Rolling upgrades

Upgrade one node at a time: stop it, start the new release on the same data directory, and wait until its /health answers 200 before you touch the next. Healthy means it knows the leader and has nearly caught up, within 1,000 entries of the leader’s commit or less than 2 s behind it (QUEEN_RAFT_READY_LAG_ENTRIES, QUEEN_RAFT_READY_LAG_MS), so by the time you stop the next node this one is back in the majority. In the Compose cluster, set QUEEN_VERSION in .env to the release you move to and let --wait hold each step until the node is healthy:

for n in 1 2 3; do docker compose up -d --wait queen-$n; done

We rehearsed that roll from 2.0.0-beta.5 to beta.6 while a writer pushed about 15 messages a second round-robin over the three nodes, retrying a failed push on the next node with the same transactionId. All 1,172 pushes were answered, and afterwards every node held all 1,172 messages, each once. Six attempts met a node while it restarted and went on to the next one, and the slowest push took 0.63 s.

A release that adds something to the log’s format has to roll without any old node meeting an entry it cannot decode. So each node tells the leader, on every append it answers, the highest format version it reads (kinds in /health), and the leader raises the cluster’s version (clusterVersion) only once every member reports the new one. Until the last node runs the new release, the leader goes on writing the old format and the feature that needs the new one stays off; then the version rises and the feature turns on, with no flag for you to flip. A node whose build reads less than the cluster’s version refuses to start on that data, and the leader refuses to add it as a learner, so a downgrade past a format change stops at the first node. The gate exists from 2.0.0-beta.2 on; the move to beta.1 was a stop-all upgrade.

Ephemeral queues in a cluster

Ephemeral queues keep their messages in memory and never in the log, and still every node of a cluster serves them. Each (queue, partition) has exactly one owner, picked by a highest-random-weight hash over the live members, so every node computes the same owner without asking anyone, and a node that does not own a partition forwards the push, pop or ack to the one that does. In the Compose cluster, six partitions pushed through node 1 landed two on node 2 and four on node 3, and every message could be popped through any node.

When the membership changes, the partitions that move are handed over with their messages and each consumer group’s position. A node stopping on SIGTERM tells its peers it is leaving and ships every partition it owns to the next owner before it hands off leadership. Stopping node 3 logged ephemeral rings handed over rings=4 delivered=4 lost=0, and its eight messages then popped through node 1. When node 3 came back, nodes 1 and 2 shipped 13 partitions back to it, none lost.

The messages live in one node’s memory, so a crash (kill -9, OOM, a lost host) loses what that node owned. Of six messages spread over the cluster, a SIGKILL of node 3 left the two that node 2 owned. Two issues are still open in 2.0.0-beta.6. A hand-over to a node that is still restarting can drop the partition’s messages, which we have seen during rolling restarts. And through the proxy, a request forwarded to the owning node does not keep its tenant, so ephemeral queues behind the proxy are safe on a single node only, for now.

In production

The Compose file holds every setting a cluster needs, and production changes what surrounds it, starting with one machine per voter and a private network between them. Every node gets the single-node settings plus these:

Variable Default What it does
QUEEN_RAFT_REPLICATOR local openraft (or raft) makes the node a cluster member. A peer list without it fails the boot
QUEEN_RAFT_NODE_ID 1 A number from 1, different on each node. ordinal reads the trailing number of HOSTNAME plus one, so pod queen-0 is node 1
QUEEN_RAFT_PEERS unset Every voter as id=raft_addr/http_addr, the same list on every node. Unset, the node is a single voter
QUEEN_RAFT_LISTEN 0.0.0.0: plus the node’s own raft port Where this node’s raft RPC server binds
QUEEN_RAFT_TOKEN unset The secret every raft RPC carries in x-queen-raft-token. Unset, anything that reaches the raft port can vote and append
QUEEN_RAFT_DIR /var/lib/queen/raft This node’s data directory, never shared between nodes
QUEEN_RAFT_JOIN off The node joins a cluster that already exists, and serves 503 until the leader adds it

The raft port (7400 by convention) carries votes, log entries, snapshots and the commands followers send the leader, over plain HTTP with the token in a header. Keep it on a private network that only the nodes reach (security). Port 6632 has no authentication unless you turn on broker JWT, so it does not face the internet either. When you add the proxy for tenants and API keys, give it its own port (QUEEN_PROXY_PORT) and keep PORT, the address in the peer list, for the nodes and your internal clients.

Every node of a new cluster starts on an empty directory with the same peer list, and each one initializes the cluster from it. After that first start the membership lives in the log, and the list is read only for this node’s own listen address, the token and a group’s preferred leader. A committed membership change therefore stands whatever QUEEN_RAFT_PEERS still says. Update the list anyway, so the next new node sees the current members. The addresses are part of the membership too: renaming hosts, ports or a namespace on a live cluster takes a forced recovery under the new names.

Give a stopping node time to finish. The ephemeral hand-over can take up to 15 s, the leadership transfer up to 3 s, and a long-poll pop in flight up to its timeout (30 s by default), which is why the Compose file sets a 60 s stop grace period and the Kubernetes manifest a 90 s one.

Leases are timed on the leader’s clock. QUEEN_RAFT_MAX_CLOCK_SKEW_MS (500) is how long a lease another node timed is held past its expiry, and it lengthens the pause before a new leader serves consumption. Lease exclusivity across a leader change holds only while the nodes’ clocks agree within that bound, so run NTP or chrony on every node.

/health answers 200 once a node knows a leader and has caught up (monitoring), so use it for readiness. For alerts use a write probe (recovery), because a node that lost its quorum keeps answering 200, as node 2 did above.

Change the membership

QUEEN_RAFT_PEERS is read once, at the first boot. From then on the membership changes only through these routes, which need the admin access level and are blocked by the proxy, so call them on PORT:

Route Body, or what it answers
GET /api/v1/system/raft/membership Answers the leader, the term, voters, learners, and each member’s lag and live
POST /api/v1/system/raft/membership/learners {"id": 4, "raft": "queen-4:7400", "http": "queen-4:6632"}
POST /api/v1/system/raft/membership/promote {"ids": [4]}
PUT /api/v1/system/raft/membership/voters {"voters": [1, 2, 4]}
DELETE /api/v1/system/raft/membership/members/:id none

The leader checks every change before it writes it. One change runs at a time, a live majority has to exist before and after (409 no_quorum otherwise), the last voter cannot be removed, and a learner is promoted only within 1,000 entries of the leader (QUEEN_RAFT_PROMOTE_MAX_LAG; 409 learner_behind, or "force": true). Every change is idempotent, so a lost answer is safe to repeat. To replace node 3 after its disk is gone:

curl -X DELETE http://queen-1:6632/api/v1/system/raft/membership/members/3
# wipe node 3's directory, then start it with QUEEN_RAFT_JOIN=true: it serves 503 until added
curl -X POST http://queen-1:6632/api/v1/system/raft/membership/learners \
  -H 'content-type: application/json' -d '{"id":3,"raft":"queen-3:7400","http":"queen-3:6632"}'
curl -X POST http://queen-1:6632/api/v1/system/raft/membership/promote \
  -H 'content-type: application/json' -d '{"ids":[3]}'

Remove a member before you wipe its disk. A voter that comes back empty under its old id is never repaired: it looks healthy and its clients hang.

Raft groups

QUEEN_RAFT_GROUPS=3 (default 1, at most 64, the same on every node) runs that many Raft groups in each node, each with its own log, leader and data directory (groups/g<N>/). A tenant lives in exactly one group, picked by a stable hash of its id or pinned with QUEEN_TENANT_GROUPS=<tenant>=<group>,..., so its transactions stay single entries. Group N listens on the raft port plus N and prefers the member at position N mod n as its leader, so with as many groups as nodes every node leads one and follows the others. It needs QUEEN_RAFT_CLIENT_OFFLOAD on, which is the default.

Raft groups were not part of the Jepsen pass, which ran one group. A tenant cannot change group while it has data, because there is no migration yet, so pin a large tenant before it writes.

Recovery

When something looks wrong, start with a write probe rather than /health, because a node that lost its quorum can keep answering 200:

curl -m 10 -X PUT http://queen-1:6632/api/v1/kv/ops/probe \
  -H 'content-type: application/json' -d '{"value":"probe","ttlSeconds":60}'

An answer within a second means the cluster commits; nothing after 10 s means it does not. From there, Recovery walks through every case we rehearsed on real clusters (a lost node, a full disk, a lost majority, a forced recovery, a poisoned entry, a rolling restart gone wrong), and Backups and restore covers snapshots of the data directory.

Limits

  • Three containers on one host, as in the Compose file, all go down with that host.
  • Losing a majority stops writes, and losing a majority of the disks for good can lose acknowledged data.
  • /health can stay 200 on a node that lost its quorum, and on an empty node that never caught up (F8 and F2 in test/recovery/FINDINGS.md). Alert on the write probe.
  • The raft port is plain HTTP with a shared token: keep it on a private network.
  • Peer addresses are fixed at the first boot. Changing them on a live cluster is a forced recovery.
  • Ephemeral messages live in their owner’s memory: a crash loses that node’s share, and a hand-over to a node that is still restarting can drop a partition’s messages. Behind the proxy they are safe on a single node only (2.0.0-beta.6).
  • One leader orders every write of a group. Adding nodes to a group adds no write throughput.

Kubernetes turns this cluster into a StatefulSet, and monitoring covers what to watch once it runs.

Navigation

Type to search…

↑↓ navigate↵ selectEsc close