Sharding and Authority Boundaries
LESSON
Sharding and Authority Boundaries
By the end of this lesson, you will be able to...
Distinguish replication, which copies an authority, from sharding, which splits authority.
Design a strong write so its routing key identifies one owner before the decision begins.
Explain why a rebalance is an ownership transfer rather than a background data copy.
Idea in one sentence: Sharding scales decisive work by giving different key ranges different owners, so routing and rebalancing become part of the correctness contract.
Core Insight
At 09:30, desk North and desk South both try to confirm the final unit in allocation bucket B-173. The current service sends both commands to one replicated primary. It can choose one winner, but the queue behind that one authority is growing as more desks trade.
Harbor Point can replicate the bucket and keep its history ordered. But all allocation writes still arrive at one authority. That authority becomes the limit even if it has many read replicas.
The tempting solution is “add more databases.” That adds capacity only if different machines are allowed to decide different writes. Sharding makes that change. It divides the keyspace into authority domains, each with its own replicated owner.
For a scarce allocation bucket, the important question is not “which server has the row?” It is “which owner may decide whether this bucket changes from held to confirmed?” That owner must be unambiguous before the decision starts.
The Promise We Need to Keep
Two desks race to confirm the same allocation bucket, B-173. Harbor Point must promise that one authority decides the outcome and that a retry reaches the same authority while the bucket is moving between shards.
POST /allocations/confirm
bucket_id=B-173
hold_id=H-8821
idempotency_key=K-44
Promise:
exactly one current owner decides the bucket transition.
That promise tells the API what information it needs. A command that arrives with only customer_id and a vague issuer description may require a lookup, fan-out, or cross-shard transaction before it can find the deciding owner. A command with bucket_id can route directly to one authority boundary.
The Naive Design
Start with one logical reservation authority:
writes
|
v
[one primary]
/ | \
replica replica replica
Replication improves availability, recovery, and read capacity. It does not make the primary's authoritative write decision parallel. Every confirm still queues behind one owner.
Now divide buckets among several shard groups:
B-000..B-255 -> shard A replica group
B-256..B-511 -> shard B replica group
B-512..B-767 -> shard C replica group
B-768..B-999 -> shard D replica group
This improves write capacity only because B-173 and B-881 can now be decided independently. Replication still exists inside each shard group; it protects the local authority. Sharding decides where local authority begins and ends.
Why It Breaks Without an Authority Key
The initial model says a router can send a request to any healthy shard and let the database find the record. That works for a broadcast report when delay is acceptable. It fails for a strong command.
confirm(bucket_id=B-173)
-> partition map says shard A owns B-173
-> shard A runs one local transaction
confirm(customer_id=C-7, bucket unknown)
-> router may query several shards
-> two owners or a stale map can enter the decision path
The missing model is an authority boundary. A shard key is not merely a distribution function. It states which facts may be decided together without cross-shard coordination.
For B-173, the bucket row, current hold, capacity count, and transition audit record should live under the same owner if the invariant is “confirm this bucket at most once.” A local transaction can then enforce the invariant.
shard A owns B-173
assert bucket.status = held and hold_id = H-8821
set bucket.status = confirmed
append audit event confirmed
all in one shard-local decision
A Better Boundary
Plain meaning:
Put the state that must be decided together under one owner, then replicate that owner for durability and failover.
In this scenario:
bucket_idselects the owner that decides holds, releases, and confirmations forB-173.
Technical names:
The stable value that selects an authority domain is the shard key. A router's mapping from that value or partition to a current owner is the partition map.
Mature designs separate a logical partition from its physical placement. Harbor Point can derive a stable partition from the bucket key, then map it to a shard group:
bucket B-173 -> logical partition 173 -> owner shard A, epoch 61
bucket B-881 -> logical partition 881 -> owner shard D, epoch 19
The application need not know the server name. It needs a stable key and a routing layer that can reject a stale epoch. This lets the system move partition 173 later without changing the meaning of bucket_id.
A Worked Ownership Transfer
Shard A is overloaded, so the control plane moves logical partition 173 to shard E. This trace uses simplified states to expose the authority handoff.
start
map: partition 173 -> A, epoch 61
A is authoritative; E has no copy
1. copy
copy a consistent base of partition 173 from A to E
2. catch up
stream later changes from A to E while A remains the only writer
3. fence
briefly drain or reject new writes for partition 173
verify E has applied A's final authoritative position
4. switch
publish: partition 173 -> E, epoch 62
5. enforce
A rejects epoch-61 writes after the switch
E accepts epoch-62 writes and becomes the only authority
6. clean up
retain A's old copy until recovery rules say it is safe to remove
The critical observation is step 4 plus step 5. Copying bytes alone does not change authority. Without an epoch and stale-owner rejection, a router with an old map could send a confirmation to A while a new router sends another to E. Two healthy shard groups would then be deciding the same bucket.
So far: a rebalance is not finished when E has the data. It is finished when exactly one epoch-authorized owner can make new decisions.
The Trade-off
Sharding improves parallel decision capacity and can keep common strong operations local. The trade-off is that the schema, APIs, router, and migration controller all become part of one distributed correctness surface.
| Choice | Helps | Costs or boundary |
|---|---|---|
| Bucket-oriented key | Keeps one allocation invariant local | Desk and issuer reports may need derived views |
| Customer-oriented key | Keeps customer history local | A scarce bucket may require cross-shard coordination |
| Missing shard key in a query | Flexible query shape | Router may scatter/gather across many shards |
| Live rebalance | Relieves hot shards and supports growth | Requires epochs, catch-up, drains, and stale-route handling |
MongoDB's documentation makes the routing part concrete: queries containing the shard key can target a shard, while queries without it can broadcast to all shards. The product-specific conclusion is not “always use one key.” Choose the key that keeps the most costly invariant local, then deliberately build projections for the queries it makes expensive.
Operational Consequences
The router and partition map need the same care as a database write path. Useful signals include:
| Signal | What it reveals |
|---|---|
| Per-shard write rate and queue time | Hot authority domains and imbalance |
| Targeted versus scatter/gather query rate | Whether APIs carry useful routing keys |
| Stale-epoch rejection rate | Cached maps or clients still targeting an old owner |
| Migration catch-up lag and drain time | Whether a move can safely switch authority |
| Cross-shard transaction or saga rate | Invariants that do not fit the chosen boundary |
Sharding does not make a multi-bucket workflow atomic. A swap between B-173 and B-881 crosses two authorities. The system must use a distributed transaction, escrow, saga, compensation, or redesign its invariant. The right choice depends on whether partial progress is acceptable and what external effects must never be duplicated.
Common Confusions
Confusion: Adding replicas is the same as sharding.
Why it is tempting: Both add machines.
Better model: Replicas copy one authority decision. Shards create several authority decisions for disjoint key ranges.
Confusion: A partition map is only a performance cache.
Why it is tempting: It looks like routing metadata.
Better model: For strong writes, the map identifies who may decide. A stale map must be rejected or refreshed, not silently followed.
Confusion: Rebalancing is complete after the data copy.
Why it is tempting: The future owner can read the copied records.
Better model: A live move needs a handoff of authority: catch-up, epoch switch, stale-route rejection, and a safe cleanup rule.
Check Your Understanding
Check: A confirm request needs one decision for B-173. Which request shape better protects the invariant: confirm(bucket_id=B-173, hold_id=H-8821) or confirm(customer_id=C-7, issuer=CA-MUNI)?
Think first, then reveal.
Answer: The first shape. It exposes the key that selects the authority before the decision begins, so the router can target one shard. The second may need lookup or fan-out and risks turning a local invariant into cross-shard coordination.
Check: E has copied all data from A, but clients can still send epoch-61 writes to A after the new map names E. What is missing?
Answer: Stale-owner fencing. A must reject the old epoch and E must accept the new one, so copied data becomes a single-authority migration rather than two live writers.
Practice: Review a Shard Move
The system moves inventory partition 173 from shard A to shard E. At the same time, a client retries confirm(bucket_id=B-173, idempotency_key=K-44) after a timeout. Describe the routing and write contract that prevents duplicate confirmation.
A good answer should mention:
- A partition map that resolves partition 173 to one owner and an epoch.
- Idempotent command identity so a retry recovers the recorded decision rather than creating another one.
- Catch-up and a fencing boundary before E becomes authoritative.
- Rejection or refresh of stale epoch requests at A after the switch.
- Signals for stale-route rejections, migration lag, and duplicate or conflicting confirmation attempts.
Connections
- Chain Replication and Ordered Failover keeps one replica group's history ordered; sharding decides how many such authority domains exist.
- Partitioning Strategies and Shard Keys compares candidate keys and the workload arguments hidden in each one.
Resources
- [DOC] MongoDB Sharding — Focus: Compare targeted operations, scatter/gather queries, shard keys, chunks, routers, and range migration with the authority model.
- [DOC] Vitess Documentation — Focus: Explore the sharding, routing, and migration material as examples of a production control plane.
- [BOOK] Designing Data-Intensive Applications — Focus: Connect partitioning, skew, replication, and query routing to the invariants your application must preserve.
Key Takeaways
- Replication copies one authority domain; sharding splits authority by key or partition.
- Strong commands should carry the key that routes them to the owner able to decide the invariant.
- A partition map and its epochs are correctness data during live ownership changes, not merely performance metadata.
- Rebalancing succeeds only when one new owner is authoritative and stale routes cannot keep writing to the old owner.