Python脚本如何告警核心模块同步延迟

wen python案例 29

本文目录导读:

Python脚本如何告警核心模块同步延迟

  1. 基础方案:日志监控 + 阈值告警
  2. 使用 Prometheus + Grafana 架构
  3. 集成钉钉/企业微信告警
  4. 使用 Redis 实现分布式同步监控
  5. 完整的生产级方案
  6. 建议的告警策略

在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("监控已停止")

建议的告警策略

  1. 多重阈值:设置 WARNING(5分钟)和 CRITICAL(10分钟)不同级别
  2. 告警抑制:避免告警风暴,设置冷却期
  3. 自动恢复:当延迟恢复正常时发送恢复通知
  4. 历史趋势:记录延迟变化趋势,用于问题排查
  5. 分级通知:根据严重程度通知不同级别的人员

选择哪种方案取决于你的具体环境:

  • 小型项目:使用基础方案或 Redis 方案
  • 中型项目:集成钉钉/企业微信
  • 大型项目:使用 Prometheus + Grafana 架构

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