Rebalance
此内容尚不支持你的语言。
When the cluster’s membership changes — a node joins or leaves — shards need to move to keep the workload distributed. The coordinator drives the process: it picks shards to relocate, tells the source region to hand off, and waits for confirmation before re-allocating.
This is intentionally conservative — buffered messages wait
for HandOffComplete before forwarding to the new owner, so
nothing races past a half-moved entity.
What triggers rebalance
Section titled “What triggers rebalance”Three paths:
- Membership-driven —
MemberUp(a new node),MemberRemoved(a leaving node), or any cluster transition that changes the candidate set. The coordinator re-runs the allocation strategy for every owned shard. - Strategy-driven — every
rebalanceIntervalMs(default 2s), the coordinator asks the allocation strategy for its rebalance recommendations.LeastShardAllocationStrategyreturns shards to drain off busy nodes;HashAllocationStrategyreturns shards whose hash-target moved. - A region stopping — a region tells the coordinator on its way down, and its shards are reallocated immediately. This covers a region stopped on a node that stays in the cluster, where no membership event fires at all.
The first two feed into the handoff protocol below. The third does not: the region is already gone, so there is nothing to hand off — the coordinator simply drops it from the registry and re-allocates.
Without that third path a stopped region left its shards orphaned: the coordinator kept naming it as their home, senders cached that home and delivered into a stopped actor, and the messages became dead letters. Nothing recovered it either — the candidate set is derived from the registry with no liveness check, so the rebalance tick saw a balanced cluster and moved nothing. It stayed that way until the node left the cluster entirely.
The notification is best-effort and idempotent: it is skipped if the
region never found a coordinator, a transport failure on the way down
is logged and swallowed, and the MemberRemoved path that follows a
whole-node shutdown finds the region already gone and does nothing.
The handoff sequence
Section titled “The handoff sequence”When the coordinator decides shard X should move from node A to
node B:
HandOff(X)is sent to node A’s region.- Node A’s region marks shard
Xashanding-off. It still receives messages for entities inX, but buffers them instead of forwarding to entity actors. - Each entity in
Xis sent the framework’s stop signal. They runpostStop, persist any final state, and die. - Once all entities in
Xare gone, node A sendsHandOffComplete(X)to the coordinator. - The coordinator runs
allocate(X, ...)to pick a new home. Suppose it picks node B. - The coordinator publishes the new owner via gossip. Buffered messages on node A start forwarding to node B’s region.
- Node B’s region receives messages for entities in
Xand spawns entities on demand (just like first-time allocation).
The whole thing usually takes sub-second to a few seconds,
depending on how many entities are in the shard and how long their
postStop work takes.
Who may order a handoff
Section titled “Who may order a handoff”A HandOff stops every entity in a shard, so the region checks
where it came from before acting on it — not just what it says.
The rule covers every directive only the coordinator may issue:
HandOff, ShardHome, RememberedEntities, ShardMapUpdate and
RegisterAcknowledgment.
Two conditions have to hold:
- The frame arrived through the region’s own envelope route, which stamps it with the peer the connection was authenticated as — never a field out of the payload, which the sender writes.
- That peer is the node the region currently believes hosts the coordinator, i.e. the current leader.
Anything else is dropped and logged at WARN. Until this rule
existed, any peer that could complete the cluster handshake could
send one small frame and take a shard’s entities down, repeatably;
the region dispatched on the message’s kind alone and had no
notion of a sender at all.
Refusing is safe by construction: a region re-registers and re-asks for every buffered shard on the next leader or membership event, so a directive dropped during a leadership handover is re-requested rather than lost.
A third condition is about the shard rather than the sender: the
region only hands off a shard it currently owns, and only once.
A HandOff it has already begun acting on, or one naming a shard
that lives somewhere else, is ignored — previously such a frame
still emptied the region’s cached “shard X lives on node N”
entries, so an authentic directive arriving twice, or arriving
late after the coordinator had already force-reallocated the shard,
cost every other node’s traffic an extra round trip to the
coordinator. Ownership is judged by allocation, not by whether the
shard actor happens to be running: a shard that passivated for being
idle is still owned and still hands off.
The coordinator does not stall on the refusal. It only ever sends
HandOff to the region its own allocation map names, so a refusal
means the two disagree about ownership — and handOffTimeoutMs
below is exactly the fallback for that, force-reallocating the shard
instead of waiting forever.
Configuration
Section titled “Configuration”const shardingOptions = StartShardingOptions.create() // ... .withRebalanceIntervalMs(2_000) // strategy-driven rebalance every 2s .withHandOffTimeoutMs(10_000); // give up + force-reallocate after 10ssharding.start(shardingOptions);| Knob | Default | What |
|---|---|---|
rebalanceIntervalMs | 2000 | How often the strategy is consulted for rebalance recommendations. |
handOffTimeoutMs | 10_000 | If the source region doesn’t send HandOffComplete within this window, the coordinator gives up waiting and force-reallocates. |
What buffered messages mean
Section titled “What buffered messages mean”While shard X is handing-off:
- Messages targeting entities in
Xarrive at node A. - Node A’s region buffers them — doesn’t forward to the (now-stopped) entity actors.
- Once the handoff completes, node A asks the coordinator where
Xwent and forwards the buffer to the new owner as soon as it hears back. - The new region spawns entities and processes the messages in order.
This means a message sent during handoff is delayed, not dropped. The cost: latency spikes during rebalance windows.
The re-ask matters more than it looks. Completing a handoff clears the old region’s cached shard home, and the coordinator only announces a new placement to the new owner and to regions with an outstanding query — so without asking, the old region would sit on its buffer until some unrelated later message for the same shard happened to trigger a lookup. On a shard that goes quiet after the rebalance, that never comes.
When entities have state
Section titled “When entities have state”For persistent entities (PersistentActor):
postStopon the source side finalizes any pending persist.- The new entity instance on the destination replays the journal on startup.
- The buffered message arrives at a fully-recovered entity.
For non-persistent entities, state is lost between incarnations — the rebalance is equivalent to a restart with the buffered message acting as the first command.
If state matters across rebalance, use persistence. This isn’t optional — and it carries a prerequisite of its own: the destination replays whatever journal its node opens, so the promise only holds when every node reads the same database. Per-node storage (a SQLite file, the in-memory defaults) turns a rebalance into a silent state fork instead — the entity recovers a different history with no error anywhere. The cluster warns about that combination; see storage locality & identity.
Force-reallocation
Section titled “Force-reallocation”If HandOffComplete doesn’t arrive within handOffTimeoutMs:
- The coordinator logs a warning (“handoff timed out”).
- It force-reallocates the shard to its new owner.
- Buffered messages forward as if handoff had completed normally.
- The old entities may still be alive on node A — they receive messages locally too, leading to split-brain entities until the orphans die.
This is rare but possible. Causes:
- A
postStopthat hangs (slow journal write, blocked external call). - A network partition during handoff.
Mitigation:
- Keep
postStopshort and non-blocking. - Tune
handOffTimeoutMsto your worst-case persist time + buffer. - For the split-brain risk, consider a lease on the sharding coordinator (see sharding with-lease).
Rebalance vs scaling
Section titled “Rebalance vs scaling”Node added → coordinator allocates some shards to itNode removed → coordinator re-homes its shardsNode ↑ in load → LeastShardAllocationStrategy drains shards offNode ↓ in load → no rebalance (only the busy direction triggers)Rebalance is work-shedding, not work-acquiring. An idle node
in LeastShardAllocationStrategy’s view receives new shards as
they’re allocated, but doesn’t get existing shards moved into it
just because it’s quiet.
For an “always-balanced” effect, restart the coordinator
periodically (which re-runs allocation from scratch) — or
configure rebalanceThreshold: 1 to react to any imbalance.
Where to next
Section titled “Where to next”- Sharding overview — the broader picture.
- Allocation strategy — what decides the new homes.
- Remember entities — affects what gets re-spawned after handoff.
- Sharding with lease — split-brain protection for the coordinator.
