Skip to main content

Replication Factor Is Not Fault Tolerance: Sizing a Zaris Cluster for Durability on Kubernetes

· 8 min read
Clustron Team
Distributed Systems Engineering

Sizing a Zaris cluster for durability on Kubernetes

There is a tempting line of reasoning when you deploy a replicated store on Kubernetes: I set replicationFactor: 2, so every partition has a backup, so I survive a node loss. It feels airtight. It is also, for a surprisingly common cluster size, wrong — and wrong in the quiet way that only shows up the first time a pod actually dies.

Zaris on Kubernetes stores data in memory and makes it durable by replication, not disk: the node data volume is an emptyDir, and a rescheduled pod rebuilds from the copy held on another pod. That design is sound — but it leans entirely on there being a copy on another pod. Whether that copy exists is not decided by replicationFactor alone. It is decided by the arithmetic between your pod count and your replication factor, and if you get that arithmetic wrong the store will run perfectly, serve every read and write, and silently hold one partition that has nowhere to replicate to.

The deployment shape​

The Zaris Helm chart deploys the node tier as a StatefulSet, one node instance per pod (zaris-0, zaris-1, …), each an in-memory shard-owner. A headless Service gives each pod a stable DNS identity, readiness is gated on cluster convergence, and a PodDisruptionBudget plus a graceful-drain termination window protect the set during voluntary disruptions.

# values.yaml (defaults)
node:
replicas: 4 # pod count — ONE node instance per pod
replicationFactor: 2 # copies of each partition, on distinct pods

The chart's own notes are blunt about where durability comes from:

Node data is emptyDir. In-memory durability is via replication (RF), not disk; a rescheduled node rebuilds from replicas.

So a node pod can be killed and rescheduled all day with no data loss — provided its partitions have a live replica on a different pod to rebuild from. That proviso is the whole ballgame.

How partitions land on pods​

A partition in Zaris is held by RF distinct node instances — one active owner plus RF−1 replicas — and the replicas are always placed on different pods (a replica on the same pod as its active would be no protection at all). The runtime places partitions greedily:

The runtime fills a partition to its full RF before forming the next one.

That one sentence is the key to the whole sizing story. Instances are consumed RF at a time. So the number of fully-replicated partitions you get is:

fully-replicated partitions = ⌊ instances ÷ RF ⌋

and any instances left over — instances mod RF of them — form a single-copy partition with no peer to replicate to. It is a real, serving partition. It just has no backup.

With four pods at RF 2 you get two partitions, each with two copies on distinct pods — lose any one pod and the surviving copy carries on while the partition re-replicates onto the rescheduled pod. With three pods at RF 2, pods 0 and 1 form partition A at full RF, but pod 2 is left over and forms partition B with a single copy. Add a fourth pod and it becomes B's replica. Until then, B is a liability.

The sizing rule, as a table​

Because each partition needs RF distinct pods and Kubernetes runs one instance per pod, the rule is just:

Run a pod count that is a multiple of RF. With RF 2, use an even number of pods.

Pods (1 instance each)RFFully-replicated partitionsSurvives one pod loss?
221Yes
321 (+ 1 single-copy)No — losing the single-copy pod loses that partition
422Yes
623Yes

The pattern generalises: N partitions at RF 2 needs 2N pods; at RF 3 it needs 3N. Scaling out means adding pods a full RF at a time if you want every new partition to be born durable.

Why three pods with RF 2 is the trap​

Three is the number that bites, because three pods looks like a healthy, highly-available cluster. It passes readiness. It serves traffic. kubectl get pods shows zaris-0..2 all 1/1 Ready. Nothing is obviously wrong.

What is wrong is invisible until the single-copy pod — the one holding partition B's only copy — is the pod that dies. At that instant B's keys are gone, permanently, because there was never a second copy to rebuild from. This is not a bug; it is the same caveat as running a single machine, just scoped to one partition instead of the whole store. The map only regenerates when all copies of a partition are lost, and for a single-copy partition "all copies" is one pod.

This was measured, not assumed. On a kind cluster:

  • At replicas=4 (2 partitions × RF 2), killing a partition's active-owner pod lost 0 of 200 keys — the replica took over and RF self-restored onto the rescheduled pod.
  • At replicas=3, killing the single-copy pod lost that partition's keys outright.

Same store, same RF setting, same chaos. The only variable was whether the pod count was a multiple of RF.

The guardrails the chart gives you​

The chart does not leave this to folklore. Install a cluster whose pod count is not a multiple of RF and helm install prints a durability warning in its notes:

⚠  DURABILITY WARNING
node.replicas (3) is not a multiple of node.replicationFactor (2), so at least one
partition is SINGLE-COPY and will LOSE its data if that pod is lost. For fault
tolerance set node.replicas to a multiple of 2 (e.g. 4).

Beyond the warning, several defaults are set so that a correctly-sized cluster behaves well under disruption:

  • Convergence-gated readiness. /readyz passes only once the partition map and ownership are live, so a Ready pod is a routable one — Kubernetes won't send traffic to a node that hasn't taken ownership yet.
  • PodDisruptionBudget. By default at most one node is voluntarily disrupted at a time, so a drain or rolling upgrade can't take out both copies of a partition at once.
  • Graceful-drain termination. A termination grace window lets a leaving node hand off cleanly rather than being torn down mid-flight.
  • Hot scale-out. With Kubernetes discovery on (the default), the roster is synthesised from pod ordinals and headless DNS, so adding pods is a plain kubectl scale statefulset — new pods join and the partition map grows without restarting the existing ones. Add them a multiple of RF at a time and every new partition is durable from birth.

None of these help the single-copy partition, though. A PDB that keeps all-but-one pod up is no protection when one of those pods is the only home of a partition. The guardrails assume you sized the cluster correctly; they don't rescue a cluster that wasn't.

Practical rules​

The whole topic compresses to a few habits:

  1. Pick RF first, then make the pod count a multiple of it. RF 2 → even pod counts; RF 3 → multiples of three. The default chart ships 4 pods at RF 2 precisely because it's the smallest fault-tolerant shape.
  2. Treat an odd pod count as evaluation-only. Three pods at RF 2 is fine for a demo on kind; it is not a production durability posture, and the chart tells you so at install time.
  3. Scale out by RF. Adding one pod to a durable cluster temporarily creates a single-copy partition until the next pod arrives. If you're growing for durability, grow a full replication set at a time.
  4. Remember the volume is emptyDir. There is no disk to fall back on. Replication is your durability, which is exactly why the copy has to actually exist on another pod.

Replication factor tells you how many copies a partition wants. The pod count decides whether it can have them. Line those two numbers up and a pod loss is a non-event — the replica serves, the rescheduled pod rebuilds, and not a key is lost. Let them drift out of alignment and you've built a cluster that is highly available right up until the moment it matters.