本文目录导读:

我来介绍几种让Python脚本适配不同业务同步时效的策略:
配置化同步间隔(推荐)
import time
import configparser
from datetime import datetime, timedelta
import logging
class SyncScheduler:
def __init__(self, config_file='sync_config.ini'):
self.config = configparser.ConfigParser()
self.config.read(config_file)
self.logger = self._setup_logger()
def _setup_logger(self):
logging.basicConfig(level=logging.INFO,
format='%(asctime)s - %(levelname)s - %(message)s')
return logging.getLogger(__name__)
def get_sync_interval(self, business_type):
"""根据业务类型获取同步间隔"""
try:
# 从配置文件获取各业务的同步策略
interval = self.config.get('sync_intervals', business_type, fallback='300')
return int(interval)
except Exception as e:
self.logger.error(f"获取同步间隔失败: {e}")
return 300 # 默认5分钟
def should_sync_now(self, last_sync_time, business_type):
"""判断是否应该执行同步"""
interval = self.get_sync_interval(business_type)
next_sync_time = last_sync_time + timedelta(seconds=interval)
return datetime.now() >= next_sync_time
# 示例使用
class BusinessSyncService:
def __init__(self):
self.scheduler = SyncScheduler()
# 记录各业务最后同步时间
self.last_sync_times = {}
def sync_data(self, business_type):
"""执行数据同步"""
current_time = datetime.now()
if business_type not in self.last_sync_times:
self.last_sync_times[business_type] = current_time - timedelta(hours=1)
if self.scheduler.should_sync_now(self.last_sync_times[business_type], business_type):
self._perform_sync(business_type)
self.last_sync_times[business_type] = current_time
return True
return False
def _perform_sync(self, business_type):
"""实际同步逻辑"""
# 不同业务的同步逻辑
sync_methods = {
'realtime': self._realtime_sync,
'near_realtime': self._near_realtime_sync,
'batch': self._batch_sync
}
method = sync_methods.get(business_type, self._default_sync)
method()
def _realtime_sync(self):
print("执行实时同步(间隔5秒)")
def _near_realtime_sync(self):
print("执行近实时同步(间隔30秒)")
def _batch_sync(self):
print("执行批量同步(间隔1小时)")
def _default_sync(self):
print("执行默认同步(间隔5分钟)")
配置文件示例 (sync_config.ini)
[sync_intervals] # 不同业务类型的同步间隔(秒) realtime = 5 near_realtime = 30 batch = 3600 hourly = 3600 daily = 86400 weekly = 604800 [business_rules] # 各业务的同步规则 order_sync = realtime inventory_sync = near_realtime report_sync = batch price_sync = 60
基于CRON表达式的灵活配置
from croniter import croniter
import datetime
class CronBasedScheduler:
def __init__(self):
self.cron_expressions = {
'order': '*/5 * * * *', # 每5秒
'inventory': '*/30 * * * *', # 每30秒
'report': '0 * * * *', # 每小时整点
'daily_stats': '0 2 * * *', # 每天凌晨2点
'weekly_report': '0 3 * * 1' # 每周一凌晨3点
}
def get_next_sync_time(self, business_type):
"""获取下一次同步时间"""
cron_expr = self.cron_expressions.get(business_type, '*/5 * * * *')
base_time = datetime.datetime.now()
cron = croniter(cron_expr, base_time)
return cron.get_next(datetime.datetime)
def is_time_to_sync(self, business_type, last_sync_time=None):
"""检查是否该同步了"""
cron_expr = self.cron_expressions.get(business_type, '*/5 * * * *')
base_time = last_sync_time or datetime.datetime.now()
cron = croniter(cron_expr, base_time)
next_sync = cron.get_next(datetime.datetime)
return datetime.datetime.now() >= next_sync
动态调整同步频率
class AdaptiveSyncScheduler:
"""自适应同步调度器"""
def __init__(self):
self.business_configs = {
'order': {'base_interval': 5, 'peak_multiplier': 2},
'inventory': {'base_interval': 30, 'peak_multiplier': 1},
'report': {'base_interval': 3600, 'peak_interval': 1800}
}
def get_adaptive_interval(self, business_type):
"""根据业务负载自适应调整间隔"""
config = self.business_configs.get(business_type, {})
base_interval = config.get('base_interval', 300)
# 获取当前业务负载情况(示例)
current_load = self._get_business_load(business_type)
if current_load > 0.8: # 高负载
return base_interval * config.get('peak_multiplier', 1)
elif current_load < 0.2: # 低负载
return base_interval // 2
else:
return base_interval
def _get_business_load(self, business_type):
"""获取业务负载(示例方法)"""
# 实际应从监控系统获取
return 0.5 # 返回0-1之间的值
多线程/异步支持
import asyncio
from concurrent.futures import ThreadPoolExecutor
class AsyncSyncManager:
def __init__(self):
self.executor = ThreadPoolExecutor(max_workers=5)
async def run_sync_tasks(self):
"""异步执行所有同步任务"""
tasks = []
business_types = ['realtime', 'near_realtime', 'batch']
for biz_type in business_types:
task = asyncio.create_task(self._run_business_sync(biz_type))
tasks.append(task)
await asyncio.gather(*tasks)
async def _run_business_sync(self, business_type):
"""异步运行单个业务同步"""
interval = self._get_business_interval(business_type)
while True:
# 在单独的线程中执行同步
loop = asyncio.get_event_loop()
await loop.run_in_executor(
self.executor,
self._sync_business,
business_type
)
await asyncio.sleep(interval)
def _get_business_interval(self, business_type):
intervals = {'realtime': 5, 'near_realtime': 30, 'batch': 3600}
return intervals.get(business_type, 300)
def _sync_business(self, business_type):
print(f"同步业务: {business_type}")
最佳实践建议
配置文件管理
import yaml
class BusinessSyncConfig:
def __init__(self, config_file='sync_config.yaml'):
with open(config_file, 'r') as f:
self.config = yaml.safe_load(f)
def get_config(self, business_type):
return self.config.get(business_type, {})
监控和告警
class SyncMonitor:
def __init__(self):
self.metrics = {}
def record_sync_duration(self, business_type, duration):
if business_type not in self.metrics:
self.metrics[business_type] = []
self.metrics[business_type].append({
'time': datetime.now(),
'duration': duration
})
def check_sync_health(self, business_type):
# 检查同步是否超时
recent_metrics = self.metrics.get(business_type, [])[-10:]
if recent_metrics:
avg_duration = sum(m['duration'] for m in recent_metrics) / len(recent_metrics)
if avg_duration > self._get_timeout_threshold(business_type):
self._send_alert(f"{business_type} 同步超时")
使用示例
# 初始化
sync_service = BusinessSyncService()
# 主循环
while True:
sync_service.sync_data('order') # 实时同步
sync_service.sync_data('inventory') # 近实时同步
sync_service.sync_data('report') # 批量同步
time.sleep(1) # 每秒检查一次
这种设计可以轻松适配不同业务的同步时效要求,通过配置化管理,无需修改代码即可调整同步策略。