skip to content

An optimistic check-and-set loop against Redis works fine in staging but in production one popular key sees so much concurrent traffic that most attempts are cancelled and the service latency rises. How do you diagnose and fix this?

level: principalimportance: should knowfreq 28%

answer

  1. attempts ≈ 1/(1−p); feedback makes it cliff-edge
  2. measure attempts/success + who writes the key
  3. atomic command > script > shard > coalesce
  4. jitter + bounded retries = guardrail, not fix
  5. hot key may be a modelling smell

basics

~20 s

Measure attempts per success and who writes the key. Then remove the contention rather than tune the loop: use a single atomic command if one exists, move the logic into one Lua script, shard the key so writers no longer collide, or batch updates. Keep bounded, jittered retries as a backstop only.

solid answer

~60 s

**Diagnose first.** Instrument attempts-per-success and the distribution of retries; if the mean is above one you are burning more work on re-reads than on the operation. Identify every writer of the watched key — including refresh jobs that rewrite identical values, which invalidate watchers just as hard as real updates. Confirm it is a single hot key with `redis-cli --hotkeys` or key-level metrics rather than a general capacity problem. **Then remove the contention, in this order:** 1. **Use an atomic command.** Many loops reimplement `INCRBY`, `SET … NX`, `ZADD … GT`, `HSETNX` or `SET … XX GET`. One command means no conflict at all. 2. **Move the decision server-side.** One Lua script or Function does read-decide-write with nothing interleaved: no retries, one round trip, flat cost under load — provided the script stays microsecond-short. 3. **Shard the key.** Split the counter or set across N sub-keys with a hash tag, write to one at random, aggregate on read. Contention drops by roughly N. 4. **Reduce write rate.** Coalesce updates in the client, move volatile fields out of the watched key. Retry caps and jittered backoff stop the meltdown; they do not fix it.

code

text · 5 lines
text
# write path: pick one of 16 shards
INCRBY counter:{page42}:7 1

# read path: aggregate (one node, thanks to the shared hash tag)
MGET counter:{page42}:0 counter:{page42}:1 ... counter:{page42}:15

go deeper

for a junior

Recognise the symptom (most attempts cancelled) and that retrying harder is not the fix; a single atomic command often replaces the loop.

for a middle

Explain the retry economics, then give concrete fixes: atomic commands, moving logic into a script, jittered bounded retries.

for a senior

Lead with instrumentation (attempts per success, writers of the key, hotkey detection) and pick the fix that removes contention, noting the script-duration and cluster constraints.

for a principal

Name the positive-feedback failure mode, order remedies by how much contention they eliminate, weigh exactness against contention in the data model, and add shedding and guardrails so the system degrades predictably.

