---
title: "Run a cluster"
description: "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."
---

> Queen MQ documentation, for AI agents
> Complete self-contained summary of Queen MQ: https://queenmq.com/llms-brief.txt
> Fetch that first when the question is about the product rather than about this page.
> Index of all pages: https://queenmq.com/llms.txt

# Run a cluster

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:

```bash
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
```

```text
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:

```bash
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" }')"
```

```json
[{"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:

```bash
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
```

```text
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:

```text
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:

```bash
docker compose up -d --wait queen-3
curl -s localhost:36632/health | jq -c .raft
```

```json
{"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:

```yaml
# 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.

**Figure.** 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.

Any node takes any request: followers forward writes to the leader, so a client can hold all three URLs without knowing which one leads.

- your app: urls: all three nodes
- queen-1: 127.0.0.1:16632
- queen-2: 127.0.0.1:26632
- queen-3: 127.0.0.1:36632
- volume queen-1: /var/lib/queen/raft
- volume queen-2: /var/lib/queen/raft
- volume queen-3: /var/lib/queen/raft
- Group "raft on :7400, one QUEEN_RAFT_TOKEN": queen-1, queen-2, queen-3
- your app → queen-1
- your app → queen-2: HTTP
- your app → queen-3
- queen-1 connects to volume queen-1
- queen-2 connects to volume queen-2
- queen-3 connects to volume queen-3

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:

```js
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.

**Figure.** 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.

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.

1. operator → node 1: SIGTERM
2. node 1 → node 2: ephemeral partitions
3. node 1 → itself: in-flight acks commit (500 ms at most)
4. node 1 → node 2: transfer leadership
   (node 2: leads within 3 s)
   (node 1: close the listener, finish requests as a follower, exit)

*restarted on the same data directory*

5. node 2 → node 1: the entries it missed
6. node 1 → operator: /health 200, caught up
   (operator: 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](/operate/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](/concepts/guarantees/).

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:

```text
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](/concepts/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](#raft-groups)) or more
clusters. What one three-node cluster carries is on the [benchmarks](/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:

```bash
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](/guides/ephemeral/) 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](/operate/tenants/), 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](/operate/) 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](/operate/security/#the-raft-token)). Port 6632 has no authentication
unless you turn on broker JWT, so it does not face the internet either. When you add the
[proxy](/operate/tenants/) 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](/operate/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](/operate/monitoring/#health)), so use it for readiness. For alerts use a write probe
([recovery](/operate/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:

```bash
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:

```bash
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](/operate/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](/operate/recovery-backups/) 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](/operate/kubernetes/) turns this cluster into a StatefulSet, and
[monitoring](/operate/monitoring/) covers what to watch once it runs.

Source: https://queenmq.com/operate/cluster/index.mdx
