PrepZone Logo
PrepZone

Hot Spots, Skew, and Rebalancing

Celebrity accounts and power-law traffic — detect skew and rebalance without downtime.

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.
Hash ring
Node Avn1, vn2
Node Bvn3, vn4
Node Cvn5, vn6

Keys hash onto the ring; walk clockwise to find owner

Virtual nodes spread keys evenly. Adding a node moves only adjacent key ranges.

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:usd row 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

SignalThresholdAction
shard_qps / avg_qps> 3.0 for 5 minInvestigate top keys
single_row_write_iops> 5k/sSalt or split key
disk_usage_growth> 2x fleet avgReshard or archive
p99_latency per shard> 2x SLOLoad shed or migrate
  • shard_qps / avg_qps

    Threshold> 3.0 for 5 min
    ActionInvestigate top keys
  • single_row_write_iops

    Threshold> 5k/s
    ActionSalt or split key
  • disk_usage_growth

    Threshold> 2x fleet avg
    ActionReshard or archive
  • p99_latency per shard

    Threshold> 2x SLO
    ActionLoad shed or migrate

ShardPay hot-spot detection signals

Walkthrough: Zero-downtime shard migration

  1. Control plane marks accountId range [hash_A, hash_B) for migration from shard-3 to shard-17.
  2. Copy phase — Background job copies historical rows; verifies row counts and checksums.
  3. Dual-write phase — Router writes every new transfer to both shard-3 and shard-17 for keys in range (7-day window).
  4. Catch-up — Destination applies replication stream until lag = 0.
  5. Cutover — Ring version v44 points range to shard-17 only; dual-write ends.
  6. Cleanup — Delete migrated rows from shard-3 after 30-day retention for rollback.
Java
// 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-5 regardless 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 429 and 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.

  1. Monitor per-partition QPS, latency, and disk — skew appears before hard outages.
  2. Hot keys need split, salt, or dedicated resources — hashing alone does not fix power-law access.
  3. 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.