本文目录导读:

我来提供几种批量刷新过期缓存数据的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'])
使用建议
- 选择合适的刷新策略:根据缓存类型(Redis、文件、数据库)选择对应方案
- 设置适当的执行频率:可以使用定时任务(cron、schedule)定期执行
- 添加错误处理和日志记录:确保刷新过程可追踪和调试
- 考虑并发控制:使用锁机制避免重复刷新
- 监控刷新性能:记录刷新时间和效率,优化策略
选择适合你应用场景的缓存刷新方案,并根据实际需求进行调整。