跳转到内容
简体中文

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.

destination regionsource regioncoordinatordestination regionsource regioncoordinatordrain entities + buffer messagesregion accepts shard,spawns entities as messages arriveHandOff for shardIdHandOffCompleteallocate new home

This is intentionally conservative — buffered messages wait for HandOffComplete before forwarding to the new owner, so nothing races past a half-moved entity.

Three paths:

  1. 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.
  2. Strategy-driven — every rebalanceIntervalMs (default 2s), the coordinator asks the allocation strategy for its rebalance recommendations. LeastShardAllocationStrategy returns shards to drain off busy nodes; HashAllocationStrategy returns shards whose hash-target moved.
  3. 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.

When the coordinator decides shard X should move from node A to node B:

  1. HandOff(X) is sent to node A’s region.
  2. Node A’s region marks shard X as handing-off. It still receives messages for entities in X, but buffers them instead of forwarding to entity actors.
  3. Each entity in X is sent the framework’s stop signal. They run postStop, persist any final state, and die.
  4. Once all entities in X are gone, node A sends HandOffComplete(X) to the coordinator.
  5. The coordinator runs allocate(X, ...) to pick a new home. Suppose it picks node B.
  6. The coordinator publishes the new owner via gossip. Buffered messages on node A start forwarding to node B’s region.
  7. Node B’s region receives messages for entities in X and 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.

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:

  1. 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.
  2. 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.

const shardingOptions = StartShardingOptions.create()
// ...
.withRebalanceIntervalMs(2_000) // strategy-driven rebalance every 2s
.withHandOffTimeoutMs(10_000); // give up + force-reallocate after 10s
sharding.start(shardingOptions);
KnobDefaultWhat
rebalanceIntervalMs2000How often the strategy is consulted for rebalance recommendations.
handOffTimeoutMs10_000If the source region doesn’t send HandOffComplete within this window, the coordinator gives up waiting and force-reallocates.

While shard X is handing-off:

  • Messages targeting entities in X arrive 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 X went 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.

For persistent entities (PersistentActor):

  • postStop on 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.

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 postStop that hangs (slow journal write, blocked external call).
  • A network partition during handoff.

Mitigation:

  • Keep postStop short and non-blocking.
  • Tune handOffTimeoutMs to your worst-case persist time + buffer.
  • For the split-brain risk, consider a lease on the sharding coordinator (see sharding with-lease).
Node added → coordinator allocates some shards to it
Node removed → coordinator re-homes its shards
Node ↑ in load → LeastShardAllocationStrategy drains shards off
Node ↓ 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.