PrepZone Logo
PrepZone

Partitioning and Data Locality

Range, hash, and composite keys — how ShardPay shards accounts without cross-shard transactions.

Why this matters

  • Partition key choice is effectively permanent — resharding is expensive and risky.
  • Cross-shard operations need sagas, denormalization, or careful API design — avoid when possible.
  • Hot partitions break horizontal scale — even perfect hash partitioning fails on skewed keys.
  • Interviewers ask you to pick a partition strategy for a payment ledger and defend trade-offs.
Hash partitionaccountId % N
Range partitionA–M, N–Z
DirectoryLookup table
Hash by key for even spread; range for scans; directory for manual control.

Partition key design

ShardPay uses hash(accountId) mod 256 logical shards, mapped to physical nodes via consistent hashing. Range partitioning was rejected early: celebrity merchant accounts and platform settlement accounts would saturate range boundaries. Composite keys tenantId + accountId isolate SaaS merchants so one tenant's growth does not reshuffle others.

The partition key determines colocation — transfers between two accounts hit one shard only when both accounts hash to the same partition.

Strategies

  • Hash partitioning — shard = hash(key) % N or consistent hash ring. Even spread for random keys; range scans across shards expensive. ShardPay hashes accountId for even merchant distribution across 256 logical shards.
  • Range partitioning — Contiguous key ranges per shard (e.g., A–M, N–Z). Efficient range queries and time-series scans. ShardPay uses range partitioning only for immutable audit archives partitioned by settledDate month.
  • Composite / compound key — Multi-column partition key for tenant isolation. ShardPay's enterprise tier uses (tenantId, accountId) so a large tenant can be migrated to dedicated physical shards without rehashing the global keyspace.
  • Directory-based partitioning — Lookup table maps key → shard. Flexible migration at cost of directory availability. ShardPay's control plane maintains a shard directory in Raft — hot tenants get directory overrides pointing to dedicated hardware.
  • Cross-shard query — Scatter-gather across all shards or avoided by denormalization. ShardPay's GET /merchants/{id}/total-volume reads a pre-aggregated counter updated by CDC — no live cross-shard sum.

Colocation for single-shard transfers

ShardPay colocates both legs of a transfer on one shard using partitionKey = hash(min(accountA, accountB)). Debit and credit run in one serializable transaction — no distributed 2PC on the happy path.

Walkthrough: Transfer shard placement

  1. Transfer from account acc-1001 to acc-2042.
  2. Router computes pk = hash(min("acc-1001", "acc-2042")) → shard ledger-42.
  3. Single-shard transaction: debit acc-1001, credit acc-2042, insert audit row.
  4. Commit on shard leader with W=3 quorum.
  5. If accounts were on different shards (legacy data), saga orchestrator runs two-phase with compensation — 10x latency, avoided by colocation rule.
Java
// ShardPay partition router
public ShardId routeTransfer(String fromAccount, String toAccount) {
    String colocationKey = fromAccount.compareTo(toAccount) < 0
        ? fromAccount : toAccount;
    return shardRing.locate(colocationKey);
}

public ShardId routeAccount(String accountId) {
    return shardRing.locate(accountId);
}

Avoiding cross-shard pain

Design patterns

  • Denormalize for queries — Store merchant-level rollups on a separate keyed table. ShardPay's merchant_daily_totals table is keyed by merchantId on the merchant's home shard, updated by outbox CDC from transfer events.
  • Scatter-gather with limits — Admin search across shards uses parallel query with 2s timeout per shard, merge top-K. ShardPay uses this only for internal support tools — not merchant APIs.
  • Sticky secondary indices — Global index by email → accountId stored in a dedicated index service (Elasticsearch), not SQL scatter-gather. Ledger remains partitioned by accountId.
  • Resharding triggers — ShardPay reshard when physical node exceeds 70% disk or 15k TPS. Logical 256 shards map to growing physical cluster without changing application hash modulo.
  • Partition key immutability — accountId never changes; merchantId migration requires explicit rekey job. ShardPay forbids partition key updates in schema migrations without platform review.

Example: Celebrity merchant hot shard

A viral merchant onboarded 2M accounts in one week. Hash partitioning spread accounts evenly, but QPS concentrated on the merchant's settlement account — one hot row, not a hot shard. ShardPay split the settlement account into sub-accounts with salted prefixes (settle-00, settle-01, …) aggregated at read time — partition strategy alone did not solve row-level hot spots (see hot-spots article).

Quick recall

Everything you need if you only revisit this box.

  1. Partition key determines shard placement — treat it as immutable.
  2. Hash for even spread; range for time-series archives; composite for tenant isolation.
  3. Colocate transfer accounts on one shard via min(accountA, accountB) hash when possible.

Test yourself

Answer these before moving on — recall is what makes it stick.