Python脚本如何批量刷新过期缓存数据

wen python案例 30

本文目录导读:

Python脚本如何批量刷新过期缓存数据

  1. Redis缓存批量刷新
  2. 本地文件缓存批量刷新
  3. 数据库查询缓存批量刷新
  4. 通用的缓存刷新管理器
  5. 使用建议

我来提供几种批量刷新过期缓存数据的Python脚本方案:

Redis缓存批量刷新

使用SCAN命令清理过期缓存

import redis
import time
from typing import List, Optional
class RedisCacheRefresh:
    def __init__(self, host='localhost', port=6379, db=0, password=None):
        self.client = redis.Redis(
            host=host,
            port=port,
            db=db,
            password=password,
            decode_responses=True
        )
    def flush_expired_keys(self, pattern: str = '*', batch_size: int = 100) -> int:
        """
        批量刷新过期缓存
        :param pattern: 匹配的键模式
        :param batch_size: 每批处理的键数量
        :return: 删除的键数量
        """
        cursor = 0
        deleted_count = 0
        while True:
            # 使用SCAN命令分批获取键
            cursor, keys = self.client.scan(
                cursor=cursor,
                match=pattern,
                count=batch_size
            )
            if keys:
                # 检查每个键是否过期
                for key in keys:
                    ttl = self.client.ttl(key)
                    if ttl <= 0:  # 已过期或不存在
                        self.client.delete(key)
                        deleted_count += 1
                        print(f"删除过期缓存: {key}")
            # 游标为0表示遍历完成
            if cursor == 0:
                break
        return deleted_count
    def refresh_with_ttl_check(self, pattern: str = '*', min_ttl: int = 3600) -> List[str]:
        """
        刷新TTL小于指定值的缓存
        :param pattern: 匹配模式
        :param min_ttl: 最小TTL阈值(秒)
        :return: 刷新的键列表
        """
        refreshed_keys = []
        cursor = 0
        while True:
            cursor, keys = self.client.scan(cursor=cursor, match=pattern, count=100)
            for key in keys:
                ttl = self.client.ttl(key)
                if 0 < ttl < min_ttl:  # 存在但TTL小于阈值
                    # 重新设置缓存(这里需要根据实际业务逻辑实现)
                    self._refresh_cache(key)
                    refreshed_keys.append(key)
                    print(f"刷新缓存: {key}, 原TTL: {ttl}s")
            if cursor == 0:
                break
        return refreshed_keys
    def _refresh_cache(self, key: str):
        """根据键重新生成缓存数据(需要根据业务逻辑实现)"""
        # 示例:重新设置TTL为1小时
        # 实际应用中应该重新计算数据并设置
        self.client.expire(key, 3600)
# 使用示例
cache_refresher = RedisCacheRefresh(host='localhost', port=6379)
deleted = cache_refresher.flush_expired_keys(pattern='user:*')
print(f"共删除 {deleted} 个过期缓存")

本地文件缓存批量刷新

import os
import json
import time
from pathlib import Path
from typing import Dict, Any
class FileCacheRefresh:
    def __init__(self, cache_dir: str = './cache'):
        self.cache_dir = Path(cache_dir)
        self.cache_dir.mkdir(exist_ok=True)
    def clean_expired_files(self, max_age_hours: int = 24) -> int:
        """
        清理过期缓存文件
        :param max_age_hours: 文件最大存活时间(小时)
        :return: 删除的文件数量
        """
        current_time = time.time()
        deleted_count = 0
        for file_path in self.cache_dir.iterdir():
            if file_path.is_file():
                # 检查文件最后修改时间
                file_age = current_time - file_path.stat().st_mtime
                max_age_seconds = max_age_hours * 3600
                if file_age > max_age_seconds:
                    file_path.unlink()
                    deleted_count += 1
                    print(f"删除过期缓存文件: {file_path}")
        return deleted_count
    def refresh_cache_files(self, pattern: str = '*.json', 
                           refresh_func=None) -> int:
        """
        批量刷新缓存文件
        :param pattern: 文件匹配模式
        :param refresh_func: 刷新函数,接收文件路径参数
        :return: 刷新的文件数量
        """
        refreshed_count = 0
        for file_path in self.cache_dir.glob(pattern):
            try:
                # 读取当前缓存数据
                with open(file_path, 'r', encoding='utf-8') as f:
                    data = json.load(f)
                # 检查是否过期
                if self._is_expired(data):
                    # 执行刷新策略
                    if refresh_func:
                        new_data = refresh_func(file_path)
                    else:
                        new_data = self._default_refresh(file_path)
                    # 保存更新后的数据
                    with open(file_path, 'w', encoding='utf-8') as f:
                        json.dump(new_data, f, ensure_ascii=False, indent=2)
                    refreshed_count += 1
                    print(f"刷新缓存文件: {file_path}")
            except Exception as e:
                print(f"刷新文件 {file_path} 失败: {e}")
        return refreshed_count
    def _is_expired(self, data: Dict[str, Any]) -> bool:
        """检查缓存数据是否过期"""
        # 假设数据中包含timestamp字段
        if 'timestamp' in data:
            current_time = time.time()
            return current_time - data['timestamp'] > 3600  # 1小时过期
        return True
    def _default_refresh(self, file_path: Path) -> Dict:
        """默认刷新策略"""
        data = {
            'timestamp': time.time(),
            'data': f'refreshed at {time.strftime("%Y-%m-%d %H:%M:%S")}'
        }
        return data
