You are given a 32-core host with 256 GB of RAM to run Redis for a high-throughput workload. Given that one Redis instance executes commands on a single thread, how do you decide the deployment shape?
answer
- one core per instance -> partition the host
- shard count from measured ceiling + headroom
- small datasets: fork, resync, restart, reshard
- memory headroom for copy-on-write, never swap
- isolate noisy workloads; replicas off-host
basics
~20 sOne instance uses about one core, so a big host means many instances, not one. Shard into multiple processes (Redis Cluster or per-service instances), keep each instance's dataset small enough for fast fork, replication and failover, and leave cores and RAM headroom for io-threads, bio threads and copy-on-write during snapshots.
solid answer
~60 sThe core constraint is that command throughput per instance is bounded by one core, so a 32-core host is a host for **many Redis processes**, not one huge one. My decisions: - **Shard count** from measured per-instance ceiling (typically 10^5 ops/sec for O(1) commands, far less for big values) with headroom - commonly 8-16 instances, leaving cores for the OS, io-threads and fork children. - **Instance size** deliberately small (tens of GB, not 200). Big instances mean slow forks for `BGSAVE`/AOF rewrite, long full-resyncs, big copy-on-write spikes, and slow failover. Small shards make every operation recoverable. - **Memory headroom**: never allocate all 256 GB - a write-heavy fork can transiently double pages, so plan `maxmemory` well below physical RAM per instance. - **Isolation**: put noisy workloads (big values, Lua, `SORT`) on their own instances so their latency does not become everyone's latency. - **Failure domain**: replicas for a shard must not live on the same host, so a single host loss cannot take out both roles. Then validate: pipelining and connection pooling usually buy more than any threading knob.
go deeper
The key recall is that one instance uses roughly one core, so a big machine runs several instances.
Add shard sizing rationale - fork, resync and restart times scale with dataset size - and memory headroom for copy-on-write.
Give a concrete plan: measured per-instance ceiling, instance count with headroom, maxmemory budget, lazyfree and persistence scheduling, replica placement.
Treat it as a partitioning and isolation problem with explicit failure domains and latency budgets, and rank interventions - client pipelining and data modelling before any server tuning.
## Start from the constraint One Redis instance executes commands on a single thread. Adding cores does not raise its command throughput. So on a large host, the first question is not "how do I tune Redis?" but **"how many Redis processes do I run, and how big is each?"** Everything else follows. ## Sizing by cores Measure the per-instance ceiling for *your* command mix, not a synthetic `SET`/`GET` benchmark: value sizes, pipelining depth, and whether Lua is involved change it by an order of magnitude. A rough shape: simple O(1) commands with pipelining can push into the hundreds of thousands of ops/sec per instance; large values or scripts drop it dramatically. Then budget cores. Each instance needs roughly one core for its main thread, plus a share for bio threads, plus whatever `io-threads` you enable, plus the OS and network softirq handling, plus **a spare core for a forked child during `BGSAVE`/AOF rewrite**. On 32 cores, 8-16 instances is a typical shape; running 32 instances leaves nothing for those overheads and produces worse tail latency than 12. ## Sizing by memory - smaller than you think The temptation is to give one instance 200 GB. Resist it, because dataset size drives several operations that are not throughput: - **Fork cost.** `BGSAVE`, AOF rewrite and full replica sync fork the process. The fork itself must copy page tables, which scales with resident memory, and it happens on the main thread - a multi-hundred-millisecond stall on a very large instance. Then copy-on-write means every page written during the snapshot is duplicated, so a write-heavy 200 GB instance can transiently need far more RAM. - **Recovery time.** Loading an RDB or replaying an AOF at startup is linear in dataset size. A 200 GB instance is a long outage; a 20 GB shard is a short one. - **Replication.** A full resync ships the whole dataset. Small shards resync fast and independently. - **Resharding.** In Redis Cluster, slot migration moves keys; big instances and especially big *keys* make that slow. So: many modest instances. Leave a real margin between the sum of `maxmemory` values and physical RAM - allocator fragmentation plus copy-on-write make "use all the RAM" a reliable way to invite the OOM killer or swapping, and swapping a Redis instance is a latency catastrophe because the single thread stalls on page faults. ## Sharding mechanism Two shapes, and the choice is architectural: - **Redis Cluster** - one logical keyspace across shards, clients follow redirects, slots move for rebalancing. Right when a single dataset must grow past one instance. - **Independent instances per workload/service** - simplest, and it gives you *isolation*, which is often the more valuable property: the analytics job that runs expensive commands cannot hurt the session store. Often you want both: cluster the big dataset, isolate the noisy tenants. ## Latency isolation as a first-class decision Because execution is serialized, latency is a shared resource within an instance. The principal-level judgment is to decide **what is allowed to share an event loop**. Workloads with any of: large values, `SORT`/big range reads, Lua scripts, pub/sub fan-out with slow consumers, or huge keys with TTLs, should not share an instance with a latency-critical path. This is an admission-control decision, not a tuning one. ## Knobs that matter, in order 1. **Pipelining and connection pooling in clients.** Usually the single biggest throughput win, because it amortizes syscalls, and it costs nothing on the server. 2. **Data modelling.** Bounded collections, no unbounded keys, values in the hundreds of bytes rather than megabytes. 3. **Sharding.** The only real answer to CPU saturation. 4. **`io-threads`** - only when measurement shows the main thread pegged on socket I/O and cores are free. 5. **lazyfree settings and `UNLINK`** - to keep deletes and expiries off the loop. 6. **`hz`, `maxclients`, output buffer limits** - fine adjustments, rarely the fix. ## Failure domains A host with 16 instances is a large blast radius. Replicas for a shard must live on a different host (and ideally a different rack/AZ), and the placement policy must prevent a shard's primary and replica from co-locating - otherwise the host loss takes both. Restart storms matter too: 16 instances all loading RDBs at once will contend for disk and CPU, so stagger persistence schedules so several instances do not fork simultaneously. ## The summary answer "One core per instance" turns hardware sizing into a partitioning problem. Choose shard count from measured per-instance throughput with headroom for forks and I/O; choose shard size from how long you can afford a fork, a resync and a restart; leave memory headroom for copy-on-write; isolate noisy workloads onto their own loops; and spend your first effort on client-side pipelining and data modelling before any server threading knob.
- Why not simply run one instance with io-threads set to 30 on this host?Because io-threads parallelize socket reads/writes and protocol work, not command execution - the single main thread still runs every command, so the throughput ceiling barely moves for CPU-bound work. You would also be paying coordination overhead and starving the OS, bio threads and any fork child of cores.
- What operational limits push you toward smaller per-instance datasets?Fork latency and copy-on-write memory during BGSAVE or AOF rewrite, full-resync duration when a replica reconnects, restart/load time from RDB or AOF, and slot-migration time in Cluster. All are roughly linear in dataset size, so a 20 GB shard degrades gracefully where a 200 GB one turns every routine event into an incident.
- How do you decide whether a workload deserves its own instance rather than a shared one?Look at its command profile and blast radius: large values, unbounded collections, Lua scripts, big range scans, or slow pub/sub consumers all convert into shared latency because execution is serialized. If its worst command would violate another workload's latency SLO, it gets its own event loop.
saying these in an interview costs you the question
- Proposing one giant instance and expecting Redis to use all cores
- Allocating maxmemory to the full physical RAM with no copy-on-write headroom
- Treating io-threads as the scaling answer instead of sharding
- Placing a shard's primary and replica on the same host
- Ignoring that pipelining and data modelling usually beat every server-side knob