本文目录导读:

我来写一个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()
这个同步监控系统的主要特点:
核心功能
-
多数据源支持
- MySQL数据库同步监控
- Redis缓存同步监控
- API接口同步监控
-
灵活的监控配置
- 可自定义检查间隔
- 可配置期望延迟阈值
- 支持告警阈值设置
-
状态管理
- 正常、延迟、失败、未知四种状态
- 自动状态判断
- 历史状态记录
-
并发检查
- 使用ThreadPoolExecutor并行检查
- 提高检查效率
- 不影响其他模块检查
-
告警系统
- 支持多种告警处理器
- 分级告警(延迟、失败)
- 可扩展的告警机制
-
报告生成
- 实时状态查看
- 健康评分计算
- 详细的模块状态信息
使用方法
# 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)
这个系统可以轻松扩展支持更多数据源类型,添加更多告警方式,以及集成到现有的监控系统中。