Skip to content

Scale out and in

Add nodes to a running Narad cluster, or drain nodes and remove them, without stranding any partition.

Before you start: Helm access to the release, every pod ready, and the narad CLI with admin credentials in its active context or in NARAD_ADDR, NARAD_USER and NARAD_PASS (CLI command reference).

Check members and moves

Two commands show where partitions live and what is moving:

narad cluster members
Output
{
  "members": [
    {
      "id": "narad-1",
      "addr": "127.0.0.1:18181",
      "status": "alive",
      "draining": false,
      "owned_partitions": 4,
      "outbound_moves": 0,
      "voter": true,
      "leader": false,
      "heartbeat_age_seconds": 3
    },
    {
      "id": "narad-2",
      "addr": "127.0.0.1:18182",
      "status": "alive",
      "draining": false,
      "owned_partitions": 4,
      "outbound_moves": 0,
      "voter": true,
      "leader": false,
      "heartbeat_age_seconds": 4
    },
    {
      "id": "narad-3",
      "addr": "127.0.0.1:18183",
      "status": "alive",
      "draining": false,
      "owned_partitions": 4,
      "outbound_moves": 0,
      "voter": true,
      "leader": true,
      "heartbeat_age_seconds": 4
    }
  ]
}

This output comes from a three-node test cluster on one machine, built from master; on Kubernetes, addr holds each pod's address. v3.0.1 prints the first six fields only.

  • status is alive, or dead once a node has sent no heartbeat for about 30 seconds.
  • draining is true while a decommission sheds the node's partitions.
  • owned_partitions counts the partitions the node owns, and outbound_moves those it is copying to another node.
  • voter and leader (v3.1.0) are the node's place in Raft, and heartbeat_age_seconds how long ago its last heartbeat was recorded.
  • decommission_blocked (v3.1.0) appears on a draining node whose decommission cannot progress, with a code and a message per reason (Troubleshooting).

narad cluster members --detail (v3.1.0) also asks every node for its own status: its dispatch backlog (messages its ingress WAL still has to hand to their owners), the partition copies it set aside, and the moves it runs. A node that cannot answer gets a status_error; a v3.0.1 node is reported as an older release.

narad cluster moves lists the partitions being moved right now, and prints an empty moves list when nothing moves. From v3.1.0, each move also shows from_status and to_status (alive, dead, draining or not_a_member) and, when it cannot progress, blocked; --detail adds the destination's own report of the move. Both commands need the admin grant.

Scale out

Raise replicaCount:

helm upgrade narad ./charts/narad -n narad --reuse-values \
  --set replicaCount=5

Each new pod starts empty and finds it is not one of the initial members, so it asks the leader to admit it instead of creating a cluster of its own. It does not need the leader in its peer list: a follower's answer names the leader, and the new pod asks it next, so this works after leadership has moved to a pod outside the pinned list. The existing pods are not restarted: the peer list in their pod template stays pinned to the initial members.

The leader admits a new pod as a Raft non-voter: it replicates the cluster's metadata and serves traffic, but does not count toward quorum, so a pod that cannot be reached costs the cluster nothing. Once its copy has caught up it asks again, and the leader promotes it to voter, normally within seconds (and never before the leader has led for 12 s). Until then narad_raft_nonvoters is above 0. If it stays above 0, the new pod logs why the leader defers its promotion (Troubleshooting).

Once a new node is admitted, the leader rebalances. It computes the fewest partition moves that even out the number of partitions per node, copies each partition to its new owner, and switches ownership over at the end. At most 8 moves run at once, so a large rebalance drains gradually. Watch it with narad cluster moves until the list is empty. How a move copies and hands over a partition is in Rebalance and decommission.

If Helm refuses with a conflict on .spec.replicas, someone changed the replica count outside Helm: see Troubleshooting.

Scale in

One copy per partition

Narad stores each partition once, on the volume of the node that owns it. A crash, restart or power loss loses nothing that got a 202; losing a volume loses the partitions on it. For a second copy, add a replica child or take volume snapshots, as Back up and replicate topics shows.

A pod removed while it still owns partitions strands them: they stay assigned to a node that no longer runs, and their messages cannot be consumed until it returns. The removed pod also stays a Raft voter, so the cluster tolerates one fewer failure. So scaling in is two steps: decommission the node until it owns nothing, then lower replicaCount. The StatefulSet always removes the highest-numbered pods, so decommission those.

Decommission a node

