本文目录导读:

在Python中规避并发缓存更新冲突,本质是要解决读写竞争和写写竞争问题,以下是几种常见且有效的策略,按推荐程度排序:
使用线程安全的缓存结构
functools.lru_cache + 锁
from functools import lru_cache
import threading
class ThreadSafeCache:
def __init__(self, maxsize=128):
self._cache = lru_cache(maxsize=maxsize)
self._lock = threading.Lock()
def __call__(self, func):
def wrapper(*args, **kwargs):
with self._lock:
return self._cache(func)(*args, **kwargs)
return wrapper
@ThreadSafeCache(maxsize=100)
def expensive_function(x):
# 耗时操作
return x * x
使用 cachetools 库的TTLCache
from cachetools import TTLCache, cached
import threading
cache = TTLCache(maxsize=100, ttl=300) # 5分钟过期
lock = threading.Lock()
@cached(cache, lock=lock)
def get_user_data(user_id):
# 从数据库获取数据
return {"id": user_id, "name": "Alice"}
分布式缓存(Redis)的原子操作
使用Redis的SETNX实现分布式锁
import redis
import time
import uuid
class RedisDistributedLock:
def __init__(self, redis_client, lock_name, expire_time=10):
self.redis = redis_client
self.lock_name = f"lock:{lock_name}"
self.expire_time = expire_time
self.lock_value = str(uuid.uuid4())
def acquire(self):
"""获取分布式锁"""
return self.redis.setnx(self.lock_name, self.lock_value) and \
self.redis.expire(self.lock_name, self.expire_time)
def release(self):
"""释放分布式锁(使用Lua脚本保证原子性)"""
lua_script = """
if redis.call('get', KEYS[1]) == ARGV[1] then
return redis.call('del', KEYS[1])
else
return 0
end
"""
self.redis.eval(lua_script, 1, self.lock_name, self.lock_value)
# 使用示例
redis_client = redis.Redis()
lock = RedisDistributedLock(redis_client, "cache_update")
def update_cache_with_lock(key, new_data):
if lock.acquire():
try:
# 安全地更新缓存
redis_client.setex(key, 3600, new_data)
finally:
lock.release()
Redis的原子操作(无需锁)
# 使用INCR进行原子计数
redis_client.incr("visitor_count")
# 使用HSETNX进行原子设置
redis_client.hsetnx("user:1", "name", "Alice")
# 使用SET的NX选项
redis_client.set("key", "value", nx=True, ex=3600)
使用读写锁(Read-Write Lock)
import threading
class ReadWriteLock:
def __init__(self):
self._read_ready = threading.Condition(threading.Lock())
self._readers = 0
def acquire_read(self):
"""获取读锁"""
with self._read_ready:
self._readers += 1
def release_read(self):
"""释放读锁"""
with self._read_ready:
self._readers -= 1
if self._readers == 0:
self._read_ready.notify_all()
def acquire_write(self):
"""获取写锁"""
self._read_ready.acquire()
while self._readers > 0:
self._read_ready.wait()
def release_write(self):
"""释放写锁"""
self._read_ready.release()
# 使用示例
cache = {}
rw_lock = ReadWriteLock()
def read_cache(key):
rw_lock.acquire_read()
try:
return cache.get(key)
finally:
rw_lock.release_read()
def write_cache(key, value):
rw_lock.acquire_write()
try:
cache[key] = value
finally:
rw_lock.release_write()
使用CAS(Compare-And-Swap)模式
import threading
class AtomicCache:
def __init__(self):
self._data = {}
self._lock = threading.Lock()
def compare_and_swap(self, key, expected_value, new_value):
"""原子地更新缓存,只有当前值为expected_value时才更新"""
with self._lock:
current = self._data.get(key)
if current == expected_value:
self._data[key] = new_value
return True
return False
def get_and_update(self, key, update_func):
"""原子地获取并更新"""
with self._lock:
current = self._data.get(key)
new_value = update_func(current)
self._data[key] = new_value
return new_value
# 使用示例
cache = AtomicCache()
def safe_update(user_id, increment=1):
def update_func(current):
if current is None:
return increment
return current + increment
return cache.get_and_update(user_id, update_func)
最佳实践:组合策略
import threading
import time
from functools import wraps
class CacheManager:
def __init__(self, redis_client, local_ttl=60, refresh_interval=30):
self.redis = redis_client
self.local_cache = {}
self.local_ttl = local_ttl
self.refresh_interval = refresh_interval
self._locks = {}
self._global_lock = threading.Lock()
def _get_lock(self, key):
"""为每个key创建独立的锁"""
with self._global_lock:
if key not in self._locks:
self._locks[key] = threading.Lock()
return self._locks[key]
def get_or_update(self, key, fetch_func, force_update=False):
"""
获取缓存,如果失效则更新
使用多层缓存:本地内存 -> Redis -> 数据库
"""
lock = self._get_lock(key)
# 1. 尝试从本地缓存获取
if not force_update:
local_data = self.local_cache.get(key)
if local_data and time.time() - local_data['timestamp'] < self.local_ttl:
return local_data['value']
# 2. 使用锁防止并发更新
with lock:
# 双重检查(防止等待锁期间其他线程已更新)
if not force_update:
local_data = self.local_cache.get(key)
if local_data and time.time() - local_data['timestamp'] < self.local_ttl:
return local_data['value']
# 3. 尝试从Redis获取
redis_data = self.redis.get(key)
if redis_data and not force_update:
self._update_local_cache(key, redis_data)
return redis_data
# 4. 从数据库获取并更新
new_data = fetch_func()
self.redis.setex(key, 3600, new_data)
self._update_local_cache(key, new_data)
return new_data
def _update_local_cache(self, key, value):
"""更新本地缓存"""
self.local_cache[key] = {
'value': value,
'timestamp': time.time()
}
def invalidate(self, key):
"""清除缓存"""
with self._get_lock(key):
self.local_cache.pop(key, None)
self.redis.delete(key)
# 使用示例
redis_client = redis.Redis()
cache_manager = CacheManager(redis_client)
def get_user_data(user_id):
"""带缓存的用户数据获取"""
return cache_manager.get_or_update(
f"user:{user_id}",
lambda: fetch_user_from_db(user_id)
)
# 并发安全地更新
def update_user_data(user_id, new_data):
cache_manager.invalidate(f"user:{user_id}")
update_user_in_db(user_id, new_data)
# 可选:主动更新缓存
cache_manager.get_or_update(
f"user:{user_id}",
lambda: fetch_user_from_db(user_id),
force_update=True
)
- 最小化锁的范围:只在必要时加锁,避免不必要的性能损失
- 使用原子操作:优先使用Redis的原子命令(SETNX, INCR等)
- 分布式场景使用分布式锁:避免单机锁在多进程/多服务场景失效
- 实现重入性:考虑锁的重入问题,避免死锁
- 设置超时:防止锁长时间占用导致系统崩溃
- 降级策略:缓存失效时提供降级方案,避免雪崩
选择哪种策略取决于你的具体场景:
- 单进程多线程:使用
threading.Lock或读写锁 - 多进程(如Gunicorn):使用文件锁或Memcached/Redis
- 分布式系统:必须是Redis/Zookeeper等分布式锁
- 高并发写入:优先考虑CAS模式或最终一致性方案