Partitioning Strategies and Shard Keys

LESSON

Consistency and Replication

013 30 min advanced

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:

  1. the invariant that must remain local;
  2. the key that routes the acceptance command;
  3. the city-facing read model and its allowed freshness;
  4. one hotspot risk and signal;
  5. 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

Resources

Key Takeaways

PREVIOUS Sharding and Authority Boundaries NEXT Rebalancing Partitions Under Live Traffic