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.

Diagram of the three reasons to shard Redis. Memory: used_memory close to maxmemory. Command thread: one saturated core while the others sit idle. Network: about 960 of 1,000 Mbps used by 40 KB values at 3,000 reads a second. Memory sharding helps used_memory maxmemory The working set no longer fits with headroom. One command thread sharding helps commands run on one core instantaneous_ops_per_sec, usec_per_call I/O threads do not parallelise execution. Network the NIC, not Redis total_net_output_bytes ≈ 960 of 1,000 Mbps 40 KB values × 3,000 reads a second.
Read four numbers before adding a node. Each ceiling has a different fix, and a slow command is none of them.

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 shardingProxy (Envoy, twemproxy, Codis)Redis Cluster
Extra network hopNoneOne; budget a fraction of a millisecond and measure yoursNone
Topology lives inEvery client, in every languageOne placeThe cluster itself
ReshardingManual, usually with downtimeProxy-managed, variesOnline, slot by slot
FailoverYou build itProxy plus sentinelBuilt in
Multi-key operationsImpossible across shardsUsually unsupportedSame-slot only
Real costTwo client libraries in two languages must hash identically, foreverThe proxy is now a tier you capacity-plan and make highly availableClients 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.

Diagram of Redis Cluster's 16,384 hash slots split across three nodes. The keys order:{8891}:items and order:{8891}:total both hash to slot 7249 on node B, drawn as one block that can never be split; the untagged key order:8891:items hashes to slot 4940 on node A. 16,384 hash slots across three nodes 0–5460 5461–10922 10923–16383 Node A Node B Node C order:{8891}:items order:{8891}:total slot 7249 never split order:8891:items slot 4940 whole key hashed slot = CRC16(key or {tag}) mod 16384
The slot is the smallest unit you can move. A hash tag is a promise never to split that group.

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.

A row of identical crates on a conveyor. Most hold a single pebble; one holds a boulder far too large for it, its sides bowing, showing how tagging keys by tenant pins the largest tenant to one slot.
Tag keys by tenant and the largest tenant is welded to one slot, on one node, for good.

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.

Two bar charts over six shards. Keys per shard are level. Requests per shard are not: one shard serves 40 per cent of requests because of a single hot key, and the other five serve about 12 per cent each. Keys per shard evenly spread S1 S2 S3 S4 S5 S6 Requests per shard one hot key 12% 12% 12% 40% 12% 12% S1 S2 S3 S4 S5 S6
Sharding distributes keys. It does not distribute traffic. Worked example from the text: one key takes 40 per cent of reads.

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.

A narrow bridge between two platforms blocked by one oversized parcel while small parcels queue behind it, showing how migrating one very large key blocks every client on the source node.
MIGRATE blocks both ends while it runs. One oversized key stalls every client on the source node.

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.