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
{
"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.
statusisalive, ordeadonce a node has sent no heartbeat for about 30 seconds.drainingistruewhile a decommission sheds the node's partitions.owned_partitionscounts the partitions the node owns, andoutbound_movesthose it is copying to another node.voterandleader(v3.1.0) are the node's place in Raft, andheartbeat_age_secondshow long ago its last heartbeat was recorded.decommission_blocked(v3.1.0) appears on a draining node whose decommission cannot progress, with acodeand amessageper 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¶
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.
-
Check that the node can be removed safely, then mark it for decommission:
narad cluster decommission narad-4 --dry-run narad cluster decommission narad-4The 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
409and 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
503andRetry-After: 1(from v3.1.0), so clients send their messages to another node. While its partitions move,narad cluster moveslists 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.
-
Wait until
narad cluster membersshowsowned_partitions: 0fornarad-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, itsdecommission_blockedsays why (Troubleshooting). -
Lower
replicaCount.allowScaleInTo=4tells 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-valuesdoes not approve a later scale-in to another size.helm upgrade narad ./charts/narad -n narad --reuse-values \ --set replicaCount=4 \ --set allowScaleInTo=4allowScaleInTois new in v3.1.0: the v3.0.1 chart takes--set allowScaleIn=trueinstead, and newer charts no longer readallowScaleIn.Before the change, the chart's scale-in guard (v3.1.0) checks the cluster itself. This hook Job runs before every
helm upgradeandhelm rollback. If the change deletes a pod thatnarad cluster membersstill lists, it refuses, and the command fails with the reason in the Job's log:kubectl logs -n narad job/narad-scale-in-guardA 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 withDNS 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 asadminwith theadmin-passwordkey 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 checkingnarad cluster membersby hand, add--no-hooksto 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
replicaCountbefore the node owns nothing deletes a pod whose data has not moved. - Do not overlap a decommission with a rolling restart. A
helm upgradethat changes the pod template restarts the pods, and a draining node that restarts has no stable source to copy from until it settles. Changing onlyreplicaCountdoes not restart the pods. - A rollback of
replicaCountis a scale-in.helm rollbackto a revision with fewer replicas deletes pods exactly like step 3, without steps 1 and 2.helm rollbackrenders no templates, so theallowScaleInTocheck 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
503and the leader logsevery live member is being decommissionedat 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_unavailableuntil 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_nonvotersstays 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¶
- Rebalance and decommission: how a partition moves without losing a record.
- Back up and replicate topics: check replica placement after the cluster changes shape.
- Capacity and disk sizing: decide how many nodes the load needs.