Scale in: drain narad-4, then delete it You run narad cluster decommission narad-4 on a five-node cluster. narad-4 is marked draining, so it receives no partitions, and the leader targets each partition it owns at the other four pods, and each destination copies its share: orders/4 to narad-0, orders/9 to narad-1, orders/14 to narad-2 and orders/19 to narad-3, at most 8 moves at once. Once narad-4 owns 0 partitions, the leader removes it from the Raft voters, never going below 3 voters. Then you run helm upgrade with replicaCount=4 and allowScaleIn=true, and the StatefulSet deletes narad-4. Lower replicaCount before step 2 and the partitions still on narad-4 are stranded. 1 narad cluster decommission narad-4 at most 8 moves at once narad-4 draining Raft orders/4 orders/9 orders/14 orders/19 narad-0 narad-1 narad-2 narad-3 2 owned_partitions: 0 leader removes it from Raft never below 3 voters narad-4 deleted by the StatefulSet 3 helm upgrade ... replicaCount=4 allowScaleIn=true Lower it before step 2 and its partitions are stranded StatefulSet narad Scale in: drain narad-4, then delete it You run narad cluster decommission narad-4 on a five-node cluster. narad-4 is marked draining, so it receives no partitions, and the leader targets each partition it owns at the other four pods, and each destination copies its share: orders/4 to narad-0, orders/9 to narad-1, orders/14 to narad-2 and orders/19 to narad-3, at most 8 moves at once. Once narad-4 owns 0 partitions, the leader removes it from the Raft voters, never going below 3 voters. Then you run helm upgrade with replicaCount=4 and allowScaleIn=true, and the StatefulSet deletes narad-4. Lower replicaCount before step 2 and the partitions still on narad-4 are stranded. StatefulSet narad narad-0 narad-1 narad-2 narad-3 orders/4 orders/9 orders/14 orders/19 narad-4 draining Raft 1 at most 8 moves at once 2 narad-4 3 deleted by the StatefulSet 1 narad cluster decommission narad-4 2 owned_partitions: 0 leader removes it from Raft never below 3 voters 3 helm upgrade ... replicaCount=4 allowScaleIn=true Lower it before step 2 and its partitions are stranded
Drain first, delete second: narad-4 hands every partition it owns to the other pods and leaves the Raft voters at owned_partitions: 0, and only then does a lower replicaCount remove the pod. The four moves are the ones from the example below.

This example takes a five-node cluster down to four.

  1. Check that the node can be removed safely, then mark it for decommission:

    narad cluster decommission narad-4 --dry-run
    narad cluster decommission narad-4
    

    The dry run (v3.1.0) changes nothing; it says whether the decommission would be accepted and, if not, every reason. A decommission that could never complete safely is refused with 409 and the same reasons, and nothing changes: fewer than three voters would remain (below_min_voters), the voters left alive would not be a majority (no_healthy_majority), no other alive node could take the node's partitions (no_receivers), or the node is dead and owns partitions (owner_dead).

    The node stops receiving partitions and its partitions start moving to the others. It also refuses new produce with 503 and Retry-After: 1 (from v3.1.0), so clients send their messages to another node. While its partitions move, narad cluster moves lists them:

    Output
    {
      "moves": [
        {
          "topic": "orders",
          "partition": 4,
          "from": "narad-4",
          "to": "narad-0"
        },
        {
          "topic": "orders",
          "partition": 9,
          "from": "narad-4",
          "to": "narad-1"
        },
        {
          "topic": "orders",
          "partition": 14,
          "from": "narad-4",
          "to": "narad-2"
        },
        {
          "topic": "orders",
          "partition": 19,
          "from": "narad-4",
          "to": "narad-3"
        }
      ]
    }
    

    This output comes from a five-node test cluster on one machine, with one topic of 20 partitions.

  2. Wait until narad cluster members shows owned_partitions: 0 for narad-4. The leader then removes it from Raft (as a voter, or as a non-voter if it was never promoted), and it drops out of the list. The pod keeps running, reports not ready, and its heartbeats are refused. Before the removal, the leader asks the node for its dispatch backlog and waits until its ingress WAL has handed every message it accepted to the partition's owner (from v3.1.0); a node on v3.0.1 cannot answer and is removed without that check, as v3.0.1 did. If the node stays listed, its decommission_blocked says why (Troubleshooting).

  3. Lower replicaCount. allowScaleInTo=4 tells the chart the pods above size 4 were decommissioned; without it, the chart refuses to lower the replica count of a running StatefulSet. It approves that one size only, so keeping it with --reuse-values does not approve a later scale-in to another size.

    helm upgrade narad ./charts/narad -n narad --reuse-values \
      --set replicaCount=4 \
      --set allowScaleInTo=4
    

    allowScaleInTo is new in v3.1.0: the v3.0.1 chart takes --set allowScaleIn=true instead, and newer charts no longer read allowScaleIn.

    Before the change, the chart's scale-in guard (v3.1.0) checks the cluster itself. This hook Job runs before every helm upgrade and helm rollback. If the change deletes a pod that narad cluster members still lists, it refuses, and the command fails with the reason in the Job's log:

    kubectl logs -n narad job/narad-scale-in-guard
    

    A change that deletes no pod passes without calling the API. The guard finds pods through cluster DNS, and a lookup that fails looks the same for a pod that is gone and for a DNS outage, so it trusts one only while cluster DNS answers for the API server's Service (kubernetes.default.svc) or a pod the change keeps; otherwise it refuses with DNS lookups fail; cannot tell which pods exist. A release with no pod yet (a first Argo CD sync, or every pod Pending without an address) therefore passes: the change deletes none. On a scale-in the guard signs in as admin with the admin-password key of the security secret, so that key must hold root's current password; without it, the guard refuses and says so. An answer that lists no members counts as unreadable too. To go ahead after checking narad cluster members by hand, add --no-hooks to that one command.

