Skip to content

Kubernetes

A StatefulSet of three voters with a volume each, a headless Service for raft, a client Service and a disruption budget: the manifest, the probes, what SIGTERM does, rolling upgrades, and what is fixed at the first boot.

Updated View as Markdown

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:

  1. 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.
  2. 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.
  3. It stops accepting connections and lets the requests in flight finish. A parked long-poll pop can wait 30 seconds by default.
  4. The Kafka facade, when on, gives its registry row back (5 seconds), and the proxy, when on, writes its last usage minute (5 seconds).
  5. 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/queen

A 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).

Navigation

Type to search…

↑↓ navigate↵ selectEsc close