Python脚本如何定时同步缓存持久化数据

wen python案例 30

本文目录导读:

Python脚本如何定时同步缓存持久化数据

  1. 使用 threading.Timer(简单定时任务)
  2. 使用 schedule 库(更灵活的定时任务)
  3. 使用 APScheduler(企业级定时调度)
  4. 带有条件检查的智能同步
  5. 推荐方案

我来介绍几种在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,支持持久化、故障恢复、多种触发器

这些方案可以根据你的具体需求选择,比如同步频率、是否需要持久化调度状态等。

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