Adaptive Shard Splitting
This page is part of the documentation for Orleans.Lattice 9.9.0 (release line 9.9), built 2026-10-04. It is also published as markdown, with every table and list, at shard-splitting.md, and llms.txt lists every page.Adaptive shard splitting allows a hot physical shard to split into two at
runtime, fully online - no shard is ever taken offline. Splits happen
automatically when an autonomic monitor detects a hot shard, and an
online reshard drives the same split to grow a tree's
shard count. Shard splitting is internal-only: ITreeShardSplitGrain is declared internal
and is not reachable from consumer assemblies.
Why
Lattice trees are sharded by hashing keys into a virtual slot space and
mapping virtual slots onto physical ShardRootGrain activations. With a
fixed shard count, a workload skewed toward a small set of keys will
saturate one shard while others sit idle. Adaptive splitting redistributes
hot virtual slots to a new physical shard so the load follows the data.
How it works
A split is driven by the internal TreeShardSplitGrain coordinator through
five persisted phases, in this order. The source shard S keeps serving reads and
writes throughout; the target shard T receives mirrored data and eventually
owns the moved slots.
stateDiagram-v2
[*] --> BeginShadowWrite : split requested for source shard S
BeginShadowWrite --> Drain : S mirrors moved-slot writes to T, prepared mutations swept across
Drain --> Swap : every moved-slot entry (live + tombstones) forwarded to T
Swap --> Reject : S sealed and rejecting, final drain, registry reassigns the moved slots to T
Reject --> Complete : reject state re-asserted
Complete --> [*] : final drain pass, then S records the moved slots in its moved-away table
- BeginShadowWrite - Coordinator takes the upper half of the virtual
slots S owns as the moved set, allocates a fresh target shard index from
the registry, persists its intent, and opens S's shadow-write window for
the moved slots. From that moment every successful write S applies to a
key in a moved virtual slot is also mirrored to T through T's batched
merge, preserving the original HLC. Before the drain begins, the
coordinator then runs a retroactive prepared-mutation sweep: it walks
S's leaf chain, snapshots every in-flight prepared saga mutation whose
key hashes into a moved virtual slot, and replays each one into T's
pending-transaction buckets under its original transaction id, HLC,
origin, vector clock, and expiry, so any prepared write that landed on S
before the window opened survives the topology change. A saga the
transaction registry already reports as committed or aborted has its
terminal applied to T directly instead. The sweep is idempotent per
(transaction, key), and a coordinator crash mid-sweep re-runs the whole sweep on recovery. Instrumentation:orleans.lattice.split.retroactive_forward.entries(counter, per replayed mutation) andorleans.lattice.split.retroactive_forward.duration(histogram, total sweep wall-clock). CRDT LWW guarantees correct convergence regardless of how the foreground write and the background drain interleave. - Drain - Coordinator walks S's leaf chain and forwards moved-slot
entries (including tombstones) to T with their original HLC timestamps.
The drain is chunked and leaf-side filtered: each leaf returns
only entries whose virtual slot is in the moved-slot set, and the
coordinator flushes to T in batches of
SplitDrainBatchSize(default 1024) entries. This bounds peak memory on the coordinator regardless of source shard size, and avoids transferring non-moved entries over the wire. Each pass is bounded by the background-drain budget (BackgroundDrainLeavesPerPassleaves orBackgroundDrainMaxDuration, whichever binds first) and persists a key cursor the next tick resumes from, so the phase advances to Swap only once the whole leaf chain has been swept. Idempotent under retry - re-running merges only converges to the same state. - Swap - Coordinator seals every source leaf for the moved slots and
puts S into its reject phase before the registry map flips. From
this point any read or write to S for a moved-slot key is refused as
stale routing, which freezes the source's committed state for the
migrating slots. The coordinator then runs one final authoritative drain
pass to T, so the destination is synchronised with the source's
now-frozen committed state before any reader can route to T. Only then
does it reassign the moved slots to T in the registry's
ShardMap, with a single registry call that re-reads the live map, applies the reassignment, and persists it under a freshVersion, so a concurrent split or shard consolidation of the same tree cannot erase either change. Reversing the order - flipping the map before the source rejects - would open a window in which a stale-routing reader could still be served the pre-split value by the source. New router activations immediately route the moved slots to T; stale activations that still cache the old map hit the source's reject gate, invalidate their cached map, fetch the fresh map from the registry, and retry against T - a single transparent retry per call. - Reject - Coordinator re-asserts S's reject phase. S already entered it during Swap and the call is idempotent, so this phase only records the transition; a crash-recovered coordinator that re-enters it changes nothing.
- Complete - Coordinator runs one final drain pass to capture any tombstones written during shadow that were not mirrored on the hot path, then tells S to complete the split and clears its own state. Completing also promotes the split's moved slots into S's persisted moved-away slot table, so even after the active reject-phase state is cleared, every subsequent operation on a moved-slot key continues to be refused as stale routing. This guarantees that stale stateless-worker router activations (which may have cached the pre-split shard map) always trigger a map refresh on first use rather than silently returning orphan data. The table is lifted only if a later shard consolidation folds T back into S - and only after T's entries for those slots have been drained onto S - at which point S's leaves re-ship the reclaimed rows to their read caches.
The coordinator state is persisted before any side effect, so a silo crash mid-split is recovered by the keepalive reminder, which resumes from the last persisted phase; every phase is idempotent.
Scan semantics during a split
This section describes the mechanism by which live operations behave
during a split. For the consistency contract each ILattice method
provides - including under concurrent splits - see
Consistency.
Point reads and writes (GetAsync, SetAsync, DeleteAsync,
SetIfVersionAsync, GetOrSetAsync, etc.) continue to serve traffic
throughout the split: every successful write is mirrored to the new
owner during the shadow phase and the reject phase (which precedes the map
swap) causes stale activations to transparently retry against the correct
shard. The
post-Complete MovedAwaySlots rejection extends this for as long as the
source does not own those slots again (only a later shard consolidation
that folds them back lifts it).
Scans (ScanKeysAsync, ScanEntriesAsync, CountAsync) reconcile against
topology changes mid-scan as described below. See
Consistency for the guarantee this reconciliation
delivers.
How the reconciliation works
Each scan uses a reconciliation algorithm coordinated against the
registry's monotonically-incrementing ShardMap.Version, but CountAsync
and the ScanKeysAsync / ScanEntriesAsync streams follow two different paths.
CountAsync / CountPerShardAsync - per-slot routing
The orchestrator reads the authoritative ShardMap. While its Version
is still 0 - no split has ever been persisted for the tree - it sums
each physical shard's own count and accepts the total only if the map is
still at version 0 afterwards. Otherwise it partitions the virtual slots
by current owner and asks each physical shard to count only the slots it
owns (bounded to the requested range), driving each shard's work-bounded
count batches to completion so no shard is held for its whole leaf chain.
Because each virtual slot is counted exactly once - against whichever
shard the map identifies as its current owner - the result is
topology-consistent by construction, independent of the source shard's
per-split phase. The map version is re-read after the fan-out; if it
moved, the count is discarded and retried on the fresh map, bounded by
LatticeOptions.MaxScanRetries (default 3). Throws
InvalidOperationException on retry exhaustion.
ScanKeysAsync / ScanEntriesAsync - in-line reconciliation
Reconciliation is driven inside the main k-way merge loop rather than as a separate pass. Each shard root reports back:
- the keys/entries of all keys not in its
MovedAwaySlotstable (entries it no longer authoritatively owns), and - the set of
MovedAwaySlotsvirtual slots it observed during the traversal (used as a topology-stability hint).
Before each priority-queue dequeue, the orchestrator checks whether
any live shard cursor has reported new MovedAwaySlots since the last
reconciliation step. If so, it queries the current owners of the
affected slots via the slot-filtered variants
GetSortedKeysBatchForSlotsAsync / GetSortedEntriesBatchForSlotsAsync,
loads the reconciled keys into memory, sorts them with the same
comparer, and injects them as an additional in-memory cursor into the
same priority queue. The merge invariant (global minimum is yielded
next) then carries ordering across the topology boundary. A per-call
HashSet<string> suppresses duplicates across pre- and post-swap
views. A final stability check after the priority queue drains catches
the edge case where a split commits after all live cursors finished -
reconciled entries from this path are also sorted and injected as a
cursor, not appended. Bounded by LatticeOptions.MaxScanRetries.
Snapshot scans - pinned-map slot ownership in the snapshot leaf
A snapshot-isolated scan (OpenSnapshotEntryCursorAsync /
OpenSnapshotKeyCursorAsync, the read-only state API, and the Explorer's
Data area) cannot use the live in-line reconciliation above, because it
must read a single, internally-consistent point in time. Instead the
snapshot coordinate pins the registry's ShardMap at open
(LatticeSnapshotCoordinate.PinnedShardMap, alongside the pinned
ShardMap.Version, the per-shard / per-partition WAL offsets, and the
registry HLC). The snapshot fan-out forces a fresh routing read so it
opens against the post-split map, then for each fan-out shard it
resolves that shard's owned virtual slots under the pinned map and
passes them to the shard's snapshot leaf.
The snapshot leaf serves a durable frozen baseline captured at open time
(see Snapshot Cursors): the per-shard projection
is materialised once by walking the shard's leaf chain and folding each
leaf's WAL tail, then persisted and seeded into the leaf with no
serve-time WAL replay. Whether a per-key record reaches the served view
is resolved by the key's virtual slot under the pinned map - not by
the mutation's stamped ShardIndex - applied both when the baseline is
folded at capture and again through the leaf's IsKeyOwned filter when
the durable rows are seeded. Resolving by slot is what makes the snapshot
view correct across a split, because the stamp records the shard that
authored a record, which is not the same as the shard that owns the
key after the split:
- A moved key's pre-split copy is physically retained on the donor shard (an orphan). Its stamp still names the donor, but the pinned map now assigns its slot to the target, so the donor's snapshot leaf drops it and only the target surfaces the key - no duplicate.
- A write routed to the donor for an already-moved slot is shadow-forwarded into the target shard's WAL but keeps the donor's source stamp. Resolving by slot keeps that forwarded record on the target's snapshot leaf, so the leaf's last-writer-wins merge applies it. A stamp-based filter would drop it (its stamp names the donor) and resurrect the pre-forward value - for example a post-split delete forwarded through the donor would be lost and the deleted key would reappear with its stale drained value.
Ownership must be resolved against the pinned map version, not the current one, so the snapshot neither over-excludes (a key not yet moved at the pinned version) nor under-excludes (a key moved after the pinned version). The snapshot k-way merge additionally collapses equal, adjacent keys as a defensive net, but value correctness comes from the leaf-side ownership filter (the merge operates on raw bytes with no HLC and cannot pick the last-writer-wins winner on its own).
Live leaf reactivation - current-map slot ownership in the live leaf
The authoritative live leaf has its own activation-time WAL replay, and
the same stamp-versus-slot distinction applies to it. When a live leaf
activates (a cold reactivation after deactivation, a silo move, or a
crash), it rebuilds its in-memory projection by replaying the WAL
through its replay filter - from its snapshot's offset when a persisted
snapshot rehydrated it, otherwise across the whole readable window (see
State Model).
Unlike the snapshot leaf the live leaf has no pinned map - it serves the current point in time -
so it resolves per-mutation shard ownership by the key's virtual slot
under the current registry ShardMap, fetched once at the start of
replay.
This matters only on a cold reactivation that replays a WAL suffix past a checkpoint taken before a shadow-forward. In steady state the live leaf applies the forwarded mutation in real time and folds it into its next checkpoint, so a warm leaf never re-evaluates the record. The trigger is a target leaf that checkpointed before a post-split write was shadow-forwarded through the donor, then reactivated cold:
- A write routed to the donor for an already-moved slot is shadow-forwarded into the target shard's WAL but keeps the donor's source stamp. Resolving by slot keeps that forwarded record on the target's live leaf (the current map routes its slot here), so the leaf's last-writer-wins merge applies it on replay. A stamp-based filter would drop it - its stamp names the donor - and resurrect the pre-forward value, for example losing a post-split delete and reappearing the deleted key with its stale drained value.
- Genuine sibling-shard data multiplexed through a shared WAL partition
(partitions are keyed by key hash, not by shard) still resolves to
another shard under the current map and is dropped, and a donor leaf's
own orphan copies of moved slots are dropped because their slot now
routes to the target - the live read path already seals those orphans
via
MovedAwaySlots.
The current map is fetched best-effort and is trusted only when it
references the leaf's own shard, so a registry hiccup at activation or a
map drawn from a foreign physical shard space falls back to the legacy
stamped-ShardIndex axis - a leaf can never reject its own writes.
Slot-less legacy leaves (pre-split-feature persisted state with a null
ShardIndex) keep their unconditional V1 single-leaf-per-shard
semantics and apply on the shard axis regardless.
Trade-offs
- Order: Keys/Entries are streamed in strict lexicographic (or reverse) order end-to-end, even when splits commit mid-scan. Reconciled entries participate in the same k-way merge as live cursors, so the ordering guarantee is preserved.
- Memory: scans allocate a
HashSet<string>for dedup that grows with the number of distinct keys observed during the scan, plus a per-reconciliation buffer proportional to the number of keys in slots that actually moved during the scan (typically small). For very large trees, prefer the range-bounded overload ofScanKeysAsync/ScanEntriesAsyncto bound memory. - Latency: when no split has ever occurred, scans take the same fast path as before (one round-trip per shard). The reconciliation passes only run when a shard actually reports moved slots.
- System trees: a reserved system tree - the lattice registry tree among them - bypasses the reconciliation path. It never participates in adaptive splits and routes by the default shard map without reading one from the registry (for the registry tree that read would deadlock), so its internal scans fan out with no moved-slot tracking. The public count and scan surface refuses system trees outright.
Autonomic detection
Each tree's hot-shard monitor is armed when the tree activates and re-armed from the point, conditional, batched, predicated and CRDT write paths (deletes, atomic batches, bulk loads and merges do not re-arm it) until an arming attempt succeeds, and a keepalive reminder re-activates it after collection. The monitor arms its sampling timer before it registers that keepalive, so a reminder service that is still initialising cannot leave it claiming to run with nothing sampling; a keepalive registration deferred that way is retried after each sampling pass until it succeeds. On each tick (default every 30 s) it:
- Polls every physical shard's
GetHotnessAsync()in parallel. - Computes ops/sec =
(reads + writes) / window.TotalSeconds. - Counts the shard migrations in flight for this tree by asking every
physical shard whether it is splitting - the source of an adaptive split
reports that it is, and so does the donor of a shard consolidation
(fold). If that count is already
MaxConcurrentAutoSplitsor more, the pass triggers nothing, so a fold in flight takes up one of the tree's autonomic split slots. BecauseHotShardMonitorGrainis keyed per-tree, the cap is enforced independently per tree - in a multi-tree cluster each tree may have up toMaxConcurrentAutoSplitsconcurrent splits running simultaneously. - Selects the top-
(MaxConcurrentAutoSplits - inFlight)hottest shards whose rate reachesHotShardOpsPerSecondThreshold(default 200 ops/s), skipping any shard already splitting, on cooldown, or owning a single virtual slot. Three shape clauses can refuse a hot shard as well: no shard is admitted once the tree hasMaxPhysicalShardsPerTreephysical shards (default 256), or while its load is uniform - the hottest shard's rate belowHotShardMinSkewRatio(default 1.5) times the median shard rate, the signature of a bulk ingest that a split cannot relieve - and a candidate holding fewer thanHotShardMinShardEntrieslive entries (default 1024) is skipped. - Triggers
ITreeShardSplitGrain.SplitAsyncon each selected shard in parallel viaTask.WhenAlland starts a per-shard cooldown.
Each split runs in its own coordinator activation, keyed
{treeId}/{sourceShardIndex}, so independent splits of different
source shards within the same tree do not contend on a single
coordinator. Concurrent target-index allocation is made collision-free
by an atomic registry-side counter, and concurrent shard-map swaps
compose because each applies its moved-slot reassignment in a single
registry call that re-reads the live map before persisting. Both
atomicity guarantees rely on the singleton registry grain running each
mutating call to completion before the next (it is not reentrant).
A split is suppressed (whole pass skipped) or a candidate is skipped individually when:
| Suppression rule | Scope | Mechanism |
|---|---|---|
AutoSplitEnabled = false |
Whole pass | Returns early. |
Tree younger than AutoSplitMinTreeAge (since monitor activation, default 60 s) |
Whole pass | Returns early. |
| Resize / reshard / merge / snapshot in progress | Whole pass | ILattice.IsResize/Reshard/Merge/SnapshotCompleteAsync() returns false. |
| Any shard has a pending bulk graft | Whole pass | IShardRootGrain.HasPendingBulkOperationAsync() returns true. |
Shard migrations in flight (adaptive splits and fold donors) already at MaxConcurrentAutoSplits |
Whole pass | Count of shards that report they are splitting. |
Cluster-wide split ceiling reached (MaxClusterConcurrentAutoSplits set) |
Per candidate | No cluster headroom left in the admission gate; the candidate is deferred to a later tick. |
Tree already has MaxPhysicalShardsPerTree physical shards (default 256) |
Per candidate (every hot shard) | Counted on orleans.lattice.split.admission.deferred with reason=shard_ceiling. |
Uniform load: hottest shard's rate below HotShardMinSkewRatio (default 1.5) times the median shard rate |
Per candidate (every hot shard) | Counted on orleans.lattice.split.admission.deferred with reason=uniform_load. |
Shard holds fewer than HotShardMinShardEntries live entries (default 1024) |
Per shard | One count per candidate that cleared every cheaper clause; counted with reason=low_occupancy. |
| Shard already splitting | Per shard | Excluded from candidate set. |
| Per-shard cooldown active (default 2 min) | Per shard | In-memory cooldown timestamp. |
| Shard owns a single virtual slot | Per shard | Cannot be subdivided further. |
Cluster-wide split concurrency (opt-in)
MaxConcurrentAutoSplits is enforced per tree: because the split monitor runs once per tree, each tree counts only its own in-flight migrations. In a multi-tenant or many-tree cluster the summed drain I/O from many trees splitting at once can saturate the storage provider even though no single tree exceeds its own cap.
MaxClusterConcurrentAutoSplits (default null = disabled) opts in to a cluster-wide admission gate - a singleton (well-known integer key 0) - that admits a new autonomic split only while the shard migrations in flight on the trees that set the ceiling, fold donors included, leave headroom under it. The ceiling is enforced in addition to each tree's MaxConcurrentAutoSplits and can only ever lower the number of splits a tree triggers, never raise it. When the option is null the monitor never requests a slot, so nothing is ever denied and the number of splits a tree starts is byte-for-byte identical to running without the option.
A tree with the ceiling unset does still publish an observation-only footprint to that singleton while it has migrations in flight, because the same grain is the cluster's readable split-activity source (see Reading split activity). Those footprints are held in a separate list and never consume admission headroom - a tree that never opted into a shared budget must not be able to throttle one that did - and publication is edge-triggered, so a tree with nothing splitting issues no call at all.
Reading split activity
Split progress is also exposed as a metric (orleans.lattice.split.in_flight), but metrics are write-only in-process: nothing can read one back to decide something. ILatticeAdmin.GetSplitActivityAsync is the readable counterpart, returning a cluster-wide SplitActivityReport (InFlight, ReportingTrees, ObservedAt, AnyInFlight) assembled from the footprints above. Both count shard migrations rather than only splits: a shard donating its slots to a shard consolidation (fold) counts as one in flight, exactly as an adaptive split's source does. It costs a single call to the singleton and never fans out across trees or shards.
Its first consumer is the Orleans.Lattice.Scaling scale-in safety gate, which suppresses scale-in while any split or fold is in flight so a silo is never drained underneath one; it is equally useful to an operator tool or a deployment guard. The count comes from the monitors' sampling passes, so it is not exact. An unsuppressed pass refreshes it once per HotShardSampleInterval; a pass suppressed while the tree is resized, resharded, merged into or snapshotted, during its age grace period, or after autonomic splitting is switched off re-publishes the tree's last count unchanged, so migrations a reshard starts are not counted while it runs; and a tree whose monitor starts with autonomic splitting disabled publishes nothing. Footprints expire, so a silo lost mid-migration cannot pin the count above zero indefinitely.
Per-tree heartbeat footprints (self-healing)
Admission uses a per-tree heartbeat model rather than long-lived permits. On every unsuppressed sampling pass an enabled monitor reports its tree's authoritative in-flight migration count (split sources and fold donors, read from each shard's own record of whether it is splitting) and how many new splits it wants; the gate drops any tree footprint whose time-to-live has lapsed, sums the live in-flight counts of the other trees, and grants new slots only up to the remaining cluster headroom. It then records this tree's footprint (its in-flight count plus any grant) with a fresh expiry of HotShardSampleInterval * 3. Because the count is re-reported from ground truth each pass, there is no permit to leak: a silo that crashes mid-split simply stops refreshing its footprint, so the stale entry lapses at its expiry and the next pass reclaims that share of the ceiling. This self-healing property is what makes the aggregate ceiling safe to enable.
Per-group override
Per-tree options resolve through named IOptionsMonitor<LatticeOptions>.Get(treeName), so a low-traffic tree group can clamp its own MaxConcurrentAutoSplits down (e.g. to 1) while a high-traffic group keeps a higher per-tree cap - all bounded in aggregate by the single global MaxClusterConcurrentAutoSplits ceiling.
Operator questions and the metrics that answer them
| Question | Metric | How to read it |
|---|---|---|
| (a) Do I need to enable this? | orleans.lattice.split.in_flight (summed across tree) and orleans.lattice.split.candidates_suppressed |
Both emit even when the gate is disabled. A high steady-state cluster sum with chronically non-zero suppression across many trees means aggregate drain pressure the per-tree cap cannot see. |
| (b) What ceiling should I pick? | orleans.lattice.split.in_flight peak / quantiles |
Size the ceiling near the aggregate your storage provider absorbs without saturating (correlate with storage-latency panels), leaving headroom above typical peak so the gate only bites during pathological bursts. |
| (c) How is the enabled gate affecting my system? | orleans.lattice.split.admission.deferred (reason=cluster_cap) |
Flat-zero means the ceiling never binds (raise it or leave the gate off). Sustained non-zero with rising hot-shard latency means the ceiling is too low and is starving legitimate elasticity. |
Tunables (LatticeOptions)
| Option | Default | Description |
|---|---|---|
AutoSplitEnabled |
true |
Master switch for autonomic splits. When false, HotShardMonitorGrain will not trigger any splits. It does not gate an explicit ReshardAsync, which dispatches splits through the same coordinator to grow the shard count (see Online Reshard). |
HotShardOpsPerSecondThreshold |
200 |
Operations/second at or above which a shard is considered hot. Intentionally low so splits occur before throughput degrades. |
HotShardSampleInterval |
30 s |
How often the monitor polls hotness counters. |
HotShardSplitCooldown |
2 min |
Minimum interval between consecutive splits of the same physical shard. |
MaxConcurrentAutoSplits |
2 |
Maximum concurrent splits per tree. The monitor counts the shard migrations already in flight on the tree against it, so a shard consolidation (fold) in flight takes a slot too. Each split runs in its own per-shard coordinator activation; the cap bounds aggregate storage I/O. |
MaxClusterConcurrentAutoSplits |
null |
Optional cluster-wide ceiling on autonomic split admission across all trees: a new split starts only while the shard migrations in flight on the trees that set it, folds included, leave headroom under it. null disables the gate (per-tree caps only, zero cost); a positive value opts in to a singleton admission gate enforced in addition to each tree's MaxConcurrentAutoSplits. |
MaxConcurrentMigrations |
4 |
Maximum concurrent splits (or, for a shrink, shard consolidations) an online reshard (ReshardAsync) keeps in flight. Separate from MaxConcurrentAutoSplits, but not additive with it: the monitor starts no split while a reshard runs, a growing reshard counts the shard migrations already in flight on the tree - autonomic splits and healing folds included - against this cap, and a shrinking reshard waits for any healing fold to finish before it starts its own. See Online Reshard. |
SplitDrainBatchSize |
1024 |
Maximum number of moved-slot entries the drain accumulates in memory before flushing to the target shard. Caps coordinator allocation regardless of source shard size. |
BackgroundDrainLeavesPerPass |
64 |
Maximum source leaves one Drain-phase pass visits before persisting its key cursor and yielding to the next tick. 0 or less disables the bound. The authoritative drains inside Swap and Complete are never bounded. Shared with the online snapshot copy and the cross-tree merge drain. |
BackgroundDrainMaxDuration |
10 s |
Wall-clock net for one Drain-phase pass, for leaves that are individually slow. TimeSpan.Zero disables it and leaves the leaf count as the only bound. |
AutoSplitMinTreeAge |
60 s |
Minimum tree age before autonomic splits are allowed; absorbs startup bursts. |
HotShardMinSkewRatio |
1.5 |
Minimum ratio of the hottest shard's rate to the tree's median shard rate before any split is admitted, so a uniformly loaded tree (a bulk ingest) is never split. A value at or below 1.0 disables the clause. It is the upper edge of the split/heal dead band whose lower edge is HotShardConsolidationSkewRatio. |
HotShardMinShardEntries |
1024 |
Minimum live entries a shard must hold to be split. 0 disables the floor and its per-candidate count probe. |
MaxPhysicalShardsPerTree |
256 |
Ceiling on the physical shard count autonomic splits may reach. An explicit ReshardAsync is not gated by it. 0 or less for no ceiling. |
MaxScanRetries |
3 |
Maximum bounded retries that a scan (CountAsync, ScanKeysAsync, ScanEntriesAsync) performs when ShardMap.Version keeps moving mid-scan due to concurrent splits. Throws InvalidOperationException on exhaustion. Increase if scans run during very-high split churn. See Consistency. |
Automatic over-split healing, which folds shards back together once a tree's load is uniform, is tuned separately - see ShardHealingEnabled, HotShardConsolidationSkewRatio, and MaxConcurrentShardConsolidations in Configuration. While a fold runs, its donor counts as a migration in flight against the split caps above, and healing admits no new fold while any shard of the tree is splitting.
Convergence guarantees
- No data loss - every write committed to S is either drained,
shadow-mirrored, or both, and
MergeManyAsyncis idempotent under LWW. - No prepared-mutation loss - the retroactive sweep at
BeginShadowWritere-stamps every in-flight prepared mutation from S's leaves onto T's_pendingTxbuckets, so aSetManyAtomicAsyncsaga whose Prepare landed on S before the split commits and completes against T with no perceived interruption. Combined withLatticeOptions.TxDecisionRetention(default 60 s), a sweep that installs a pending bucket after the saga's terminal fan-out has already broadcast can still resolve the verdict via the registry tombstone window. See Atomic Writes - Phase 4 Complete. - No mixed-round batch across the swap - the coordinator's drain
copies a shadowing saga's pre-saga value into T with a migration
marker, so between the shard-map swap and the arrival of the saga's
backstop terminal on T a reader could otherwise see that key's old
value while every sibling key already showed the new one. The
coordinator therefore installs a per-key marker naming the shadowing
saga, and T's read gate resolves it against the registry: a saga the
registry reports as in-flight or aborted is safe (the pre-saga value is
the correct answer), and so is any saga whose backstop terminal has
already landed on T. Otherwise the read raises
StaleShardRoutingExceptionand the deadline-bounded retry loop re-fans once the backstop lands. A saga whose decision has aged out ofTxDecisionRetentionreads asIndeterminateand takes the same conservative arm as a committed one - serving the migrated pre-saga value there would be an affirmative claim that the saga did not commit, which is exactly what the registry has stopped vouching for. See Atomic Writes - After the retention window. - No duplicate authority - after the swap, only T is reachable for moved slots via the public API; orphan entries on S are unreachable and reclaimed when the tree is purged, or when a later shard consolidation retires S and releases its storage, unless a later consolidation folds T back into S, which drains T's entries onto S before lifting S's seal so the survivor's copy is authoritative again, and then releases T's storage.
- Geometric convergence on a single hot slot - if all heat is in one
virtual slot, successive autonomic splits subdivide S's slot set in
half each pass, isolating the hot slot in
O(log virtualSlotsPerShard)splits.
Scope
Shard splitting is an autonomic concern. ITreeShardSplitGrain is internal
infrastructure: once AddLatticeAuth has installed its trust-boundary call
filter, starting a split asserts that the call originated inside the
cluster, so an external client call to start one is rejected with
LatticeAuthorizationDeniedException (a cluster without that filter does
not enforce the assertion). There is no public
API to trigger or control an individual split. ILattice.ReshardAsync grows a
tree's shard count by dispatching splits through the same coordinator (see
Online Reshard); autonomic splitting is tuned
exclusively through the LatticeOptions listed above.