本文目录导读:

在Python中监控核心模块同步延迟并告警,可以根据你的技术栈选择不同的方案,以下是几种常见且实用的方法:
基础方案:日志监控 + 阈值告警
import time
import logging
from datetime import datetime, timedelta
class SyncMonitor:
def __init__(self, max_delay_seconds=300, alert_threshold=3):
self.max_delay = max_delay_seconds
self.alert_count = 0
self.alert_threshold = alert_threshold
def check_sync_delay(self, last_sync_time):
"""检查同步延迟"""
now = datetime.now()
delay = (now - last_sync_time).total_seconds()
if delay > self.max_delay:
self.alert_count += 1
self.trigger_alert(delay)
else:
self.alert_count = 0
def trigger_alert(self, delay):
"""触发告警"""
if self.alert_count >= self.alert_threshold:
# 发送告警(邮件、短信、钉钉等)
print(f"[ALERT] 核心模块同步延迟: {delay:.2f}秒")
# 具体告警实现...
self.send_alert_email(delay)
def send_alert_email(self, delay):
"""发送邮件告警示例"""
import smtplib
from email.mime.text import MIMEText
msg = MIMEText(f"核心模块同步延迟: {delay:.2f}秒")
msg['Subject'] = '核心模块同步延迟告警'
msg['From'] = 'monitor@example.com'
msg['To'] = 'admin@example.com'
# 实际发送代码...
使用 Prometheus + Grafana 架构
from prometheus_client import Gauge, start_http_server
import time
import random
# 定义指标
sync_delay = Gauge('core_module_sync_delay_seconds',
'Core module synchronization delay')
def check_sync_delay():
"""检查同步延迟并更新指标"""
# 这里实现实际的延迟检查逻辑
delay = get_actual_sync_delay() # 你的实际检查方法
sync_delay.set(delay)
if delay > 300: # 5分钟
# Prometheus 会处理告警规则
print(f"High sync delay detected: {delay}s")
def get_actual_sync_delay():
"""获取实际同步延迟"""
# 示例:查询数据库中的最后同步时间
# last_sync = db.query("SELECT max(sync_time) FROM sync_log")
# return (datetime.now() - last_sync).total_seconds()
return random.uniform(0, 500)
if __name__ == '__main__':
start_http_server(8000)
while True:
check_sync_delay()
time.sleep(60) # 每分钟检查一次
配合 prometheus.yml 告警规则:
groups:
- name: sync_alerts
rules:
- alert: CoreModuleSyncDelay
expr: core_module_sync_delay_seconds > 300
for: 5m
annotations:
summary: "核心模块同步延迟过高"
集成钉钉/企业微信告警
import requests
import json
from datetime import datetime
class DingTalkAlert:
def __init__(self, webhook_url):
self.webhook_url = webhook_url
def send_alert(self, module_name, delay, severity='critical'):
"""发送钉钉告警"""
current_time = datetime.now().strftime('%Y-%m-%d %H:%M:%S')
message = {
"msgtype": "markdown",
"markdown": {
"title": f"[{severity.upper()}] 核心模块同步延迟告警",
"text": f"### 同步延迟告警\n"
f"- **模块**: {module_name}\n"
f"- **延迟时间**: {delay:.2f}秒\n"
f"- **告警时间**: {current_time}\n"
f"- **严重级别**: {severity}\n"
f"⚠️ 请立即检查同步服务状态!"
}
}
response = requests.post(
self.webhook_url,
data=json.dumps(message),
headers={'Content-Type': 'application/json'}
)
return response.json()
# 集成监控逻辑
class SyncDelayMonitor:
def __init__(self, dingtalk_webhook_url):
self.alert = DingTalkAlert(dingtalk_webhook_url)
self.delay_threshold = 300 # 5分钟
def monitor(self, module, last_sync_time):
now = datetime.now()
delay = (now - last_sync_time).total_seconds()
if delay > self.delay_threshold:
self.alert.send_alert(
module_name=module,
delay=delay,
severity='critical' if delay > 600 else 'warning'
)
使用 Redis 实现分布式同步监控
import redis
import time
from datetime import datetime
class RedisSyncMonitor:
def __init__(self, redis_host='localhost', redis_port=6379):
self.redis = redis.StrictRedis(
host=redis_host,
port=redis_port,
decode_responses=True
)
self.sync_key = 'core_module:last_sync_time'
self.alert_key = 'core_module:alert_count'
self.max_delay = 300
def update_sync_time(self):
"""更新同步时间"""
self.redis.set(self.sync_key, datetime.now().isoformat())
def check_delay(self):
"""检查延迟"""
last_sync_str = self.redis.get(self.sync_key)
if not last_sync_str:
self.trigger_alert("从未同步")
return
last_sync = datetime.fromisoformat(last_sync_str)
delay = (datetime.now() - last_sync).total_seconds()
if delay > self.max_delay:
alert_count = self.redis.incr(self.alert_key)
if alert_count >= 3: # 连续告警3次触发
self.trigger_alert(f"同步延迟: {delay:.2f}秒")
else:
# 恢复正常,重置计数
self.redis.set(self.alert_key, 0)
def trigger_alert(self, message):
"""触发告警"""
print(f"[ALERT] 核心模块同步异常: {message}")
# 这里集成具体的告警方式
完整的生产级方案
import threading
import time
from collections import deque
from datetime import datetime, timedelta
class ProductionSyncMonitor:
def __init__(self):
self.config = {
'check_interval': 60, # 检查间隔(秒)
'max_delay': 300, # 最大允许延迟(秒)
'alert_cooldown': 300, # 告警冷却时间(秒)
'history_window': 3600 # 历史窗口(秒)
}
self.alert_history = deque(maxlen=100)
self.last_alert_time = datetime.min
self.monitoring = True
def start_monitoring(self):
"""启动监控线程"""
monitor_thread = threading.Thread(target=self._monitor_loop)
monitor_thread.daemon = True
monitor_thread.start()
def _monitor_loop(self):
"""监控循环"""
while self.monitoring:
try:
self._check_sync_status()
time.sleep(self.config['check_interval'])
except Exception as e:
print(f"监控异常: {e}")
def _check_sync_status(self):
"""检查同步状态"""
# 1. 获取核心模块同步状态
sync_status = self.get_core_module_status()
# 2. 计算延迟
delay = self.calculate_sync_delay(sync_status)
# 3. 记录历史
self.alert_history.append({
'time': datetime.now(),
'delay': delay,
'status': 'normal' if delay < self.config['max_delay'] else 'delayed'
})
# 4. 判断是否需要告警
if delay > self.config['max_delay']:
self._handle_delay_alert(delay)
def _handle_delay_alert(self, delay):
"""处理延迟告警"""
now = datetime.now()
cooldown = timedelta(seconds=self.config['alert_cooldown'])
# 冷却期检查
if (now - self.last_alert_time) < cooldown:
return
self.last_alert_time = now
# 发送多重告警
self.send_multi_channel_alert({
'level': 'critical',
'message': f'核心模块同步延迟: {delay:.2f}秒',
'time': now,
'recommendation': '请检查网络连接和同步服务状态'
})
def get_core_module_status(self):
"""获取核心模块状态(具体实现)"""
# 这里实现实际的状态获取逻辑
pass
def calculate_sync_delay(self, status):
"""计算同步延迟"""
# 实现延迟计算逻辑
pass
def send_multi_channel_alert(self, alert_data):
"""多通道告警"""
# 邮件
self.send_email(alert_data)
# 短信
self.send_sms(alert_data)
# 即时通讯
self.send_im_message(alert_data)
def send_email(self, data):
"""发送邮件告警"""
print(f"[EMAIL] {data['message']}")
def send_sms(self, data):
"""发送短信告警"""
print(f"[SMS] {data['message']}")
def send_im_message(self, data):
"""发送即时消息"""
print(f"[IM] {data['message']}")
# 使用示例
if __name__ == '__main__':
monitor = ProductionSyncMonitor()
monitor.start_monitoring()
# 保持主线程运行
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
monitor.monitoring = False
print("监控已停止")
建议的告警策略
- 多重阈值:设置 WARNING(5分钟)和 CRITICAL(10分钟)不同级别
- 告警抑制:避免告警风暴,设置冷却期
- 自动恢复:当延迟恢复正常时发送恢复通知
- 历史趋势:记录延迟变化趋势,用于问题排查
- 分级通知:根据严重程度通知不同级别的人员
选择哪种方案取决于你的具体环境:
- 小型项目:使用基础方案或 Redis 方案
- 中型项目:集成钉钉/企业微信
- 大型项目:使用 Prometheus + Grafana 架构