Sharding and Authority Boundaries

LESSON

Consistency and Replication

012 30 min intermediate

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_id selects the owner that decides holds, releases, and confirmations for B-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:

Connections

Resources

Key Takeaways

PREVIOUS Chain Replication and Ordered Failover NEXT Partitioning Strategies and Shard Keys