本文目录导读:

我来介绍几种在Python中定时同步缓存数据到持久化存储的方法:
使用 threading.Timer(简单定时任务)
import threading
import time
import json
from typing import Any, Dict
class CacheSync:
def __init__(self, sync_interval: int = 60):
"""
:param sync_interval: 同步间隔(秒)
"""
self.cache: Dict[str, Any] = {}
self.sync_interval = sync_interval
self.running = False
def start_sync(self):
"""启动定时同步"""
self.running = True
self._schedule_sync()
def _schedule_sync(self):
"""安排下一次同步"""
if self.running:
timer = threading.Timer(self.sync_interval, self._sync_to_disk)
timer.daemon = True
timer.start()
def _sync_to_disk(self):
"""同步缓存到磁盘"""
try:
with open('cache_backup.json', 'w', encoding='utf-8') as f:
json.dump(self.cache, f, ensure_ascii=False, indent=2)
print(f"[{time.strftime('%Y-%m-%d %H:%M:%S')}] 缓存同步完成")
except Exception as e:
print(f"同步失败: {e}")
finally:
self._schedule_sync() # 安排下一次同步
def stop_sync(self):
"""停止同步"""
self.running = False
# 使用示例
cache = CacheSync(sync_interval=30)
cache.start_sync()
cache.cache["user_1"] = {"name": "张三", "score": 95}
try:
time.sleep(120) # 运行2分钟
finally:
cache.stop_sync()
使用 schedule 库(更灵活的定时任务)
import schedule
import time
import pickle
from typing import Any, Dict
from datetime import datetime
class ScheduledCacheSync:
def __init__(self, backup_file: str = "cache_backup.pkl"):
self.cache: Dict[str, Any] = {}
self.backup_file = backup_file
self._load_from_disk()
def _load_from_disk(self):
"""启动时从磁盘加载缓存"""
try:
with open(self.backup_file, 'rb') as f:
self.cache = pickle.load(f)
print(f"已加载缓存: {len(self.cache)} 个条目")
except FileNotFoundError:
print("未找到缓存文件,使用空缓存")
self.cache = {}
def sync_to_disk(self):
"""同步到磁盘(可被 schedule 调用)"""
try:
with open(self.backup_file, 'wb') as f:
pickle.dump(self.cache, f)
timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
print(f"[{timestamp}] 缓存同步完成 - {len(self.cache)} 个条目")
except Exception as e:
print(f"同步失败: {e}")
def start_scheduled_sync(self, interval_minutes: int = 5):
"""启动定时同步任务"""
schedule.every(interval_minutes).minutes.do(self.sync_to_disk)
# 也可以添加其他时间规则
# schedule.every().hour.do(self.sync_to_disk)
# schedule.every().day.at("00:00").do(self.sync_to_disk)
print(f"定时同步已启动,间隔: {interval_minutes} 分钟")
while True:
schedule.run_pending()
time.sleep(1)
# 使用示例
sync = ScheduledCacheSync()
sync.cache["temperature"] = 25.6
sync.cache["humidity"] = 60
# 启动定时同步(在单独的线程中运行)
import threading
sync_thread = threading.Thread(
target=sync.start_scheduled_sync,
args=(2,), # 每2分钟同步一次
daemon=True
)
sync_thread.start()
# 主程序继续运行
time.sleep(10)
sync.cache["pressure"] = 1013.25
使用 APScheduler(企业级定时调度)
from apscheduler.schedulers.background import BackgroundScheduler
from apscheduler.triggers.interval import IntervalTrigger
from datetime import datetime
import json
import atexit
class AdvancedCacheSync:
def __init__(self, db_path: str = "cache.db"):
self.cache = {}
self.db_path = db_path
self.scheduler = BackgroundScheduler()
self.scheduler.start()
# 注册退出时同步
atexit.register(lambda: self.sync_to_disk())
def sync_to_disk(self):
"""同步缓存到文件"""
try:
with open(self.db_path, 'w', encoding='utf-8') as f:
json.dump(self.cache, f, ensure_ascii=False, default=str)
timestamp = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
print(f"[{timestamp}] 缓存同步完成 - {len(self.cache)} 个条目")
except Exception as e:
print(f"同步出错: {e}")
def schedule_sync(self, seconds: int = 30):
"""
多种调度方式
"""
# 1. 固定间隔
self.scheduler.add_job(
self.sync_to_disk,
IntervalTrigger(seconds=seconds),
id='cache_sync',
replace_existing=True
)
# 2. 可以添加 cron 表达式
# self.scheduler.add_job(
# self.sync_to_disk,
# 'cron',
# hour='0', # 每天午夜执行
# id='midnight_sync'
# )
# 3. 一次性定时任务
# self.scheduler.add_job(
# self.sync_to_disk,
# 'date',
# run_date=datetime(2024, 12, 31, 23, 59)
# )
print(f"AP调度器已启动,同步间隔: {seconds} 秒")
def stop(self):
"""停止调度器"""
self.sync_to_disk() # 最后同步一次
self.scheduler.shutdown(wait=False)
# 使用示例
cache_sync = AdvancedCacheSync("my_app_cache.json")
cache_sync.cache["last_login"] = datetime.now().isoformat()
cache_sync.cache["user_count"] = 42
# 启动定时同步(每10秒)
cache_sync.schedule_sync(seconds=10)
# 主程序逻辑
try:
# 模拟应用程序运行
import time
time.sleep(60) # 运行1分钟
finally:
cache_sync.stop()
带有条件检查的智能同步
import time
from typing import Any, Dict
from datetime import datetime, timedelta
class SmartCacheSync:
"""智能缓存同步 - 只在缓存变化时同步"""
def __init__(self, sync_interval: int = 60, max_dirty_time: int = 300):
"""
:param sync_interval: 检查间隔(秒)
:param max_dirty_time: 最大脏数据存活时间(秒)
"""
self.cache: Dict[str, Any] = {}
self.last_sync_time = time.time()
self.is_dirty = False
self.sync_interval = sync_interval
self.max_dirty_time = max_dirty_time
def set(self, key: str, value: Any):
"""设置缓存值并标记为脏"""
self.cache[key] = value
self.is_dirty = True
def sync_if_needed(self):
"""按需同步"""
current_time = time.time()
# 只在以下情况同步:
# 1. 缓存有变化且距离上次同步超过间隔
# 2. 距离上次同步超过最大脏数据存活时间
should_sync = (
self.is_dirty and
(current_time - self.last_sync_time >= self.sync_interval)
) or (
self.is_dirty and
(current_time - self.last_sync_time >= self.max_dirty_time)
)
if should_sync:
self.sync_to_disk()
def sync_to_disk(self):
"""执行同步"""
try:
# 同步逻辑
with open('smart_cache.json', 'w') as f:
json.dump(self.cache, f)
self.last_sync_time = time.time()
self.is_dirty = False
print(f"智能同步完成 - {len(self.cache)} 个条目")
except Exception as e:
print(f"同步失败: {e}")
def start(self):
"""启动智能同步循环"""
while True:
self.sync_if_needed()
time.sleep(5) # 每5秒检查一次
# 使用示例
smart_cache = SmartCacheSync(sync_interval=30, max_dirty_time=120)
# 在独立线程中运行
import threading
sync_thread = threading.Thread(target=smart_cache.start, daemon=True)
sync_thread.start()
# 主线程中正常使用
smart_cache.set("temperature", 26.5)
smart_cache.set("humidity", 55)
time.sleep(10)
smart_cache.set("temperature", 27.0) # 脏数据更新
推荐方案
- 简单场景:使用
threading.Timer,代码量小,易于理解 - 中等复杂度:使用
schedule库,提供更多时间表达式支持 - 企业级应用:使用
APScheduler,支持持久化、故障恢复、多种触发器
这些方案可以根据你的具体需求选择,比如同步频率、是否需要持久化调度状态等。