feat(cache-manager): add stampede protection to get_or_set - #345
ShivanshShukla wants to merge 1 commit into
Conversation
allen0099
left a comment
There was a problem hiding this comment.
Thanks @ShivanshShukla, this is a solid first cut. The explicit lock switch with per-call None inheriting the manager setting, the contextvar re-entrancy guard, the winner's re-check and store-before-release all match what we agreed on #66. The full suite passes against live Redis and Memcached. A few things need to change before merge, one of them a real gap.
1. A crashed winner still leads to a stampede (blocking)
With the default wait_timeout (= lock_ttl), every waiter times out at the moment the dead winner's lock expires. The timeout check runs before the waiter tries to take the lock, so they all take the fallback and run factory. The takeover path described in the PR is never reached in that case.
Repro (10 concurrent misses, lock_ttl=2; the crashed winner is simulated by acquiring the lock and never releasing it):
import asyncio, time
from fastapi_cachex.backends import MemoryBackend
from fastapi_cachex.manager import CacheManager
from fastapi_cachex.lock import CacheLock
async def run(factory_secs, crashed):
backend = MemoryBackend()
m = CacheManager(backend, lock=True, lock_ttl=2)
calls = 0
async def factory():
nonlocal calls
calls += 1
await asyncio.sleep(factory_secs)
return "v"
if crashed:
await CacheLock(m._cache_key("k"), ttl=2, backend=backend).acquire(blocking=False)
t = time.monotonic()
await asyncio.gather(*(m.get_or_set("k", factory) for _ in range(10)))
print(f"factory={factory_secs}s crashed={crashed}: ran {calls}x in {time.monotonic()-t:.1f}s")
await backend.aclose()
async def main():
await run(0.5, False) # ran 1x
await run(0.5, True) # ran 10x
await run(1.5, True) # ran 10x
asyncio.run(main())The existing crashed-winner tests use a single waiter, so they can't see this.
A suggestion (happy to hear alternatives): when wait_timeout is not given, waiters have no deadline of their own. They keep waiting while someone holds the lock, since each holder is bounded by lock_ttl, and take over as soon as it is free. An explicit wait_timeout is then the caller's own latency budget, after which it raises or computes, as documented. Please add a regression test with several waiters and a crashed winner that asserts factory runs once.
2. One way to inject the clock
There are two mechanisms: the private _sleep/_monotonic constructor parameters and the module-level _sleep/_monotonic. Please keep only the module-level hooks (the clock fixture in conftest.py already patches _monotonic; a sleep hook can be patched the same way), so the public __init__ signature has no underscore parameters. Most of the new tests still sleep in real time (wait_timeout=0.08, etc.); with the hooks they can run on the fake clock (#185).
3. Failure cases on the live backends
test_stampede_protection_across_backends covers only the happy path. As agreed on #66 (and listed as gates in #280), please cover a crashed winner, a raising factory and a wait timeout on memory, Redis and Memcached.
4. Smaller points
lockis not validated.CacheManager(lock="no")is accepted and, being truthy, turns locking on. The docstring promisesTypeError, so please reject non-bools in__init__andget_or_set.- If
lock_instance.release()raises (a backend error), it replaces the computed result, or the factory's own exception. Consider logging a failed release rather than propagating it, since the lease expires on its own anyway. - The timeout raises the builtin
TimeoutError, while the package already hasLockTimeoutError(aCacheXError). What do you think about raising that, or aCacheXErrorsubclass that also subclassesTimeoutError, soexcept CacheXErrorcatches it? - Docs:
$\pm 10\%$won't render on the docs site; please write±10%. The APP_CACHE bullet "preventing concurrent misses from running factory simultaneously" should mention the timeout fallback. It is also worth saying that the lock key islock:<prefix><key>(lock:cache:user:42by default). - Please add a changelog fragment,
changelog.d/66.added.md(seechangelog.d/README.md: bold one-line summary first, no leading-, no issue link), and rebase on the currentmaster.
I've updated the PR description: the links pointed at local file:/// paths, the polling numbers now match the code (50 ms base, 1.5x, 500 ms cap), and it has Closes #66.
6d59d9b to
5bd7a73
Compare
|
Thanks for the thorough review and catching the crashed-winner stampede! All feedback has been addressed and rebased on top of the latest 1. Crashed Winner Stampede Gap
2. Clock Injection Cleanup
3. Failure Cases Across Live BackendsParameterised the failure cases across
4. Robustness & Details
|
|
Thanks for the update — the crashed-winner stampede is fixed. I re-ran the concurrency repro from my earlier review (10 concurrent misses, including a crashed winner with a slow factory) and the factory now runs exactly once in every scenario. The full suite also passes against live Redis and Memcached. I also mutation-checked the new tests. Removing the winner's cache re-check and releasing the lock before storing are both caught. A few small things remain before this is ready:
Once these are in, I think this is good to go. |
Introduce opt-in lock-based stampede protection in CacheManager.get_or_set() with double-check, polling backoff, re-entrancy prevention, crashed winner recovery, and comprehensive multi-backend tests.
5bd7a73 to
d3310db
Compare
Summary
Addresses cache stampede (dog-piling) when concurrent misses occur for the same key in
CacheManager.get_or_set().When an expensive factory (e.g., slow database query or rate-limited upstream API) is protected by$N$ simultaneous factory executions. This PR introduces opt-in distributed/local lock-based stampede protection built on
get_or_set(), a key expiration previously resulted inCacheLock.Key Changes
lock,lock_ttl,wait_timeout, andraise_on_timeoutparameters toCacheManager.get_or_set().lock: bool = Falseandlock_ttl: int = 60toCacheManager.__init__()for opting in globally while allowing per-call overrides.factoryto handle race conditions where another worker populated the cache just before lock acquisition.factorythemselves) uponwait_timeout, or raisingTimeoutErrorwhenraise_on_timeout=True.contextvars.ContextVartracking acquired lock keys to safely bypass locking if a factory recursively callsget_or_set()on the same key within the same task/context.docs/APP_CACHE.mdandi18n/zh-TW/docs/APP_CACHE.md.Verification
15/15 passed).pytest tests/test_cache_manager.py).Closes #66