Python脚本如何容忍非核心数据延迟同步

wen python案例 29

本文目录导读:

Python脚本如何容忍非核心数据延迟同步

  1. 异步处理 + 重试机制
  2. 本地缓存 + 批量延迟同步
  3. 基于检查点的容忍策略
  4. 配置化策略选择
  5. 关键设计原则

对于Python脚本容忍非核心数据延迟同步,有几种常见且实用的策略:

异步处理 + 重试机制

import asyncio
import time
from functools import wraps
def retry_on_failure(max_retries=3, delay=1):
    """重试装饰器,仅对非核心操作使用"""
    def decorator(func):
        @wraps(func)
        async def wrapper(*args, **kwargs):
            for attempt in range(max_retries):
                try:
                    return await func(*args, **kwargs)
                except Exception as e:
                    if attempt == max_retries - 1:
                        print(f"非核心数据同步失败(已重试{max_retries}次): {e}")
                        return None  # 或返回默认值
                    await asyncio.sleep(delay * (2 ** attempt))  # 指数退避
            return None
        return wrapper
    return decorator
class DataSyncManager:
    def __init__(self):
        self.pending_queue = asyncio.Queue()
        self.is_running = True
    @retry_on_failure(max_retries=5, delay=2)
    async def sync_non_critical_data(self, data):
        """非核心数据同步操作"""
        # 模拟可能失败的同步
        if data.get('should_fail'):
            raise ConnectionError("模拟同步失败")
        print(f"同步数据: {data}")
        return True
    async def process_with_tolerance(self, data):
        """容忍非核心数据延迟的处理流程"""
        # 1. 立即处理核心数据
        core_result = await self.process_core_data(data)
        # 2. 非核心数据异步处理,不阻塞主流程
        task = asyncio.create_task(
            self.sync_non_critical_data(data.get('non_critical', {}))
        )
        task.add_done_callback(lambda t: print("非核心数据同步完成" if t.result() else "同步失败但已容忍"))
        return core_result
    async def process_core_data(self, data):
        """核心数据处理,即时完成"""
        await asyncio.sleep(0.1)
        print(f"处理核心数据: {data.get('core')}")
        return {"status": "success", "data": data.get('core')}
# 使用示例
async def main():
    manager = DataSyncManager()
    # 正常数据
    await manager.process_with_tolerance({
        "core": "重要数据",
        "non_critical": {"user_stats": [1,2,3]}
    })
    # 非核心数据同步失败的情况
    await manager.process_with_tolerance({
        "core": "重要数据2",
        "non_critical": {"should_fail": True}
    })
asyncio.run(main())

本地缓存 + 批量延迟同步

import threading
import time
from collections import defaultdict
import pickle
class DelayedSyncBuffer:
    """延迟同步缓冲区 - 线程安全"""
    def __init__(self, sync_interval=60, max_batch_size=100):
        self.buffer = defaultdict(list)
        self.sync_interval = sync_interval
        self.max_batch_size = max_batch_size
        self._start_background_sync()
        self.lock = threading.Lock()
    def add(self, data_type, data):
        """添加数据到缓冲区"""
        with self.lock:
            self.buffer[data_type].append(data)
            if len(self.buffer[data_type]) >= self.max_batch_size:
                self._flush_sync(data_type)
    def _start_background_sync(self):
        """启动后台同步线程"""
        def sync_loop():
            while True:
                time.sleep(self.sync_interval)
                self._flush_all()
        thread = threading.Thread(target=sync_loop, daemon=True)
        thread.start()
    def _flush_sync(self, data_type):
        """同步特定类型的数据"""
        with self.lock:
            if data_type not in self.buffer or not self.buffer[data_type]:
                return
            batch = self.buffer[data_type]
            self.buffer[data_type] = []
        try:
            # 批量同步到远程
            self._sync_to_remote(batch)
            print(f"成功同步 {len(batch)} 条 {data_type} 数据")
        except Exception as e:
            print(f"批量同步失败: {e},数据将在下次重试")
            # 重新加入缓冲区
            with self.lock:
                self.buffer[data_type].extend(batch)
    def _flush_all(self):
        """同步所有待处理数据"""
        with self.lock:
            data_types = list(self.buffer.keys())
        for dt in data_types:
            self._flush_sync(dt)
    def _sync_to_remote(self, data_batch):
        """实际执行远程同步"""
        # 模拟可能失败
        if any(d.get('should_fail') for d in data_batch):
            raise ConnectionError("模拟同步失败")
        print(f"远程同步成功: {len(data_batch)} 条记录")
