You must double the shard count of a heavily loaded Redis Cluster that backs a latency-sensitive service, with no maintenance window. How do you plan and pace the slot migration, and what would make you refuse to do it online at all?
answer
- MIGRATE blocks both masters — big keys are the hazard
- stage add-node + replicas early, zero slots = zero traffic
- small chunks, modest --cluster-pipeline, watch p99 between
- cluster-node-timeout > worst stall, or self-inflicted failover
- skew is a key-design problem; a slot cannot be split
basics
~20 sSurvey key sizes first, add empty masters plus replicas ahead of time, then move slots in small paced chunks with a modest MIGRATE batch size, watching p99 and SLOWLOG between chunks. Refuse online if huge keys would block masters past cluster-node-timeout and trigger spurious failovers.
solid answer
~60 s**Plan before touching slots.** Inventory the keyspace: `--bigkeys`/`MEMORY USAGE` sampling to find multi-hundred-megabyte values, `CLUSTER COUNTKEYSINSLOT` to find fat slots, and hot-key data from `MONITOR` sampling or `redis-cli --hotkeys`. A single enormous key is the main hazard, because `MIGRATE` blocks both the source and target event loops while transferring it. **Stage the change.** Join the new masters and their replicas first — zero-slot masters take no traffic, so this step is reversible and can happen days early. Confirm every client actually refreshes topology, since a fleet that does not converge turns a reshard into a permanent latency tax. **Pace the migration.** Move slots in chunks of a few hundred, with a modest `--cluster-pipeline`, pausing between chunks to check p99, `SLOWLOG` and error rates. Ensure `cluster-node-timeout` comfortably exceeds the worst MIGRATE stall, or a blocked master gets declared failed and fails over mid-reshard. **Refuse online** when big keys cannot be split first, when clients are known-stale, or when the error budget cannot absorb ASK/TRYAGAIN retries — then fix key design or rebuild into a new cluster and cut over instead.
code
text · 16 lines# 1. find the hazards (run against a replica)
redis-cli -h replica --memkeys
for s in $(seq 0 16383); do echo "$s $(redis-cli cluster countkeysinslot $s)"; done | sort -k2 -nr | head
# 2. headroom check: stall tolerance
redis-cli config get cluster-node-timeout
# 3. paced chunks with checks between
for i in 1 2 3 4 5 6 7 8; do
redis-cli --cluster reshard 10.0.0.1:6379 \
--cluster-from <src-ids> --cluster-to <dst-id> \
--cluster-slots 256 --cluster-pipeline 10 --cluster-yes
redis-cli -h 10.0.0.1 slowlog get 10
redis-cli --cluster check 10.0.0.1:6379 | grep -i 'migrating\|importing\|covered'
sleep 60
donego deeper
Know that resharding happens online but costs latency, and that it is done in chunks with redis-cli tooling.
Explain the MIGRATE blocking cost, the value of adding nodes ahead of time, and the need for clients to refresh topology.
Own the runbook: survey, pacing knobs, node-timeout headroom, monitoring between chunks, and the repair path for open slots.
Decide whether to reshard at all — diagnose capacity versus skew, weigh in-place migration against build-and-cut-over, and set the risk and error-budget terms for the change.
## Framing: resharding is a live availability event Moving slots is not a background maintenance task. Every slot moved runs `MIGRATE`, which is synchronous and blocking on **both** the source and the target: the source serialises the value, ships it, waits for the target's `RESTORE`, then deletes locally. Both event loops are stalled for that duration, and Redis's single-threaded command execution means everything else on those two masters waits. Meanwhile clients absorb `ASK` redirects (an extra round trip), occasional `-TRYAGAIN` on split multi-key commands, and a burst of `MOVED` as each slot commits. So the question is not "can it be done online" — it can — but "how much of the latency budget does it consume, and is that bounded". ## Step 1 — survey before you plan - **Big keys.** `redis-cli --memkeys` / `--bigkeys` (sampling, run against replicas to spare masters) and `MEMORY USAGE <key>` on suspects. Anything in the hundreds of megabytes is a hard blocker: transferring it stalls two masters for seconds. - **Slot skew.** `CLUSTER COUNTKEYSINSLOT` over the slot range shows whether keys are evenly spread. Aggressive hash tagging concentrates whole tenants into one slot, and a slot cannot be split — a single overloaded slot is a design defect that resharding cannot fix. - **Hot keys.** `redis-cli --hotkeys` (needs an LFU policy) or sampled `MONITOR`. A hot key migrating means its clients take ASK redirects at peak rate for the duration. - **Client fleet.** Which libraries, which versions, and is periodic/adaptive topology refresh actually enabled? A service whose client never refreshes will pay double round trips indefinitely after the reshard; that is the most common post-reshard regression. ## Step 2 — stage what is free Joining new masters (`--cluster add-node`) and attaching their replicas costs nothing in traffic terms — a zero-slot master serves no keys. Do this early and separately so that the risky change window contains only slot movement. Verify `cluster_state:ok` and that the new nodes appear in every client's refreshed view before proceeding. ## Step 3 — pace the migration The knobs: - **Chunk size.** Instead of one `--cluster-slots 8192` run, do repeated small runs (a few hundred slots) so you can abort between chunks. Each chunk is independently committed; there is no all-or-nothing rollback anyway. - **`--cluster-pipeline`.** Keys per `MIGRATE` batch (default 10). Higher is faster but each blocking call is longer. With large values, lower it. - **`--cluster-timeout`.** The MIGRATE timeout; must exceed the transfer time of your largest key or the migration errors out mid-slot and leaves an open slot. - **`cluster-node-timeout`.** The failure detector's threshold. If a MIGRATE stall exceeds it, peers declare the master failed and promote its replica **in the middle of your reshard** — a self-inflicted failover. Verify headroom, and raise it temporarily if necessary. - **Time of day.** Trough traffic reduces both the ASK-redirect volume and the impact of each stall. Between chunks, look at: p99/p999 from the application, `INFO commandstats` and `SLOWLOG GET`, `latency` events, client error counters for MOVED/ASK/TRYAGAIN, and replica lag on both masters involved. ## Step 4 — know your abort and repair paths An interrupted reshard leaves slots flagged MIGRATING/IMPORTING, which keeps producing ASK and TRYAGAIN indefinitely. `redis-cli --cluster check` detects them; `--cluster fix` completes or reverts; `CLUSTER SETSLOT <slot> STABLE` is the manual clear. Rehearse this before the real run, because discovering it under pressure is how a resize becomes an incident. ## When to refuse the online path Refuse, and choose a different strategy, when: 1. **Unsplittable giant keys exist.** Fix the data model first — shard the giant hash under a hash tag, or move the blob out of Redis. Migrating it online will stall masters for seconds regardless of pacing. 2. **Clients cannot be trusted to converge.** If part of the fleet uses a client that does not refresh topology, or is not cluster-aware at all behind a proxy that caches badly, the reshard degrades them permanently. Fix clients first. 3. **The error budget cannot absorb retries.** TRYAGAIN and redirect-cap errors will occur; if the caller has no retry path and the operation is user-facing, that is real user-visible failure. 4. **The cluster is already at its limits.** Memory pressure with eviction running, or CPU saturation on the masters, means there is no headroom to absorb the extra work; resolve capacity first. The alternative in those cases is **build-and-cut-over**: stand up a correctly sized new cluster, populate it (dual-write from the application, or replicate and promote per shard), verify, then flip clients. It costs more resources and coordination but concentrates all risk into a single, reversible cut. For a pure cache with a tolerable miss storm, a third option is to accept a cold new cluster and let it refill — cheap, but only if the origin can survive the miss rate. ## The strategic question underneath Doubling shards is usually a response to a symptom. Before spending the risk, ask whether the pressure is memory (maybe TTL/eviction policy or value encoding is the real fix), CPU on one shard (a hot key or a hash tag concentrating a tenant — resharding will not help, since a slot cannot be split), or connection count (client pooling). Resharding fixes aggregate capacity, not skew; skew is fixed in the key design.
- One shard is hot because a single tenant's keys all share a hash tag. Will doubling the shard count help?No. A hash tag forces those keys into one slot, and a slot is the atomic unit of placement — it cannot be split across masters no matter how many shards exist. The load stays on whichever single master owns that slot. The fix is in the key design: drop or narrow the hash tag so the tenant's keys spread across many slots, accepting that multi-key operations across them are no longer possible, or move that tenant to a dedicated cluster.
- How could a resharding trigger an unwanted failover, and how do you prevent it?MIGRATE blocks the master's event loop, so it stops answering cluster bus pings. If a stall exceeds cluster-node-timeout, enough peers mark it PFAIL then FAIL and its replica is promoted, in the middle of the migration. Prevent it by keeping the worst-case MIGRATE duration well under cluster-node-timeout: split or exclude giant keys, lower --cluster-pipeline, and raise cluster-node-timeout temporarily if the margin is thin.
- When would you build a new cluster and cut over instead of resharding in place?When the online path's risk cannot be bounded: unsplittable multi-gigabyte keys, a client fleet that will not converge on the new topology, or no error budget for the TRYAGAIN and redirect retries the migration produces. A parallel cluster lets you populate and verify with zero impact on the live one and makes the change a single reversible flip, at the cost of temporary double capacity and a dual-write or replication path. For a pure cache with a tolerant origin, simply cutting over to a cold cluster and accepting the miss storm can be the cheapest option of all.
saying these in an interview costs you the question
- Treating resharding as a background operation with no latency impact
- Ignoring large keys before starting, then blaming the cluster for stalls
- Running the whole reshard in one unpaced command at peak traffic
- Assuming more shards fixes a hot key or a hash-tag-concentrated tenant
- Not knowing that an aborted reshard leaves slots open and keeps emitting ASK/TRYAGAIN