To stop a decommission before it finishes, run narad cluster decommission narad-4 --cancel. The node starts receiving partitions again, and the next rebalance evens the load out.

Keep these rules while you scale in:

  • Wait for zero partitions. Lowering replicaCount before the node owns nothing deletes a pod whose data has not moved.
  • Do not overlap a decommission with a rolling restart. A helm upgrade that changes the pod template restarts the pods, and a draining node that restarts has no stable source to copy from until it settles. Changing only replicaCount does not restart the pods.
  • A rollback of replicaCount is a scale-in. helm rollback to a revision with fewer replicas deletes pods exactly like step 3, without steps 1 and 2. helm rollback renders no templates, so the allowScaleInTo check never runs; the scale-in guard, a pre-rollback hook, refuses it while a pod being deleted is still a member. A rollback to a revision rendered by an older chart runs no guard (the v3.0.1 chart refuses nothing on rollback). Decommission first.
  • Keep a node that is not draining alive. New partitions never go to a draining node. While every live node is draining, for example two draining nodes while the others restart, topic creates and partition increases answer 503 and the leader logs every live member is being decommissioned at error level; they succeed again once a node that is not draining is back.
  • Three voters is the floor. The leader never removes a voter from Raft if that would leave fewer than three voters, and from v3.1.0 also never if the voters left alive would not be a majority of the rest: with dead voters around, removing a live one could leave a configuration that can never elect a leader. From v3.1.0 such a decommission is refused up front; on v3.0.1 it moves the partitions away, but the node stays a member. If the leader itself is decommissioned, it hands leadership to another node first. A node that is still a non-voter has no vote, so the floor does not hold it back.
  • Remove dead voters first, and bring them back to do it. From v3.1.0 the leader removes a dead draining voter before a live one, but only once it has read the node's dispatch backlog, which a dead node cannot report: its decommission waits with node_status_unavailable until the node comes back and hands off its WAL. Bring it back, or cancel its decommission. A dead node that still owns partitions cannot be drained at all (owner_dead): its data is only on its disk.
  • A move aimed at the node holds it back. From v3.1.0 the leader clears moves that target a draining node and removes the node only once none is left. A move you want gone sooner can be aborted with narad cluster moves abort <topic> <partition> (Rebalance and decommission). The leader does not clear a move whose source is dead or no longer a member: the draining node may hold the only live copy, so it waits for that copy to be force-promoted and moved off, and logs the wait at error. Abort such a move only if dropping that copy is really intended.
  • Finish a rolling upgrade from 3.0.x before you decommission a non-voter. A 3.0.x leader removes only voters, so a non-voter it decommissions loses its member record but stays in the Raft configuration, and narad_raft_nonvoters stays above 0 (Troubleshooting).

Reuse a decommissioned name

The leader remembers a decommissioned node's ID. If that pod restarts with its old volume, it is refused and logs that it was decommissioned (Troubleshooting). To bring a pod back under that name later, for example when you scale out again, delete its PersistentVolumeClaim first so it starts empty. If the leader logged the decommissioned node holds quarantined partition copies for it, check those copies first (Troubleshooting): deleting the claim deletes them, and they may hold the only instance of some records.

kubectl delete pvc data-narad-4 -n narad

A node that starts with an empty data directory under a decommissioned ID is admitted again as a new node.

Source node failure during a move

If a node dies while its partitions are still being copied away, the destinations that had caught up with it promote their copies after 2 minutes instead of waiting for it. The 2 minutes count from the source's last heartbeat and, from v3.1.0, also on each destination's own clock from when it first saw the source dead, so a destination that restarts meanwhile waits up to 2 minutes more. A destination whose copy is behind cannot promote it: it waits for the source, logs that once at error and counts the move in narad_moves_blocked (Troubleshooting). This force-promote can deliver again messages that consumers acked on the dead node in its last moments. From v3.1.0 it never skips one; on v3.0.1 a force-promote could leave records committed after it undelivered (see Rebalance and decommission). What can be lost, and why, is in the failure matrix.

Next steps