# 使用示例
class ApplicationWithDelayedSync:
    def __init__(self):
        self.sync_buffer = DelayedSyncBuffer(
            sync_interval=30,  # 每30秒同步一次
            max_batch_size=50   # 或达到50条即触发同步
        )
    def process_request(self, request):
        """处理请求,核心逻辑立即执行"""
        # 核心处理
        result = self._core_processing(request)
        # 非核心日志/统计信息 -> 延迟同步
        self.sync_buffer.add(
            data_type="user_analytics",
            data={
                "user_id": request.user_id,
                "action": request.action,
                "timestamp": time.time()
            }
        )
        return result
    def _core_processing(self, request):
        """核心处理逻辑"""
        return {"success": True, "message": "请求已处理"}
# 使用
app = ApplicationWithDelayedSync()
for i in range(10):
    app.process_request(some_request_object)

基于检查点的容忍策略

import threading
import time
import enum
from dataclasses import dataclass
from typing import Optional
@dataclass
class SyncState:
    """同步状态管理"""
    last_sync_time: float = 0
    sync_point: str = ""
    error_count: int = 0
    is_synced: bool = False
class SyncToleranceLevel(enum.Enum):
    STRICT = 0     # 实时同步,失败则重试
    RELAXED = 1    # 延迟同步,可容忍短期失败
    OPTIONAL = 2   # 仅尽力同步,失败可忽略
class TolerantDataSyncManager:
    """容忍性数据同步管理器"""
    def __init__(self):
        self.sync_states = {}
        self.checkpoint_data = {}
        self.lock = threading.RLock()
    def register_sync_point(self, name: str, tolerance: SyncToleranceLevel):
        """注册同步检查点"""
        with self.lock:
            self.sync_states[name] = {
                "tolerance": tolerance,
                "state": SyncState(),
                "pending_data": []
            }
    def sync_with_tolerance(self, name: str, data: dict) -> bool:
        """根据容忍级别同步数据"""
        with self.lock:
            if name not in self.sync_states:
                raise ValueError(f"未知的同步点: {name}")
            config = self.sync_states[name]
            level = config["tolerance"]
            state = config["state"]
        if level == SyncToleranceLevel.STRICT:
            return self._strict_sync(name, data)
        elif level == SyncToleranceLevel.RELAXED:
            return self._relaxed_sync(name, data)
        else:
            return self._optional_sync(name, data)
    def _strict_sync(self, name: str, data: dict) -> bool:
        """严格同步 - 必须成功"""
        max_retries = 3
        for attempt in range(max_retries):
            try:
                if self._do_sync(data):
                    self._update_sync_state(name, True)
                    return True
            except Exception as e:
                if attempt == max_retries - 1:
                    self._update_sync_state(name, False)
                    raise  # 重新抛出,让调用者处理
                time.sleep(1 * (2 ** attempt))
        return False
    def _relaxed_sync(self, name: str, data: dict) -> bool:
        """宽松同步 - 可容忍短期失败"""
        try:
            if self._do_sync(data):
                self._update_sync_state(name, True)
                return True
        except Exception as e:
            # 保存到待同步队列
            with self.lock:
                self.sync_states[name]["pending_data"].append({
                    "data": data,
                    "error": str(e),
                    "timestamp": time.time()
                })
            self._update_sync_state(name, False)
        return False
    def _optional_sync(self, name: str, data: dict) -> bool:
        """可选同步 - 失败完全可忽略"""
        try:
            result = self._do_sync(data)
            if result:
                self._update_sync_state(name, True)
            return result
        except Exception as e:
            # 记录日志但不抛出异常
            self._update_sync_state(name, False)
            return False
    def _do_sync(self, data: dict) -> bool:
        """实际的同步操作"""
        # 模拟可能失败
        if data.get('should_fail'):
            raise ConnectionError("模拟同步失败")
        return True
    def _update_sync_state(self, name: str, success: bool):
        """更新同步状态"""
        with self.lock:
            state = self.sync_states[name]["state"]
            if success:
                state.last_sync_time = time.time()
                state.error_count = 0
                state.is_synced = True
            else:
                state.error_count += 1
                state.is_synced = False
    def retry_pending_syncs(self):
        """重试待同步的数据"""
        with self.lock:
            for name, config in self.sync_states.items():
                if config["pending_data"]:
                    pending = config["pending_data"]
                    config["pending_data"] = []
        # 在锁外重试
        for item in pending:
            self._relaxed_sync(name, item["data"])
