Why this matters
- Horizontal scale fails if one shard absorbs 40% of traffic — adding nodes does not help skewed keys.
- Celebrity accounts, platform settlement wallets, and payroll batch windows are classic ShardPay hot spots.
- Rebalancing is an operational discipline with dual-write cutovers — not a one-time cluster setup task.
- Interviewers ask how you detect skew and migrate without downtime or data loss.
Keys hash onto the ring; walk clockwise to find owner
Detect and fix skew
ShardPay monitors per-shard QPS, p99 latency, and disk growth every minute. A merchant payroll day spiked shard ledger-42 to 18k TPS while fleet median was 3k — PagerDuty fired on shard_qps_skew_ratio > 3. Root cause: enterprise tenant's 50,000 employee accounts hashed across shards, but all debits targeted one platform settlement account on ledger-42.
Fix: split the settlement account into 16 salted sub-accounts (settle-{00..0f}) with scatter-gather aggregation on admin reads.
Skew vocabulary
- Skew — Uneven distribution of keys, traffic, or storage across partitions. ShardPay dashboards rank shards by
qps / fleet_avg_qps— anything above 2.0 enters investigation. - Hot key — Single key receiving disproportionate access. ShardPay's
platform:settlement:usdrow received 40% of write IOPS on one shard despite even account hashing. - Hot partition / shard — Entire shard overloaded because many hot keys or one dominant tenant landed there. Directory override moved a SaaS tenant to dedicated hardware after composite key isolation was insufficient.
- Split — Divide a partition or key into children (sub-keys, sub-shards). ShardPay splits hot logical shards by doubling vnode density on the ring segment and migrating half the key range to new hardware.
- Migrate — Copy key range from source to destination with dual-write and verification. ShardPay's migration controller runs copy → dual-write → verify checksum → cutover → retire source.
- Salting — Append random prefix to spread logical key across multiple physical keys (
acc-123#03). Read path merges — ShardPay uses salting only when write QPS exceeds single-row limits; reads pay scatter-gather cost.
Detection signals
| Signal | Threshold | Action |
|---|---|---|
| shard_qps / avg_qps | > 3.0 for 5 min | Investigate top keys |
| single_row_write_iops | > 5k/s | Salt or split key |
| disk_usage_growth | > 2x fleet avg | Reshard or archive |
| p99_latency per shard | > 2x SLO | Load shed or migrate |
shard_qps / avg_qps
Threshold> 3.0 for 5 minActionInvestigate top keyssingle_row_write_iops
Threshold> 5k/sActionSalt or split keydisk_usage_growth
Threshold> 2x fleet avgActionReshard or archivep99_latency per shard
Threshold> 2x SLOActionLoad shed or migrate
ShardPay hot-spot detection signals
Walkthrough: Zero-downtime shard migration
- Control plane marks
accountIdrange[hash_A, hash_B)for migration fromshard-3toshard-17. - Copy phase — Background job copies historical rows; verifies row counts and checksums.
- Dual-write phase — Router writes every new transfer to both
shard-3andshard-17for keys in range (7-day window). - Catch-up — Destination applies replication stream until lag = 0.
- Cutover — Ring version v44 points range to
shard-17only; dual-write ends. - Cleanup — Delete migrated rows from
shard-3after 30-day retention for rollback.
// ShardPay dual-write router during migration
public void writeTransfer(Transfer t, MigrationState mig) {
Shard primary = ring.locate(t.fromAccount());
primary.execute(t);
if (mig.isInDualWriteRange(t.fromAccount())) {
Shard destination = mig.destinationShard();
destination.execute(t); // idempotent by transferId
}
}
Mitigation strategies
Fix playbook
- Key splitting — Break one logical account into N physical sub-accounts. ShardPay's payroll settlement uses 16 sub-accounts; batch job round-robins credits to spread writes.
- Dedicated shard (directory override) — Control plane maps
tenantId=enterprise-99→shard-dedicated-5regardless of hash. Used when one tenant pays for isolated capacity. - Read offloading — Hot key writes may be unavoidable; serve reads from replicas. ShardPay's settlement balance admin UI reads from async replica — writes still bottleneck until split.
- Caching layer — Read-heavy hot keys fronted by Redis with TTL. ShardPay caches public merchant profile, not balances — financial reads stay on ledger.
- Rate limiting per key — Protect shard from abuse. ShardPay throttles single-account write rate at 1k TPS with
429and merchant notification. - Proactive rebalancing — ShardPay auto-triggers migration planning when forecast shows shard crossing 70% disk in 14 days — not after outage.
Example: Viral merchant Black Friday
Merchant mega-shop drove 8x normal QPS. Hash partitioning spread their buyer accounts evenly, but their merchantId rollup counter lived on one shard. ShardPay's CDC aggregator fell behind. Mitigation: moved rollup to a dedicated Kafka consumer partition keyed by merchantId, materialized view on separate store — decoupled hot counter from ledger shard write path.
Quick recall
Everything you need if you only revisit this box.
- Monitor per-partition QPS, latency, and disk — skew appears before hard outages.
- Hot keys need split, salt, or dedicated resources — hashing alone does not fix power-law access.
- Rebalance with copy → dual-write → verify → cutover; idempotency keys make dual-write safe.
Test yourself
Answer these before moving on — recall is what makes it stick.