Partitioning Strategies and Shard Keys
LESSON
Partitioning Strategies and Shard Keys
By the end of this lesson, you will be able to...
choose a shard key by starting from an invariant, traffic shape, and request path;
compare range, hash, compound-key, and directory-based placement under named constraints;
identify the read models and signals required by a key that keeps a decisive write local.
Idea in one sentence: A shard key is a promise about which work one authority can decide locally, plus a bet about where traffic will gather.
Core Insight
Harbor Point must confirm the last available allocation in bucket B-173 exactly once. Two desks can request that confirmation at the same time. The product also wants quick desk histories and issuer-wide risk reports.
The tempting key is desk_id. Traders often open a desk page, so putting each desk on one shard seems useful. It works while a desk only changes its own records. It breaks when desk North and desk South both try to confirm the last allocation in B-173. If the two commands go to different desk shards, neither shard can decide the scarce bucket alone.
The key decision is not “Which field appears most often in a query?” It is:
Which facts must one authority decide together before it acknowledges this operation?
For this scenario, that unit is the allocation bucket. Harbor Point uses allocation_bucket_id as its authority key. A bucket may be derived from issuer, instrument, settlement window, and a capacity slice, but it has one stable identity once a hold is created.
This does not make every query local. It makes the no-double-confirm decision local. Desk history and issuer exposure become separate read problems, with explicit indexes or projections rather than accidental cross-shard writes.
Separate Authority, Partition, and Placement
The word “shard” often hides three different things. Keep them separate.
| Layer | Harbor Point example | What it answers |
|---|---|---|
| Authority key | allocation_bucket_id=B-173 |
Which state must be decided together? |
| Logical partition | hash(B-173) -> 173 |
Which stable slice of the keyspace owns that key? |
| Physical shard group | partition 173 currently maps to group B |
Which replica group serves that slice today? |
The numbers in this lesson are illustrative. A system might use 256, 4,096, or a dynamically split range layout. The important design property is that application code names a stable authority key, while a versioned partition map is free to place the resulting partition on different replica groups.
def route_confirmation(command):
authority_key = command.allocation_bucket_id
partition = stable_partition(authority_key)
return partition_map.owner_for(partition)
The function is a teaching model. Real systems may use ranges, hashes, lookup tables, or a combination. But the router must receive the key before the authoritative operation begins. A lookup performed after an arbitrary shard accepts the command is already too late for a one-owner decision.
This distinction also prepares the next lesson. Rebalancing can move a logical partition from one physical shard group to another without redefining the business fact that B-173 has one owner.
Start with the Strongest Write, Not the Nicest Read
Compare two routing paths for the same scarce bucket.
Candidate: desk_id
desk North confirms B-173 -> desk shard N
desk South confirms B-173 -> desk shard S
Result: one bucket decision crosses two authorities.
Candidate: allocation_bucket_id
desk North confirms B-173 -> bucket partition 173 -> shard B
desk South confirms B-173 -> bucket partition 173 -> shard B
Result: both commands reach the same authority before it chooses a winner.
The second route does not guarantee that the confirmation succeeds. It gives one shard the chance to enforce the right condition in its local transaction.
def confirm(command):
shard = route_confirmation(command)
return shard.transaction(
assert_hold(command.bucket_id, command.hold_id, status="held"),
assert_bucket_available(command.bucket_id),
mark_bucket_confirmed(command.bucket_id),
record_command(command.idempotency_key),
)
This pseudocode is a teaching model, not a database recipe. It assumes the shard transaction can make these checks atomic for one bucket. Its point is visible locality: the availability check, state change, and idempotency record share one authority boundary.
Now consider a desk page. It asks for every reservation handled by desk North. The authority-key layout does not make that query local. Harbor Point can maintain a desk-index projection:
authoritative bucket B-173 confirms a hold
-> publish a desk-history event
-> desk projection adds the reservation for North
-> desk page reads the projection with its stated freshness
That projection is a trade-off, not a defect. The confirmation endpoint reads the authoritative bucket. The desk page may read a delayed derived view if its product contract permits it. It must not reuse that view to decide whether the scarce bucket is available.
So far, the shard key has done one precise job: it has kept the decisive write local. The cost is that other access patterns need their own read path and freshness promise.
Score Candidate Keys with Workload Evidence
An invariant is necessary, but it is not enough. A key can keep the right decision local and still create one overloaded partition. Harbor Point scores each candidate with the same evidence.
| Question | Why it matters | Evidence to inspect |
|---|---|---|
| Does every decisive command carry the key? | A router cannot target an owner from missing information | Request shapes and command logs |
| Does the key contain the invariant? | Otherwise a local transaction becomes cross-shard coordination | State-transition diagram and invariants |
| Is there enough cardinality? | Too few possible values limit useful distribution | Distinct-key count and expected growth |
| Is traffic concentrated on a few values? | High cardinality does not prevent hot values | Top-key write share and per-partition latency |
| Is the key stable? | Changing it may mean moving ownership | Product lifecycle and update history |
| Which reads lose locality? | Every expensive read needs an honest serving plan | Query traces and latency budgets |
Cardinality means how many distinct key values are available. It helps a system spread work, but it does not prove that work is spread. A million possible issuer IDs still produce a hotspot if most new commands target one issuer. Likewise, a high-cardinality timestamp can create a hot range when every new write lands at the same end of the keyspace.
The measurements are more useful than a key-name debate. Observe the busiest logical partitions, the share of writes held by the hottest one, targeted versus scatter/gather reads, and cross-shard attempts for the decisive command. Then decide whether the key matches the traffic that exists rather than the traffic the team hopes will exist.
Compare Placement Strategies
The authority key and the placement strategy solve related but separate problems. Harbor Point can choose different placement styles for different data shapes.
| Strategy | Good fit | Main cost or boundary |
|---|---|---|
| Range by business value | Ordered scans and locality, such as an issuer region | Hot ranges, especially for monotonic or popular values |
| Hash of a stable key | Evenly spread independent operations | Loses range locality; range reads may fan out |
| Compound key | One field establishes ownership and another improves spread or query targeting | Every write and query must carry and interpret the same key order |
| Directory lookup | Explicit tenant or workload placement | Metadata becomes a versioned, monitored routing dependency |
Range placement is useful when nearby values are read together. It is risky when new values arrive in order, such as an increasing timestamp, because writes can collect at one edge. A hash of a stable key often spreads independent writes more evenly, but an issuer-wide scan may then need many partitions.
A compound key can make a useful compromise. For example, (issuer_id, allocation_bucket_id) can retain issuer context while distinguishing individual buckets. It is not automatically better: the first components, key order, and actual command path determine whether it helps placement or recreates an issuer hotspot.
A directory can explicitly map a tenant or bucket to an owner. This is attractive when a few customers require special placement or when a controller needs to move work independently of the data key. Its cost is operational: the mapping needs an epoch or version, cache behavior, observability, and stale-route rejection. In this track, the directory's live handoff belongs to the rebalancing lesson; here, the design question is whether the extra routing surface buys a needed placement choice.
These are situated preferences. A reporting-first product might select issuer ranges and accept a coordinated confirmation workflow. Harbor Point selects the bucket boundary because an incorrect confirmation is more costly than maintaining derived report views.
Failure Boundaries and Signals
Shard-key design cannot make every invariant local. A transfer between B-173 and B-881 still spans two owners. That operation needs a distributed transaction, an escrow scheme, a saga, compensation, or a redesigned business rule. The key merely tells the team where that boundary begins.
Changing a shard key is also not a harmless schema edit. It can require moving records, rebuilding indexes, updating request contracts, and validating that old and new routes do not both decide the same item. Treat a mutable field such as desk ownership, priority, or region as projection data unless the product is prepared to manage those ownership moves.
Harbor Point watches signals that expose the chosen key's boundary:
| Signal | What it reveals |
|---|---|
| Write rate and queue time by logical partition | A hot authority boundary, not just a busy cluster |
| Largest-key share of writes | Whether a few values defeat the expected distribution |
| Targeted versus scatter/gather read rate | Whether request shapes support the chosen key |
| Cross-shard confirmation attempt count | An invariant that escaped its intended owner |
| Projection lag and stale-read age | The user cost of read locality sacrificed for write correctness |
| Partition-map version misses | Routers or clients using stale placement data |
No single cluster-average throughput graph answers these questions. The design is healthy only if the invariant remains local, the hot partitions stay within their budget, and the derived reads state their freshness honestly.
Check Your Understanding
Check: A desk dashboard is the most frequent query, but the last-unit confirmation must never be decided by two shards. Which field should define the primary authority key?
Think first, then reveal.
Answer: The allocation bucket identifier. It routes competing confirmations to one owner before the decision. The dashboard can use a derived desk index; query frequency alone does not justify splitting the decisive invariant.
Check: A key has millions of distinct values. Does that prove it will avoid hotspots?
Think first, then reveal.
Answer: No. A few values may receive most of the traffic, and monotonic values can concentrate new writes in one range. Inspect frequency and per-partition load as well as cardinality.
Practice: Review a New Shard Key
A delivery service wants to shard package data by city_id because dispatch maps are city-local. However, the service must prevent two couriers from accepting the same package. The package can be offered to couriers in different cities near a border.
Choose an authority key and a serving plan. State:
- the invariant that must remain local;
- the key that routes the acceptance command;
- the city-facing read model and its allowed freshness;
- one hotspot risk and signal;
- whether a cross-city handoff belongs in the local transaction or an explicit workflow.
A strong answer identifies package_id or a stable offer/allocation identifier as the acceptance authority, not city_id alone. It may maintain city dispatch projections with a stated lag budget. It should watch hot-package or hot-zone concentration, not only total city traffic. A package that changes service city is an ownership change or workflow boundary; it should not silently create two authorities that can accept it.
Connections
- Sharding and Authority Boundaries establishes why an authoritative command needs one owner; this lesson chooses the key that makes that ownership concrete.
- Rebalancing Partitions Under Live Traffic uses the stable logical partitions introduced here to move placement without changing the authority key.
- Secondary Indexes Across Shards develops the read-side cost of this choice: convenient lookup paths can cross authority boundaries and need explicit freshness semantics.
Resources
- [DOC] MongoDB: Choose a Shard Key — Focus: Compare cardinality, frequency, monotonicity, and query routing before judging a candidate key.
- [DOC] Cloud Spanner: Schema Design Best Practices — Focus: Inspect how primary-key order, monotonic values, and logical hashing affect hotspots and locality.
- [DOC] Amazon DynamoDB: Partition Key Design — Focus: Relate uniform activity and access patterns to a concrete partition-key design review.
- [BOOK] Designing Data-Intensive Applications — Focus: Connect partitioning, skew, secondary indexes, and application-level routing to the chosen invariant.
Key Takeaways
- A shard key should first identify the smallest state one authority must decide together, then be tested against real traffic and query evidence.
- Cardinality, value frequency, and monotonicity are separate properties; none alone proves that load will be balanced.
- Range, hash, compound, and directory strategies trade locality, distribution, and routing complexity under different constraints.
- Derived views can make convenience reads useful without becoming the authority for a decisive write.
- Stable authority keys and versioned placement make later rebalancing possible, but do not remove cross-shard workflow boundaries.