Python脚本如何规避并发缓存更新冲突

wen python案例 36

本文目录导读:

Python脚本如何规避并发缓存更新冲突

  1. 使用线程安全的缓存结构
  2. 分布式缓存(Redis)的原子操作
  3. 使用读写锁(Read-Write Lock)
  4. 使用CAS(Compare-And-Swap)模式
  5. 最佳实践:组合策略

在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
    )
  1. 最小化锁的范围:只在必要时加锁,避免不必要的性能损失
  2. 使用原子操作:优先使用Redis的原子命令(SETNX, INCR等)
  3. 分布式场景使用分布式锁:避免单机锁在多进程/多服务场景失效
  4. 实现重入性:考虑锁的重入问题,避免死锁
  5. 设置超时:防止锁长时间占用导致系统崩溃
  6. 降级策略:缓存失效时提供降级方案,避免雪崩

选择哪种策略取决于你的具体场景:

  • 单进程多线程:使用 threading.Lock 或读写锁
  • 多进程(如Gunicorn):使用文件锁或Memcached/Redis
  • 分布式系统:必须是Redis/Zookeeper等分布式锁
  • 高并发写入:优先考虑CAS模式或最终一致性方案

抱歉,评论功能暂时关闭!