Almost every Redis sharding project I have reviewed made its worst decision before a second node existed. The team hit a wall, concluded they had outgrown one instance, and went to a cluster. Six months later they had the same wall, now spread over six machines, plus a rebalancing problem and a key space they could no longer change.
Sharding Redis is not hard. Sharding it in a way you can still operate two years later is, and almost all of that difficulty is front-loaded into two decisions most teams make casually: which ceiling they are actually hitting, and how their keys are named.
First: work out which ceiling you have hit
There are exactly three reasons to shard a Redis instance, and they have different fixes.
Memory. The working set no longer fits, or it fits with so little headroom that a background save fails. This one is real and sharding solves it.
CPU on the command loop. Redis executes commands on a single thread. The I/O threading added in Redis 6 parallelises reading from and writing to sockets, not command execution. So one core is your entire command budget, and when it saturates, sharding genuinely helps, because each shard gets its own core.
Network. A single instance pushing large values will saturate a 1 Gbps link long before it saturates that core. If your average value is 40 KB and you serve 3,000 reads a second, you are asking for roughly 960 Mbps of egress and the problem is the NIC, not Redis.
Before adding a node, read four things: used_memory against maxmemory, instantaneous_ops_per_sec, INFO commandstats for per-command usec_per_call, and total_net_output_bytes over a known interval. They will tell you which of the three you have, and often something more useful: that you have none of them. A single modern core handles on the order of 80,000 to 120,000 simple GET/SET operations per second at small value sizes, and several hundred thousand when clients pipeline. If you are struggling at 25,000 operations a second you do not have a capacity problem, and sharding will not fix it.
What sharding does not fix
It does not fix an O(N) command. A KEYS scan, an HGETALL over a 40,000-field hash, a SMEMBERS on a large set, an LRANGE 0 -1 — each of these occupies the single command thread for its whole duration while every other client waits. After sharding you have the same stall, on one of six nodes, for the unlucky sixth of your traffic that lands there. Fix the command.
It does not fix missing pipelining. A hundred sequential round trips at 0.4 ms each is 40 ms of wall-clock that a single MGET would have done in one. That latency is in your network, and no topology change touches it.
It does not fix a hot key. More on this below, because it is the failure that most often survives a sharding project intact.
And before any of it: if one instance serves both your cache and your durable working data (sessions, rate limiters, queues, dedupe sets), separate those first. They have opposite migration properties. A cache can be re-pointed at new topology and allowed to warm up cold; a data store cannot. Teams who shard the combined instance find that the easy half of their data was holding the hard half hostage.
Choosing the mechanism
Three approaches, and the trade-offs are not close once you write them down.
| Client-side sharding | Proxy (Envoy, twemproxy, Codis) | Redis Cluster | |
|---|---|---|---|
| Extra network hop | None | One; budget a fraction of a millisecond and measure yours | None |
| Topology lives in | Every client, in every language | One place | The cluster itself |
| Resharding | Manual, usually with downtime | Proxy-managed, varies | Online, slot by slot |
| Failover | You build it | Proxy plus sentinel | Built in |
| Multi-key operations | Impossible across shards | Usually unsupported | Same-slot only |
| Real cost | Two client libraries in two languages must hash identically, forever | The proxy is now a tier you capacity-plan and make highly available | Clients must be cluster-aware; key design becomes permanent |
Client-side sharding looks appealingly simple and is the one I argue against hardest in a polyglot shop. The moment a Python service and a Node service must agree on a hashing algorithm, a ring position count and a seed, you have an undocumented protocol holding your data together, and it breaks silently during a library upgrade.
Redis Cluster is the default answer for most teams, and the rest of this article assumes it. Everything here applies equally to Valkey, the fork that followed the 2024 licence change; the slot model and the migration mechanics are identical.
Key design is topology design
Redis Cluster hashes each key with CRC16 modulo 16,384 to get a slot, and assigns slot ranges to nodes. Operations touching more than one key must land in the same slot, or you get a CROSSSLOT error. Lua scripts and transactions have the same constraint.
The escape hatch is the hash tag: only the text between the first { and the first } after it is hashed, and only if that text is non-empty. No closing brace, or an empty {}, and the whole key is hashed instead — which is the quiet way a tagging scheme stops working. So order:{8891}:items and order:{8891}:total are guaranteed to share a slot and can be read in one MGET or written in one MULTI.
This is the single most consequential decision in the whole exercise, and here is the part that is rarely said out loud: a slot is the smallest unit you can ever move between nodes. A hash tag is therefore a permanent commitment that the tagged group will never be split across machines.
That is exactly what you want at the granularity of one order, one session, one user. It is a trap at the granularity of a tenant.
We run infrastructure across a large institutional estate, and tenant sizes there do not vary by a factor of two. They vary by a factor of a thousand: a school with 400 students sits in the same system as a district body with hundreds of thousands. Tag on tenant ID because it makes your pipelines tidy, and the largest tenant's entire key space is pinned to one slot on one node, permanently. No amount of rebalancing will ever relieve that node. The only fix is a key-space migration, which is an application change and a data migration, not a topology change.

