本文目录导读:

对于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 BatchSyncStrategy:
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"}
关键设计原则
- 区分核心与非核心:明确哪些数据必须实时同步,哪些可以延迟
- 写入优先:非核心数据的读取可以读取缓存/本地存储,写入则延迟同步
- 幂等性设计:确保延迟同步的重复执行不会产生副作用
- 监控告警:即使容忍延迟,也要监控失败次数和同步延迟
- 优雅降级:当非核心同步组件完全不可用时,脚本应继续运行
选择哪种策略取决于你的具体业务场景:延迟容忍度、数据一致性要求、以及系统架构。