# 使用示例
file_cache = FileCacheRefresh(cache_dir='./data_cache')
deleted_files = file_cache.clean_expired_files(max_age_hours=48)
print(f"清理了 {deleted_files} 个过期缓存文件")

数据库查询缓存批量刷新

import sqlite3
import time
from datetime import datetime, timedelta
from typing import List, Tuple, Optional
class DatabaseCacheRefresh:
    def __init__(self, db_path: str = 'cache.db'):
        self.conn = sqlite3.connect(db_path)
        self.cursor = self.conn.cursor()
        self._init_cache_table()
    def _init_cache_table(self):
        """初始化缓存表"""
        self.cursor.execute('''
            CREATE TABLE IF NOT EXISTS cache_data (
                key TEXT PRIMARY KEY,
                value BLOB,
                created_at TIMESTAMP,
                expires_at TIMESTAMP,
                last_accessed TIMESTAMP
            )
        ''')
        self.conn.commit()
    def clean_expired_cache(self) -> int:
        """
        清理所有过期缓存
        :return: 清理的缓存数量
        """
        current_time = datetime.now()
        self.cursor.execute(
            'DELETE FROM cache_data WHERE expires_at < ?',
            (current_time,)
        )
        deleted_count = self.cursor.rowcount
        self.conn.commit()
        return deleted_count
    def refresh_stale_cache(self, stale_hours: int = 24) -> List[str]:
        """
        刷新陈旧的缓存(长时间未访问)
        :param stale_hours: 缓存被认为陈旧的小时数
        :return: 刷新的缓存键列表
        """
        stale_time = datetime.now() - timedelta(hours=stale_hours)
        # 查找陈旧缓存
        self.cursor.execute('''
            SELECT key, value FROM cache_data 
            WHERE last_accessed < ? OR last_accessed IS NULL
        ''', (stale_time,))
        stale_records = self.cursor.fetchall()
        refreshed_keys = []
        for key, old_value in stale_records:
            # 执行刷新逻辑
            new_value = self._refresh_cache_value(key, old_value)
            # 更新缓存
            current_time = datetime.now()
            new_expiry = current_time + timedelta(hours=48)  # 新的过期时间
            self.cursor.execute('''
                UPDATE cache_data 
                SET value = ?, created_at = ?, expires_at = ?, last_accessed = ?
                WHERE key = ?
            ''', (new_value, current_time, new_expiry, current_time, key))
            refreshed_keys.append(key)
            print(f"刷新缓存: {key}")
        self.conn.commit()
        return refreshed_keys
    def _refresh_cache_value(self, key: str, old_value: bytes) -> bytes:
        """
        根据业务逻辑刷新缓存值
        这里需要根据实际应用场景实现
        """
        # 示例:模拟新的缓存数据
        import json
        try:
            old_data = json.loads(old_value.decode('utf-8'))
            # 更新数据逻辑
            old_data['refreshed_at'] = datetime.now().isoformat()
            return json.dumps(old_data).encode('utf-8')
        except:
            return old_value
    def batch_refresh_by_pattern(self, pattern: str = '%') -> int:
        """
        按模式批量刷新缓存
        :param pattern: SQL LIKE模式
        :return: 刷新的缓存数量
        """
        self.cursor.execute('''
            SELECT key FROM cache_data WHERE key LIKE ?
        ''', (pattern,))
        keys = [row[0] for row in self.cursor.fetchall()]
        refreshed_count = 0
        for key in keys:
            # 执行全量刷新
            current_time = datetime.now()
            new_value = f'refreshed_{key}_{current_time.timestamp()}'.encode()
            new_expiry = current_time + timedelta(hours=48)
            self.cursor.execute('''
                UPDATE cache_data 
                SET value = ?, created_at = ?, expires_at = ?, last_accessed = ?
                WHERE key = ?
            ''', (new_value, current_time, new_expiry, current_time, key))
            refreshed_count += 1
        self.conn.commit()
        return refreshed_count
    def close(self):
        """关闭数据库连接"""
        self.conn.close()
