On Kubernetes a Queen cluster is a StatefulSet of three (or five) pods. Each pod gets its own
PersistentVolumeClaim, the pods find each other through a headless Service, a client Service
carries the HTTP port, and a disruption budget lets voluntary evictions take one pod at a time. The
repository ships no Helm chart, so the manifest below is the complete deployment, in namespace
queen, to adapt or template in your own chart. It is the shape a production deployment of 2.0
runs, and the fields that are not obvious are explained after it.
The manifest
kubectl create namespace queen
kubectl -n queen create secret generic queen \
--from-literal=QUEEN_RAFT_TOKEN="$(openssl rand -hex 32)" \
--from-literal=QUEEN_ENCRYPTION_KEY="$(openssl rand -hex 32)"apiVersion: v1
kind: Service
metadata: { name: queen }
spec:
selector: { app: queen }
ports: [{ name: http, port: 6632, targetPort: http }]
---
apiVersion: v1
kind: Service
metadata: { name: queen-headless }
spec:
clusterIP: None
publishNotReadyAddresses: true
selector: { app: queen }
ports:
- { name: http, port: 6632, targetPort: http }
- { name: raft, port: 7400, targetPort: raft }
---
apiVersion: policy/v1
kind: PodDisruptionBudget
metadata: { name: queen }
spec:
maxUnavailable: 1
selector: { matchLabels: { app: queen } }
---
apiVersion: apps/v1
kind: StatefulSet
metadata: { name: queen }
spec:
replicas: 3
serviceName: queen-headless
podManagementPolicy: Parallel
minReadySeconds: 30
updateStrategy: { type: RollingUpdate }
selector: { matchLabels: { app: queen } }
template:
metadata: { labels: { app: queen } }
spec:
terminationGracePeriodSeconds: 90
securityContext: { runAsNonRoot: true, runAsUser: 65532, runAsGroup: 65532, fsGroup: 65532 }
affinity:
podAntiAffinity:
requiredDuringSchedulingIgnoredDuringExecution:
- topologyKey: kubernetes.io/hostname
labelSelector: { matchLabels: { app: queen } }
topologySpreadConstraints:
- { maxSkew: 1, topologyKey: topology.kubernetes.io/zone, whenUnsatisfiable: ScheduleAnyway,
labelSelector: { matchLabels: { app: queen } } }
containers:
- name: queen
image: ghcr.io/queen-mq/queen:latest
ports:
- { name: http, containerPort: 6632 }
- { name: raft, containerPort: 7400 }
env:
- { name: QUEEN_RAFT_REPLICATOR, value: openraft }
- { name: QUEEN_RAFT_DIR, value: /var/lib/queen/raft }
- { name: QUEEN_RAFT_NODE_ID, value: ordinal }
- name: QUEEN_RAFT_PEERS
value: "1=queen-0.queen-headless.queen.svc.cluster.local:7400/queen-0.queen-headless.queen.svc.cluster.local:6632,2=queen-1.queen-headless.queen.svc.cluster.local:7400/queen-1.queen-headless.queen.svc.cluster.local:6632,3=queen-2.queen-headless.queen.svc.cluster.local:7400/queen-2.queen-headless.queen.svc.cluster.local:6632"
- name: QUEEN_RAFT_TOKEN
valueFrom: { secretKeyRef: { name: queen, key: QUEEN_RAFT_TOKEN } }
- name: QUEEN_ENCRYPTION_KEY
valueFrom: { secretKeyRef: { name: queen, key: QUEEN_ENCRYPTION_KEY } }
- name: QUEEN_SERVER_ID
valueFrom: { fieldRef: { fieldPath: metadata.name } }
- { name: QUEEN_LOG_JSON, value: "true" }
resources:
requests: { cpu: "1", memory: 4Gi }
limits: { memory: 4Gi }
startupProbe: { tcpSocket: { port: http }, periodSeconds: 5, timeoutSeconds: 5, failureThreshold: 360 }
readinessProbe: { httpGet: { path: /health, port: http }, periodSeconds: 5, timeoutSeconds: 5, failureThreshold: 3 }
livenessProbe: { tcpSocket: { port: http }, periodSeconds: 20, timeoutSeconds: 5, failureThreshold: 3 }
lifecycle: { preStop: { exec: { command: ["/bin/sleep", "5"] } } }
securityContext: { allowPrivilegeEscalation: false, readOnlyRootFilesystem: true, capabilities: { drop: ["ALL"] } }
volumeMounts:
- { name: data, mountPath: /var/lib/queen/raft }
- { name: tmp, mountPath: /tmp }
volumes: [{ name: tmp, emptyDir: {} }]
volumeClaimTemplates:
- metadata: { name: data }
spec:
accessModes: ["ReadWriteOnce"]
resources: { requests: { storage: 50Gi } }kubectl -n queen rollout status statefulset/queen
kubectl -n queen exec queen-0 -- curl -s localhost:6632/api/v1/system/raft/membership \
| jq -c '.membership | {leader, term, voters, members: [.members[] | {nodeId, live, lag}]}'{"leader":1,"term":4,"voters":[1,2,3],"members":[{"nodeId":1,"live":true,"lag":0},{"nodeId":2,"live":true,"lag":0},{"nodeId":3,"live":true,"lag":0}]}The answer lists three voters, all live and none behind, so the cluster has formed. Clients use
the queen Service and may land on any pod, because every node serves reads and writes itself and
sends only the prepared command to the leader.
The probes
The startup and liveness probes are a TCP connect on 6632. The node binds its port only after the store and the replicator have opened, so an open port means the boot finished, and 360 probes of 5 seconds give a long log replay 30 minutes.
Readiness is /health, which answers 503 settling while the node knows no leader or while its
apply lags behind the leader’s commits. The Service stops sending clients to a pod that is
settling, and minReadySeconds gives a rolling update time to catch up before it moves to the next
pod. Never point liveness at /health: an election, or a lost majority, would then restart every
pod at once, which is the one thing that cannot help.
What SIGTERM does
When the kubelet stops a pod, the preStop sleep gives the endpoints a few seconds to drop it, and
then the node shuts down in an order chosen so that the cluster loses as little as possible:
- It tells its peers it is leaving and hands its ephemeral partitions, with their contents, to their next owners. The hand-over gives up after 15 seconds.
- If it leads, it lets the acks already in flight commit (at most half a second) and transfers leadership to the most caught-up voter, waiting up to 3 seconds for it to take over, so a rolling restart costs the cluster one transfer instead of an election timeout without a leader.
- It stops accepting connections and lets the requests in flight finish. A parked long-poll pop can wait 30 seconds by default.
- The Kafka facade, when on, gives its registry row back (5 seconds), and the proxy, when on, writes its last usage minute (5 seconds).
- If leadership came back to it during the drain, it hands it off once more and exits.
terminationGracePeriodSeconds: 90 covers the worst case of that sequence with room to spare.
Shorter, and the kubelet’s SIGKILL can land in the middle of it, which turns a graceful stop into a
crash: the ephemeral contents still on the node are lost, and the cluster waits for an election.
The fields that are not obvious
podManagementPolicy: Parallel and publishNotReadyAddresses: true go together. Without both, a
full restart deadlocks, because no pod becomes ready without a leader and no leader is elected
until the peers’ names resolve.
QUEEN_RAFT_NODE_ID=ordinal makes pod queen-0 node 1: the node id is the trailing number of the
hostname plus one. The peer list is id=raft address/HTTP address for every pod, and it is the
same string on all of them.
The pods have no CPU limit. The heartbeat and the commit path are sensitive to latency, and CFS throttling turns a burst of work into a missed heartbeat; the memory request equals the limit. Hard anti-affinity puts each voter on its own machine, because two voters on one machine turn that machine’s failure into a lost majority, and the zone spread does the same for zones where you have them.
A pod that exits with code 75 has received a snapshot after a long absence and restarts to load it. The kubelet restarts it in place on the same volume. Once is normal; a loop is not.
Writes that grow storage answer 507 above 85% of a pod’s volume. Size the claim for your
retention and give every pod the same size, because the gate protects a node only from its own
clients: a smaller volume keeps filling from replication.
Upgrades
The manifest above runs latest, the current release. For a cluster you keep, pin the version it
runs, because a pod that restarts on latest pulls whatever release is current by then, and an
upgrade should be a step you take. Change the image tag and apply. The RollingUpdate replaces one
pod at a time, highest ordinal first,
and waits for each to be Ready plus minReadySeconds before it moves on, which is the order a
rolling upgrade of the cluster needs (run a cluster).
kubectl -n queen set image statefulset/queen queen=ghcr.io/queen-mq/queen:<new-version>
kubectl -n queen rollout status statefulset/queenA new release that changes the data format writes the new format only once every member reads it,
so old and new pods coexist through the rollout. For a repair that changes an environment variable
on some pods only, switch updateStrategy to OnDelete first and delete exactly the pods that
must take the change, because a RollingUpdate waits forever for pods that cannot become Ready
(recovery).
Fixed at the first boot
The peer addresses are stored in the raft log. Scaling the StatefulSet does not add a voter, and these settings cannot change on a live cluster without a recovery:
replicas(add or remove voters through the membership API, one at a time);- the StatefulSet, headless Service and namespace names, which are part of the peer names;
- the raft port and the HTTP port.
A bigger volume is a PVC expansion per pod followed by a restart of that pod. The
volumeClaimTemplates of a StatefulSet cannot be edited in place, so delete the StatefulSet with
--cascade=orphan and apply it again with the new size, which new pods then get too.
The proxy and Kafka
For tenants and API keys, add the proxy on a port of its own:
- { name: QUEEN_PROXY_EMBEDDED, value: "true" }
- { name: QUEEN_PROXY_PORT, value: "6711" }
- { name: QUEEN_PROXY_SPOOL_DIR, value: /var/lib/queen/raft/proxy-spool }
- name: QUEEN_PROXY_JWT_SECRET
valueFrom: { secretKeyRef: { name: queen, key: QUEEN_PROXY_JWT_SECRET } }
- name: QUEEN_PROXY_CP_TOKEN
valueFrom: { secretKeyRef: { name: queen, key: QUEEN_PROXY_CP_TOKEN } }Every pod needs the same session secret, and the spool directory has to be on the volume because the root filesystem is read-only. Expose 6711 through a Service of its own and your Ingress. Port 6632 stays inside the cluster: with the proxy on 6711 it is the broker itself, without authentication, and it is also the port the pods use to reach each other.
For Kafka clients, add QUEEN_KAFKA_EMBEDDED=true and, per pod,
QUEEN_KAFKA_ADVERTISED_ADDR=$(QUEEN_SERVER_ID).queen-headless.queen.svc.cluster.local:9092. With
several pods, the facades present themselves to clients as one Kafka cluster: each pod needs an
integer QUEEN_KAFKA_NODE_ID from 1 to 64, derived from the pod name in a small wrapper because it
does not take ordinal, and a QUEEN_TOKEN shared by all of them. In this cluster mode the facade
refuses Kafka transactions (a transaction lives in one node’s memory), so run Flink exactly-once or
transactional producers against a single node for now. See Kafka and its
ceilings.
Limits
Ready is not proof of a quorum. A pod that lost its majority, and an empty pod that never caught
up, can keep answering /health with 200, so alert on a write probe as well
(monitoring).
Delete one PVC at a time, and remove the member first. In a rehearsal, deleting two voters’ volumes
at once left a cluster that could not commit while every pod reported Ready (F11 in
test/recovery/FINDINGS.md); recovery has the procedure that brings such a
cluster back.
Two ephemeral-queue issues are open in 2.0.0-beta.6. A hand-over to a pod that is still starting can drop that partition’s messages, which we have seen during rolling restarts. And with the proxy on, an ephemeral request forwarded to the owning pod does not keep its tenant, so ephemeral queues behind the proxy are safe on a single node only, for now (ephemeral queues).
Port 7400 belongs to the pods. Keep it off the client Service and close it to everything else with a NetworkPolicy (security).