The rule I use: tag at the smallest granularity your atomic operations genuinely require, and never on a group whose maximum size you do not control.
Sharding distributes keys, not traffic
Hash slots spread keys evenly. They do nothing about the fact that 40 percent of your reads hit one key: the feature-flag blob, the global config document, a leaderboard, the roster of your biggest customer. After sharding, one node does 40 percent of the work while the other five split the rest, and you have paid for five machines without touching the bottleneck.
Find them before they find you. redis-cli --hotkeys works but needs an LFU policy in maxmemory-policy before OBJECT FREQ returns anything useful. Do not switch a data-bearing instance to allkeys-lfu to get it — that makes every key evictable; volatile-lfu gives the same frequency data and only evicts keys you gave a TTL. MONITOR shows everything and costs real throughput, so treat it as a last resort. I prefer sampling client-side: one percent of requests into a per-key counter, rolled up every minute. It costs almost nothing and it is always on, which matters because hot keys arrive with a product launch, not with a capacity plan.
Then match the remedy to the shape of the heat.
Read-hot, slow-changing: config, flags, reference data. Cache it in the application process with a short TTL, or use Redis 6 client-side caching via CLIENT TRACKING, which pushes invalidations to the client. On a key read 10,000 times a second, a five-second local TTL turns 50,000 reads into one, and five seconds of staleness on a feature flag is usually acceptable. Write that down in the design rather than discovering it in an incident.
Read-hot, fast-changing: replicate the key under N suffixed names, leaderboard:{0} through leaderboard:{7}, and have readers pick one at random. The tags put them in eight different slots, which is not the same as eight different nodes — check with CLUSTER KEYSLOT and the slot map rather than assuming. You pay 8x on writes and accept a small divergence window. That trade is often right for something written once a second and read ten thousand times.
Write-hot: a counter, a rate limiter on a shared resource. Split it into N sub-counters on different slots, increment a random one, sum on read. The read gets more expensive and the write stops being a bottleneck.
Resharding without downtime
Redis Cluster migrates a slot by marking it MIGRATING on the source and IMPORTING on the destination, then moving keys with MIGRATE until CLUSTER COUNTKEYSINSLOT reaches zero, then reassigning ownership. During the window, the source answers for keys it still holds and returns ASK redirections for keys it has already handed over.
Three things decide whether this is invisible or an incident.
MIGRATE is synchronous and blocking on both ends. It serialises the value, ships it, and waits. A 500 MB hash blocks the source node — every client on it — for as long as that takes. This is the most common cause of "the rebalance caused an outage", and it is entirely preventable: find your large keys with redis-cli --bigkeys and split them before you reshard, never during. Read that output carefully — --bigkeys ranks by element count, not by bytes, so a 200-field hash of large strings can outweigh a 50,000-field hash of integers. --memkeys, or MEMORY USAGE on the candidates, gives you the number that actually predicts the stall.

Batch by bytes, not by key count. MIGRATE accepts multiple keys in one call, which cuts round trips substantially. But a batch of 100 keys is fine when keys are 2 KB and catastrophic when one of them is 200 MB. Cap the batch on estimated total bytes, in the low single-digit megabytes, rather than on count.
Throttle deliberately. A rebalance that completes in twenty minutes while producing p99 spikes is worse than one that takes four hours and nobody notices. Move a slot, pause, watch the latency percentiles, move the next. Run it during your trough, and afterwards verify cluster_state:ok with no slots left in MIGRATING or IMPORTING — a half-finished migration is a live cluster with a permanent redirection tax.
One client-side detail worth auditing before you start: MOVED means the slot has moved for good, so update your cached slot map; ASK means this one key has already gone ahead of the slot, so retry that single request on the target prefixed with ASKING and do not update the map. Client libraries that conflate the two will thrash their topology cache through the entire migration. Check yours handles it, because the symptom looks like a network problem.
Size the shard by recovery time, not by RAM
The question teams ask is "how much fits on this node". The better question is "how long does this node take to come back".
A replica performing a full synchronisation transfers the whole dataset. At 1 Gbps line rate, 25 GB is a shade over three minutes of pure transfer, before you add the fork, the copy-on-write pressure during it, and the writes accumulating meanwhile. Size repl-backlog-size against that window: if the backlog wraps before a briefly disconnected replica reconnects, it cannot partially resync and you pay for the whole transfer again. On a write-heavy instance, copy-on-write during a fork can push memory use towards twice the dataset size in the worst case. That is why the standard advice is to keep maxmemory near half of physical RAM, and why people who ignore it meet the OOM killer during a routine save.
So pick a per-shard dataset size from the failover time you can tolerate, then derive the shard count from it. In most systems I have worked on that lands between 10 and 25 GB per shard, far smaller than the machines people instinctively buy.
What I got wrong
I have signed off on a hash-tag scheme keyed on tenant identifier. It made the application code clean: every tenant's data co-located, pipelines simple, transactions working. It was also unfixable by the time the largest tenants had grown into it, because their slots could not be split and rebalancing had nothing to work with.
And I have started a rebalance without auditing key sizes first, on a cluster holding a handful of very large hashes. The migration stalled the source node long enough to time out live traffic. redis-cli --bigkeys takes minutes and would have caught it. It is now the first line of the runbook rather than the first line of the retrospective.
The short version
Diagnose which of the three ceilings you hit, and check it is not a command problem wearing a capacity problem's clothes. Split cache from durable data before you split anything else. Choose Redis Cluster unless you have a specific reason not to. Tag keys at the smallest granularity your atomic operations need, never on an unbounded group. Instrument hot keys before the launch, not after. Audit big keys before you reshard. Size shards by how long you can afford a resynchronisation.
Almost all of that is decided in the week before the second node exists. That week is the whole strategy.
Illustrations generated with AI.