# 使用示例
db_cache = DatabaseCacheRefresh('my_cache.db')
deleted = db_cache.clean_expired_cache()
print(f"清理了 {deleted} 个过期缓存")
refreshed = db_cache.refresh_stale_cache(stale_hours=12)
print(f"刷新了 {len(refreshed)} 个陈旧缓存")
db_cache.close()

通用的缓存刷新管理器

import asyncio
import logging
from typing import Callable, Dict, Any
from datetime import datetime
class CacheRefreshManager:
    """通用缓存刷新管理器"""
    def __init__(self):
        self.refresh_strategies: Dict[str, Callable] = {}
        self.logger = logging.getLogger(__name__)
    def register_strategy(self, name: str, strategy: Callable):
        """注册刷新策略"""
        self.refresh_strategies[name] = strategy
    async def refresh_cache(self, strategy_name: str, **kwargs) -> bool:
        """执行单个缓存刷新"""
        if strategy_name not in self.refresh_strategies:
            raise ValueError(f"未知策略: {strategy_name}")
        try:
            start_time = datetime.now()
            result = await self.refresh_strategies[strategy_name](**kwargs)
            elapsed = (datetime.now() - start_time).total_seconds()
            self.logger.info(f"策略 {strategy_name} 刷新完成,耗时: {elapsed:.2f}s")
            return result
        except Exception as e:
            self.logger.error(f"策略 {strategy_name} 刷新失败: {e}")
            return False
    async def refresh_all(self, **kwargs) -> Dict[str, bool]:
        """刷新所有注册的缓存策略"""
        results = {}
        tasks = []
        for strategy_name in self.refresh_strategies:
            tasks.append(
                asyncio.create_task(
                    self.refresh_cache(strategy_name, **kwargs)
                )
            )
        completed = await asyncio.gather(*tasks, return_exceptions=True)
        for i, strategy_name in enumerate(self.refresh_strategies.keys()):
            results[strategy_name] = not isinstance(completed[i], Exception)
        return results
    def run_batch_refresh(self, strategies: list = None, **kwargs):
        """运行批量的缓存刷新任务"""
        if strategies:
            selected_strategies = {
                name: self.refresh_strategies[name] 
                for name in strategies 
                if name in self.refresh_strategies
            }
        else:
            selected_strategies = self.refresh_strategies
        results = {}
        for name, strategy in selected_strategies.items():
            try:
                result = strategy(**kwargs)
                results[name] = result
                print(f"策略 {name} 刷新成功: {result}")
            except Exception as e:
                results[name] = None
                print(f"策略 {name} 刷新失败: {e}")
        return results
# 自定义刷新策略示例
def refresh_user_cache(**kwargs):
    """刷新用户缓存策略"""
    # 实现具体的用户缓存刷新逻辑
    print("正在刷新用户缓存...")
    return True
def refresh_product_cache(**kwargs):
    """刷新商品缓存策略"""
    # 实现具体的商品缓存刷新逻辑
    print("正在刷新商品缓存...")
    return True
# 使用示例
manager = CacheRefreshManager()
manager.register_strategy('user_cache', refresh_user_cache)
manager.register_strategy('product_cache', refresh_product_cache)
# 批量刷新所有缓存
results = manager.run_batch_refresh()
print("批量刷新结果:", results)
# 选择特定缓存刷新
results = manager.run_batch_refresh(strategies=['user_cache'])

使用建议

  1. 选择合适的刷新策略:根据缓存类型(Redis、文件、数据库)选择对应方案
  2. 设置适当的执行频率:可以使用定时任务(cron、schedule)定期执行
  3. 添加错误处理和日志记录:确保刷新过程可追踪和调试
  4. 考虑并发控制:使用锁机制避免重复刷新
  5. 监控刷新性能:记录刷新时间和效率,优化策略

选择适合你应用场景的缓存刷新方案,并根据实际需求进行调整。

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