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.
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) % Nor consistent hash ring. Even spread for random keys; range scans across shards expensive. ShardPay hashesaccountIdfor 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
settledDatemonth. - 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-volumereads 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
- Transfer from account
acc-1001toacc-2042. - Router computes
pk = hash(min("acc-1001", "acc-2042"))→ shardledger-42. - Single-shard transaction: debit
acc-1001, creditacc-2042, insert audit row. - Commit on shard leader with W=3 quorum.
- If accounts were on different shards (legacy data), saga orchestrator runs two-phase with compensation — 10x latency, avoided by colocation rule.
// 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_totalstable is keyed bymerchantIdon 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 —
accountIdnever changes;merchantIdmigration 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).
Databases trackSee sharding strategies in the Databases track
Quick recall
Everything you need if you only revisit this box.
- Partition key determines shard placement — treat it as immutable.
- Hash for even spread; range for time-series archives; composite for tenant isolation.
- 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.