Python脚本如何适配不同业务同步时效

wen python案例 29

本文目录导读:

Python脚本如何适配不同业务同步时效

  1. 配置化同步间隔(推荐)
  2. 配置文件示例 (sync_config.ini)
  3. 基于CRON表达式的灵活配置
  4. 动态调整同步频率
  5. 多线程/异步支持
  6. 最佳实践建议
  7. 使用示例

我来介绍几种让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)  # 每秒检查一次

这种设计可以轻松适配不同业务的同步时效要求,通过配置化管理,无需修改代码即可调整同步策略。

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