# 使用示例
def main():
    manager = TolerantDataSyncManager()
    # 注册不同类型的同步点
    manager.register_sync_point("critical_logs", SyncToleranceLevel.STRICT)
    manager.register_sync_point("user_stats", SyncToleranceLevel.RELAXED)
    manager.register_sync_point("cache_warmup", SyncToleranceLevel.OPTIONAL)
    # 核心日志 - 必须同步成功
    try:
        manager.sync_with_tolerance("critical_logs", {"event": "user_login", "user": "123"})
    except Exception:
        print("关键日志同步失败,需要人工干预")
    # 用户统计 - 可以容忍延迟
    success = manager.sync_with_tolerance("user_stats", {"user_id": "123", "clicks": 5})
    if not success:
        print("用户统计延迟同步")
    # 缓存预热 - 完全可忽略失败
    manager.sync_with_tolerance("cache_warmup", {"api": "user_profile"})
    # 定期重试
    time.sleep(5)
    manager.retry_pending_syncs()
if __name__ == "__main__":
    main()

配置化策略选择

import yaml
# sync_config.yaml
"""
sync_strategies:
  user_analytics:
    mode: batch
    batch_interval: 60
    max_batch_size: 100
    retry_policy:
      max_retries: 5
      backoff_type: exponential
  clickhouse_export:
    mode: async
    failure_tolerance: 2  # 容忍连续失败的次数
    fallback_to_local: true
  third_party_api:
    mode: fire_and_forget
    criticality: low
"""
class ConfigurableSyncStrategy:
    def __init__(self, config_path):
        with open(config_path, 'r') as f:
            self.config = yaml.safe_load(f)
        self.strategies = {
            "batch": BatchSyncStrategy,
            "async": AsyncSyncStrategy,
            "fire_and_forget": FireAndForgetStrategy
        }
    def get_strategy(self, data_type):
        strategy_config = self.config['sync_strategies'].get(data_type)
        if strategy_config:
            strategy_class = self.strategies[strategy_config['mode']]
            return strategy_class(strategy_config)
        return DefaultSyncStrategy()
class Batc‌hSyncStrategy:
    def __init__(self, config):
        self.batch_interval = config.get('batch_interval', 60)
        self.max_batch_size = config.get('max_batch_size', 100)
        # ... 实现
# 使用示例
def sync_data(data_type, data):
    strategy = configurable_sync.get_strategy(data_type)
    result = strategy.sync(data)
    if not result and strategy.is_tolerable():
        return {"status": "accepted", "sync": "delayed"}
    return {"status": "synced"}

关键设计原则

  1. 区分核心与非核心:明确哪些数据必须实时同步,哪些可以延迟
  2. 写入优先:非核心数据的读取可以读取缓存/本地存储,写入则延迟同步
  3. 幂等性设计:确保延迟同步的重复执行不会产生副作用
  4. 监控告警:即使容忍延迟,也要监控失败次数和同步延迟
  5. 优雅降级:当非核心同步组件完全不可用时,脚本应继续运行

选择哪种策略取决于你的具体业务场景:延迟容忍度、数据一致性要求、以及系统架构。

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