## Why the loop looks fine until it does not Optimistic concurrency has a non-linear failure mode. With conflict probability *p*, expected attempts are about 1/(1−p). At p=0.05 you retry once in twenty operations and nobody notices. At p=0.6 you average 2.5 attempts, each costing a read and a transaction setup, and — crucially — the extra load you generate *raises* p further, because more clients are now sitting in the window between WATCH and EXEC. That positive feedback is why the graph looks flat and then goes vertical. Staging, with one tenth the traffic, sits on the flat part. ## Step 1 — measure the right things Three numbers make the diagnosis unambiguous: - **Attempts per successful operation**, as a distribution, not a mean. A p99 of 8 with a mean of 1.3 means a small set of keys is starving while everything else is fine. - **Write rate on the watched key**, and its sources. Watch invalidation counts *any* write, so a job rewriting an unchanged value at 500/s is indistinguishable from 500 real updates. Audit who writes the key before you redesign anything. - **Key skew.** `redis-cli --hotkeys` (with an LFU policy) or sampled key metrics tell you whether this is one hot key or a broad load problem. The fixes are completely different. Also check granularity: if you are watching a hash and the invalidating writes touch a field your logic never reads, the whole problem is an artefact of key layout. ## Step 2 — prefer removing the conflict over managing it **Atomic command.** The cheapest fix is discovering the operation already exists: counters (`INCRBY`), claim-once (`SET k v NX EX`), monotonic maxima (`ZADD key GT`), conditional set (`SET k v XX`), take-and-replace (`SET k v GET`, `GETDEL`). Zero conflicts, one round trip, no code. **Server-side script.** When the decision is real logic, one Lua script or Redis Function performs read, decision and write with nothing interleaved. Conflicts stop existing; clients simply serialise on the execution thread, which is bounded and fair, instead of racing and discarding work. The constraint is duration: while the script runs nothing else does, so it must be short and loop-free over unbounded collections. This converts wasted work into orderly queuing — usually the single biggest win available. **Shard the key.** For aggregate structures, split `counter` into `counter:{shard}:0..N-1`, have each writer pick a shard (by client id or at random) and have readers sum them. Contention per key falls by roughly N; the cost is a fan-out on read and a slightly more complex model. The same trick applies to sets and to rate limiters. In a cluster, decide deliberately whether the shards share a hash tag (co-located, aggregatable in one script) or spread across nodes (more throughput, aggregation done client-side). **Cut the write rate.** Coalesce: buffer updates in the client for 50 ms and write once. Move volatile fields out of the watched key so refreshers no longer invalidate readers. Question whether the write needs to happen synchronously at all. ## Step 3 — make the loop survivable regardless Even with the right design, keep the guardrails: - **Bounded attempts.** After N tries, fail fast and surface a specific error rather than looping. An unbounded loop under a traffic spike becomes a self-inflicted denial of service, consuming connections and CPU on both ends. - **Jittered backoff.** Without randomisation, all conflicting clients retry in lockstep and re-collide; jitter spreads them out and is often worth several points of success rate on its own. - **Load shedding.** When retries exceed a threshold, shed or queue the work; it is better to reject 2% of requests quickly than to melt the instance for everyone. - **Release watches on every path.** In pooled clients, an early return that skips EXEC/DISCARD/UNWATCH hands the next borrower stale watches and mysterious nil results. ## Step 4 — reconsider the model Sometimes the honest answer is that the hot key is a design smell: a global counter that every request touches, a single sorted set acting as a work queue, a shared config blob rewritten on every deploy. Ask whether the value must be exact and immediate. Approximate counters, per-partition aggregates, or a stream that a single consumer folds into the aggregate remove the contention entirely rather than distributing it. ## What a strong answer sounds like Lead with measurement, name the feedback loop that makes optimistic CAS fail suddenly, then propose fixes in order of how much contention they *remove* rather than manage — atomic command, server-side script, sharding, write-rate reduction — and finish with the guardrails and the question of whether the hot key should exist at all.

  • Why does jittered backoff help an optimistic retry loop more than a longer fixed delay?
    Conflicting clients that fail at the same instant also retry at the same instant under a fixed delay, so they collide again as a group — the collisions merely repeat more slowly. Randomising the delay spreads attempts across the window so only a subset overlap each round, which raises the success rate per attempt without adding latency for everyone.
  • When is sharding a hot key the wrong fix?
    When reads need an exact, immediate aggregate and are far more frequent than writes — you then move the cost onto every read, fanning out to N keys. It is also wrong when the operation is a genuine invariant rather than an accumulation, such as a unique claim, since the claim must be evaluated against one place. There a single atomic command or a script is the answer.
  • How would you decide between fixing the loop and redesigning the data model?
    By asking what the hot key represents. If it is an accumulation (counts, rates, sets of ids), it can almost always be sharded or approximated and the contention disappears. If it is a shared decision point that genuinely must serialise, then serialising it deliberately server-side with a script is honest and the model stays. The redesign is warranted when the exactness the hot key provides is not actually required by the product.

Everyone racing for the same doorway and backing off when they collide: adding polite retries does not help nearly as much as adding doors or letting a steward admit people one at a time.

saying these in an interview costs you the question

  • Fixing a retry storm by raising the retry limit
  • Adding a distributed lock so 'only one client retries at a time', converting waste into a queue with worse failure modes
  • Not auditing which writers invalidate the watched key, including no-op refresh jobs
  • Sharding a key whose semantics require a single evaluation point, such as a unique claim
  • Moving logic into a script without bounding its runtime, turning a retry problem into a latency problem for all clients
  • Treating a hot key as a capacity problem and scaling the instance instead

context