Python脚本如何监控各模块同步时效状态

wen python案例 26

本文目录导读:

Python脚本如何监控各模块同步时效状态

  1. 核心功能
  2. 使用方法

我来写一个Python脚本来监控各模块的同步时效状态,这个脚本会检查不同模块的数据同步是否在预期的时间内完成。

#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""
模块同步时效监控系统
监控各模块数据同步的时效性和状态
"""
import time
import json
import logging
import threading
from datetime import datetime, timedelta
from typing import Dict, List, Optional, Tuple
from dataclasses import dataclass, asdict
from enum import Enum
import requests
import mysql.connector
from mysql.connector import Error
import redis
from concurrent.futures import ThreadPoolExecutor, as_completed
# 配置日志
logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s - %(name)s - %(levelname)s - %(message)s',
    handlers=[
        logging.FileHandler('sync_monitor.log'),
        logging.StreamHandler()
    ]
)
logger = logging.getLogger(__name__)
class SyncStatus(Enum):
    """同步状态枚举"""
    NORMAL = "正常"
    DELAYED = "延迟"
    FAILED = "失败"
    UNKNOWN = "未知"
@dataclass
class ModuleSyncConfig:
    """模块同步配置"""
    module_name: str
    source_type: str  # 'mysql', 'redis', 'api'
    source_config: dict
    check_interval: int  # 检查间隔(秒)
    expected_delay: int  # 期望延迟(秒)
    alert_threshold: int  # 告警阈值(秒)
    query_template: str  # 查询模板或API路径
@dataclass
class SyncStatusRecord:
    """同步状态记录"""
    module_name: str
    status: SyncStatus
    current_delay: int  # 当前延迟(秒)
    last_sync_time: Optional[datetime]
    check_time: datetime
    message: str
    metrics: dict = None
class DatabaseConnector:
    """数据库连接管理"""
    def __init__(self):
        self.connections = {}
        self.lock = threading.Lock()
    def get_mysql_connection(self, config: dict):
        """获取MySQL连接"""
        conn_key = f"mysql_{config.get('host')}_{config.get('database')}"
        with self.lock:
            if conn_key not in self.connections:
                try:
                    connection = mysql.connector.connect(
                        host=config.get('host', 'localhost'),
                        port=config.get('port', 3306),
                        database=config.get('database'),
                        user=config.get('user'),
                        password=config.get('password'),
                        connection_timeout=10
                    )
                    self.connections[conn_key] = connection
                    logger.info(f"MySQL连接创建成功: {conn_key}")
                except Error as e:
                    logger.error(f"MySQL连接失败 {conn_key}: {e}")
                    return None
            return self.connections[conn_key]
    def get_redis_connection(self, config: dict):
        """获取Redis连接"""
        conn_key = f"redis_{config.get('host')}_{config.get('db', 0)}"
        with self.lock:
            if conn_key not in self.connections:
                try:
                    connection = redis.Redis(
                        host=config.get('host', 'localhost'),
                        port=config.get('port', 6379),
                        db=config.get('db', 0),
                        decode_responses=True,
                        socket_connect_timeout=5
                    )
                    # 测试连接
                    connection.ping()
                    self.connections[conn_key] = connection
                    logger.info(f"Redis连接创建成功: {conn_key}")
                except redis.ConnectionError as e:
                    logger.error(f"Redis连接失败 {conn_key}: {e}")
                    return None
            return self.connections[conn_key]
    def close_all(self):
        """关闭所有连接"""
        for key, conn in self.connections.items():
            try:
                if 'mysql' in key:
                    conn.close()
                elif 'redis' in key:
                    conn.close()
                logger.info(f"关闭连接: {key}")
            except Exception as e:
                logger.error(f"关闭连接失败 {key}: {e}")
        self.connections.clear()
class ModuleSyncMonitor:
    """模块同步监控器"""
    def __init__(self, configs: List[ModuleSyncConfig]):
        self.configs = configs
        self.db_connector = DatabaseConnector()
        self.executor = ThreadPoolExecutor(max_workers=10)
        self.status_history = {}  # 存储历史状态
        self.alert_handlers = []  # 告警处理器列表
        self.running = False
    def add_alert_handler(self, handler):
        """添加告警处理器"""
        self.alert_handlers.append(handler)
    def check_mysql_sync(self, config: ModuleSyncConfig) -> SyncStatusRecord:
        """检查MySQL同步状态"""
        try:
            conn = self.db_connector.get_mysql_connection(config.source_config)
            if not conn:
                return SyncStatusRecord(
                    module_name=config.module_name,
                    status=SyncStatus.FAILED,
                    current_delay=-1,
                    last_sync_time=None,
                    check_time=datetime.now(),
                    message="数据库连接失败"
                )
            cursor = conn.cursor(dictionary=True)
            # 获取最新同步时间
            sync_query = config.query_template
            cursor.execute(sync_query)
            result = cursor.fetchone()
            if result and 'last_sync_time' in result:
                last_sync_time = result['last_sync_time']
                if isinstance(last_sync_time, datetime):
                    current_delay = (datetime.now() - last_sync_time).total_seconds()
                    # 确定状态
                    if current_delay <= config.expected_delay:
                        status = SyncStatus.NORMAL
                    elif current_delay <= config.alert_threshold:
                        status = SyncStatus.DELAYED
                    else:
                        status = SyncStatus.FAILED
                    return SyncStatusRecord(
                        module_name=config.module_name,
                        status=status,
                        current_delay=int(current_delay),
                        last_sync_time=last_sync_time,
                        check_time=datetime.now(),
                        message=f"检查完成,延迟{int(current_delay)}秒",
                        metrics=result
                    )
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.UNKNOWN,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message="无法获取同步时间"
            )
        except Exception as e:
            logger.error(f"检查MySQL同步失败 {config.module_name}: {e}")
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.FAILED,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message=f"检查异常: {str(e)}"
            )
        finally:
            if 'cursor' in locals():
                cursor.close()
    def check_redis_sync(self, config: ModuleSyncConfig) -> SyncStatusRecord:
        """检查Redis同步状态"""
        try:
            redis_conn = self.db_connector.get_redis_connection(config.source_config)
            if not redis_conn:
                return SyncStatusRecord(
                    module_name=config.module_name,
                    status=SyncStatus.FAILED,
                    current_delay=-1,
                    last_sync_time=None,
                    check_time=datetime.now(),
                    message="Redis连接失败"
                )
            # 获取同步时间戳
            last_sync_timestamp = redis_conn.get(config.query_template)
            if last_sync_timestamp:
                last_sync_time = datetime.fromtimestamp(float(last_sync_timestamp))
                current_delay = (datetime.now() - last_sync_time).total_seconds()
                if current_delay <= config.expected_delay:
                    status = SyncStatus.NORMAL
                elif current_delay <= config.alert_threshold:
                    status = SyncStatus.DELAYED
                else:
                    status = SyncStatus.FAILED
                # 获取更多指标
                metrics = {
                    'last_sync_timestamp': last_sync_timestamp,
                    'key': config.query_template
                }
                return SyncStatusRecord(
                    module_name=config.module_name,
                    status=status,
                    current_delay=int(current_delay),
                    last_sync_time=last_sync_time,
                    check_time=datetime.now(),
                    message=f"检查完成,延迟{int(current_delay)}秒",
                    metrics=metrics
                )
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.UNKNOWN,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message="Redis中未找到同步时间戳"
            )
        except Exception as e:
            logger.error(f"检查Redis同步失败 {config.module_name}: {e}")
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.FAILED,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message=f"检查异常: {str(e)}"
            )
    def check_api_sync(self, config: ModuleSyncConfig) -> SyncStatusRecord:
        """检查API同步状态"""
        try:
            url = config.source_config.get('url', '') + config.query_template
            headers = config.source_config.get('headers', {})
            response = requests.get(
                url,
                headers=headers,
                timeout=10
            )
            if response.status_code == 200:
                data = response.json()
                # 解析响应数据
                if 'last_sync_time' in data:
                    last_sync_time = datetime.fromisoformat(data['last_sync_time'])
                    current_delay = (datetime.now() - last_sync_time).total_seconds()
                    if current_delay <= config.expected_delay:
                        status = SyncStatus.NORMAL
                    elif current_delay <= config.alert_threshold:
                        status = SyncStatus.DELAYED
                    else:
                        status = SyncStatus.FAILED
                    return SyncStatusRecord(
                        module_name=config.module_name,
                        status=status,
                        current_delay=int(current_delay),
                        last_sync_time=last_sync_time,
                        check_time=datetime.now(),
                        message=f"API检查完成,延迟{int(current_delay)}秒",
                        metrics=data
                    )
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.UNKNOWN,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message=f"API返回异常: {response.status_code}"
            )
        except Exception as e:
            logger.error(f"检查API同步失败 {config.module_name}: {e}")
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.FAILED,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message=f"API检查异常: {str(e)}"
            )
    def check_module_sync(self, config: ModuleSyncConfig) -> SyncStatusRecord:
        """检查单个模块同步状态"""
        check_functions = {
            'mysql': self.check_mysql_sync,
            'redis': self.check_redis_sync,
            'api': self.check_api_sync
        }
        check_func = check_functions.get(config.source_type)
        if not check_func:
            return SyncStatusRecord(
                module_name=config.module_name,
                status=SyncStatus.UNKNOWN,
                current_delay=-1,
                last_sync_time=None,
                check_time=datetime.now(),
                message=f"不支持的数据源类型: {config.source_type}"
            )
        return check_func(config)
    def check_all_modules(self):
        """检查所有模块同步状态"""
        futures = []
        results = []
        for config in self.configs:
            future = self.executor.submit(self.check_module_sync, config)
            futures.append(future)
        for future in as_completed(futures):
            try:
                result = future.result()
                results.append(result)
                # 更新历史记录
                if result.module_name not in self.status_history:
                    self.status_history[result.module_name] = []
                self.status_history[result.module_name].append(result)
                # 保留最近100条记录
                if len(self.status_history[result.module_name]) > 100:
                    self.status_history[result.module_name] = \
                        self.status_history[result.module_name][-100:]
                # 触发告警
                if result.status in [SyncStatus.DELAYED, SyncStatus.FAILED]:
                    self._trigger_alerts(result)
                logger.info(f"模块 {result.module_name} 状态: {result.status.value}")
            except Exception as e:
                logger.error(f"检查模块同步失败: {e}")
        return results
    def _trigger_alerts(self, record: SyncStatusRecord):
        """触发告警"""
        for handler in self.alert_handlers:
            try:
                handler(record)
            except Exception as e:
                logger.error(f"告警处理器执行失败: {e}")
    def get_status_summary(self) -> Dict:
        """获取状态摘要"""
        summary = {
            'total_modules': len(self.configs),
            'normal': 0,
            'delayed': 0,
            'failed': 0,
            'unknown': 0,
            'modules': []
        }
        for config in self.configs:
            module_history = self.status_history.get(config.module_name, [])
            if module_history:
                latest_status = module_history[-1]
                summary['modules'].append(asdict(latest_status))
                if latest_status.status == SyncStatus.NORMAL:
                    summary['normal'] += 1
                elif latest_status.status == SyncStatus.DELAYED:
                    summary['delayed'] += 1
                elif latest_status.status == SyncStatus.FAILED:
                    summary['failed'] += 1
                else:
                    summary['unknown'] += 1
        summary['health_score'] = (summary['normal'] / summary['total_modules']) * 100 \
            if summary['total_modules'] > 0 else 0
        return summary
    def generate_report(self) -> str:
        """生成报告"""
        summary = self.get_status_summary()
        report = []
        report.append("=" * 60)
        report.append("模块同步时效性监控报告")
        report.append(f"生成时间: {datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
        report.append("=" * 60)
        report.append(f"总模块数: {summary['total_modules']}")
        report.append(f"正常: {summary['normal']}")
        report.append(f"延迟: {summary['delayed']}")
        report.append(f"失败: {summary['failed']}")
        report.append(f"未知: {summary['unknown']}")
        report.append(f"健康评分: {summary['health_score']:.1f}%")
        report.append("-" * 60)
        for module in summary['modules']:
            report.append(f"\n模块: {module['module_name']}")
            report(f"状态: {module['status']['value']}")
            report(f"当前延迟: {module['current_delay']}秒")
            report(f"最后同步: {module['last_sync_time']}")
            report(f"检查时间: {module['check_time']}")
            report(f"信息: {module['message']}")
        report.append("=" * 60)
        return "\n".join(report)
    def start_monitoring(self, interval: int = 60):
        """启动持续监控"""
        self.running = True
        logger.info(f"开始监控,检查间隔: {interval}秒")
        while self.running:
            try:
                start_time = time.time()
                # 执行检查
                results = self.check_all_modules()
                # 生成并保存报告
                report = self.generate_report()
                logger.info(f"\n{report}")
                # 计算等待时间
                elapsed = time.time() - start_time
                wait_time = max(0, interval - elapsed)
                time.sleep(wait_time)
            except KeyboardInterrupt:
                logger.info("收到中断信号,停止监控")
                self.stop_monitoring()
                break
            except Exception as e:
                logger.error(f"监控循环异常: {e}")
                time.sleep(10)  # 异常后等待10秒重试
    def stop_monitoring(self):
        """停止监控"""
        self.running = False
        self.executor.shutdown(wait=False)
        self.db_connector.close_all()
        logger.info("监控已停止")
# 告警处理器示例
def email_alert_handler(record: SyncStatusRecord):
    """邮件告警处理器"""
    if record.status == SyncStatus.FAILED:
        logger.warning(f"[邮件告警] 模块 {record.module_name} 同步失败!")
    elif record.status == SyncStatus.DELAYED and record.current_delay > 300:
        logger.warning(f"[邮件告警] 模块 {record.module_name} 延迟严重: {record.current_delay}秒")
def sms_alert_handler(record: SyncStatusRecord):
    """短信告警处理器"""
    if record.status == SyncStatus.FAILED and record.current_delay > 600:
        logger.critical(f"[短信告警] 模块 {record.module_name} 同步失败超过10分钟!")
# 配置示例
def get_sample_configs() -> List[ModuleSyncConfig]:
    """获取示例配置"""
    return [
        ModuleSyncConfig(
            module_name="订单同步模块",
            source_type="mysql",
            source_config={
                'host': 'localhost',
                'port': 3306,
                'database': 'orders_db',
                'user': 'monitor_user',
                'password': 'monitor_pass'
            },
            check_interval=60,
            expected_delay=300,  # 期望5分钟内同步
            alert_threshold=600,  # 10分钟告警
            query_template="""
                SELECT MAX(updated_at) as last_sync_time, 
                       COUNT(*) as total_records,
                       MAX(id) as max_id
                FROM orders 
                WHERE status = 'synced'
            """
        ),
        ModuleSyncConfig(
            module_name="用户同步模块",
            source_type="redis",
            source_config={
                'host': 'localhost',
                'port': 6379,
                'db': 0
            },
            check_interval=120,
            expected_delay=60,
            alert_threshold=180,
            query_template="user:sync:timestamp"
        ),
        ModuleSyncConfig(
            module_name="数据仓库同步",
            source_type="api",
            source_config={
                'url': 'http://localhost:8080',
                'headers': {'Authorization': 'Bearer your-token'}
            },
            check_interval=300,
            expected_delay=3600,  # 期望1小时内同步
            alert_threshold=7200,  # 2小时告警
            query_template="/api/v1/sync/status"
        )
    ]
def main():
    """主函数"""
    # 获取配置
    configs = get_sample_configs()
    # 创建监控器
    monitor = ModuleSyncMonitor(configs)
    # 添加告警处理器
    monitor.add_alert_handler(email_alert_handler)
    monitor.add_alert_handler(sms_alert_handler)
    # 启动监控(每60秒检查一次)
    monitor.start_monitoring(interval=60)
if __name__ == "__main__":
    main()

这个同步监控系统的主要特点:

核心功能

  1. 多数据源支持

    • MySQL数据库同步监控
    • Redis缓存同步监控
    • API接口同步监控
  2. 灵活的监控配置

    • 可自定义检查间隔
    • 可配置期望延迟阈值
    • 支持告警阈值设置
  3. 状态管理

    • 正常、延迟、失败、未知四种状态
    • 自动状态判断
    • 历史状态记录
  4. 并发检查

    • 使用ThreadPoolExecutor并行检查
    • 提高检查效率
    • 不影响其他模块检查
  5. 告警系统

    • 支持多种告警处理器
    • 分级告警(延迟、失败)
    • 可扩展的告警机制
  6. 报告生成

    • 实时状态查看
    • 健康评分计算
    • 详细的模块状态信息

使用方法

# 1. 配置模块
configs = [
    ModuleSyncConfig(
        module_name="订单同步",
        source_type="mysql",
        source_config={...},
        check_interval=60,
        expected_delay=300,
        alert_threshold=600,
        query_template="SELECT MAX(updated_at) as last_sync_time FROM orders"
    ),
    # 更多模块配置...
]
# 2. 创建监控器
monitor = ModuleSyncMonitor(configs)
# 3. 添加告警处理器
monitor.add_alert_handler(your_alert_handler)
# 4. 启动监控
monitor.start_monitoring(interval=60)

这个系统可以轻松扩展支持更多数据源类型,添加更多告警方式,以及集成到现有的监控系统中。

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