第二十九节:自动化运维:定时任务调度、异常监控报警、日志管理

做量化交易,尤其是基差套利,最怕什么?

怕策略跑着跑着停了,怕半夜数据源断了没人知道,怕服务器重启后脚本没起来。我见过太多团队,策略逻辑写得漂漂亮亮,最后死在运维上。

说白了,自动化运维就是给你的策略上个“保险”。今天咱们就聊聊,怎么把这套东西搭起来。

一、定时任务调度:让机器替你干活

基差套利策略,通常需要定时拉取数据、计算价差、执行交易。手动操作?不存在的。我个人习惯用 APScheduler 来做定时任务调度,比 crontab 更灵活,也更容易和 Python 代码集成。

核心思路:把任务拆成“触发器 + 执行器”。触发器决定什么时候跑,执行器决定怎么跑。

1.1 基础调度示例

from apscheduler.schedulers.blocking import BlockingScheduler
from datetime import datetime

def fetch_market_data():
    print(f"{datetime.now()} - 拉取市场数据...")
    # 这里写你的数据拉取逻辑

def calculate_basis():
    print(f"{datetime.now()} - 计算基差...")
    # 这里写基差计算逻辑

def execute_trade():
    print(f"{datetime.now()} - 执行交易...")
    # 这里写交易执行逻辑

scheduler = BlockingScheduler()

# 每天开盘前拉取数据
scheduler.add_job(fetch_market_data, 'cron', hour=8, minute=55)
# 每5分钟计算一次基差
scheduler.add_job(calculate_basis, 'interval', minutes=5)
# 每天收盘后执行一次交易
scheduler.add_job(execute_trade, 'cron', hour=15, minute=5)

scheduler.start()

嗯,这里要注意:BlockingScheduler 会阻塞主线程。如果你的策略本身是事件驱动的,建议用 BackgroundScheduler,它会在后台跑。

我的经验:曾经有个策略,我用了 interval 模式每30秒跑一次。结果某次数据源响应慢了,任务堆积,最后内存爆了。后来我加了 max_instances=1 参数,确保同一时间只有一个实例在跑。

1.2 任务依赖与异常处理

实际项目中,任务之间往往有依赖关系。比如:先拉数据,再算基差,最后执行交易。我习惯用 任务链 的方式处理:

def job_chain():
    try:
        fetch_market_data()
        calculate_basis()
        execute_trade()
    except Exception as e:
        print(f"任务链执行失败: {e}")
        # 触发报警

scheduler.add_job(job_chain, 'cron', hour=9, minute=0)

这样写的好处是:要么全成功,要么全失败。不会出现“数据拉了但没算基差”这种尴尬情况。

二、异常监控报警:别等亏钱了才知道

我记得刚入行那会儿,有个策略跑了一周,收益曲线挺漂亮。结果有一天突然暴跌,查日志才发现——数据源第三天就断了,策略一直在用旧数据交易。

从那以后,我养成了一个习惯:任何可能出问题的地方,都要有报警

2.1 报警分级

级别 触发条件 通知方式 响应时间
P0(致命) 策略崩溃、数据源全断 电话 + 短信 + 微信 5分钟内
P1(严重) 单数据源异常、交易延迟 微信 + 邮件 15分钟内
P2(警告) 基差偏离阈值、持仓超限 邮件 1小时内
P3(信息) 定时任务完成、版本更新 日志记录 无需响应

2.2 实现一个简单的报警器

import smtplib
import requests
from email.mime.text import MIMEText

class AlertManager:
    def __init__(self):
        self.webhook_url = "https://qyapi.weixin.qq.com/cgi-bin/webhook/send?key=YOUR_KEY"
        self.email_config = {
            "smtp_server": "smtp.qq.com",
            "smtp_port": 587,
            "sender": "your_email@qq.com",
            "password": "your_password",
            "receivers": ["ops@example.com"]
        }
    
    def send_wechat(self, message):
        """发送企业微信消息"""
        data = {
            "msgtype": "text",
            "text": {"content": f"[基差套利报警] {message}"}
        }
        try:
            requests.post(self.webhook_url, json=data, timeout=5)
        except Exception as e:
            print(f"微信报警失败: {e}")
    
    def send_email(self, subject, body):
        """发送邮件报警"""
        msg = MIMEText(body, 'plain', 'utf-8')
        msg['Subject'] = subject
        msg['From'] = self.email_config['sender']
        msg['To'] = ','.join(self.email_config['receivers'])
        
        try:
            with smtplib.SMTP(self.email_config['smtp_server'], 
                             self.email_config['smtp_port']) as server:
                server.starttls()
                server.login(self.email_config['sender'], 
                            self.email_config['password'])
                server.send_message(msg)
        except Exception as e:
            print(f"邮件报警失败: {e}")
    
    def alert(self, level, message):
        """统一报警入口"""
        if level == 'P0':
            self.send_wechat(f"[致命] {message}")
            self.send_email(f"[P0] 基差套利系统异常", message)
        elif level == 'P1':
            self.send_wechat(f"[严重] {message}")
        elif level == 'P2':
            self.send_email(f"[警告] {message}", message)
        # P3 只记录日志,不报警

alert = AlertManager()
alert.alert('P0', '策略进程意外退出,请立即检查!')

避坑指南:我曾经把报警逻辑写在了主进程里。结果主进程挂了,报警也发不出去。后来我单独起了一个监控进程,用心跳检测的方式,每隔30秒检查主进程是否活着。

三、日志管理:复盘时的救命稻草

日志这东西,平时觉得没用,真出问题的时候,它就是唯一的线索。我见过有人把日志全打在一个文件里,几百万行,查个问题得 grep 半天。

3.1 日志分级与配置

Python 自带的 logging 模块其实够用,关键是要用好:

import logging
import logging.handlers
from datetime import datetime

def setup_logger(name='basis_trading'):
    logger = logging.getLogger(name)
    logger.setLevel(logging.DEBUG)
    
    # 格式:时间 - 级别 - 模块 - 消息
    formatter = logging.Formatter(
        '%(asctime)s - %(levelname)s - %(name)s - %(message)s'
    )
    
    # 按天轮转,保留30天
    file_handler = logging.handlers.TimedRotatingFileHandler(
        f'logs/{name}_{datetime.now().strftime("%Y%m%d")}.log',
        when='midnight',
        interval=1,
        backupCount=30
    )
    file_handler.setLevel(logging.DEBUG)
    file_handler.setFormatter(formatter)
    
    # 控制台输出,只显示INFO及以上
    console_handler = logging.StreamHandler()
    console_handler.setLevel(logging.INFO)
    console_handler.setFormatter(formatter)
    
    logger.addHandler(file_handler)
    logger.addHandler(console_handler)
    
    return logger

logger = setup_logger()
logger.info("策略启动成功")
logger.warning("基差偏离阈值,请注意")
logger.error("数据源连接超时")

3.2 结构化日志

纯文本日志查起来太痛苦了。我后来改用 JSON 格式,方便用 ELK 或 Splunk 做分析:

import json

class JsonFormatter(logging.Formatter):
    def format(self, record):
        log_record = {
            'timestamp': self.formatTime(record, self.datefmt),
            'level': record.levelname,
            'logger': record.name,
            'message': record.getMessage(),
        }
        # 如果有额外字段,也加进去
        if hasattr(record, 'extra'):
            log_record.update(record.extra)
        return json.dumps(log_record, ensure_ascii=False)

# 使用示例
logger = logging.getLogger('trade')
handler = logging.StreamHandler()
handler.setFormatter(JsonFormatter())
logger.addHandler(handler)

# 记录交易信息时带上额外字段
logger.info('执行交易', extra={
    'symbol': 'BTC-USDT',
    'side': 'buy',
    'price': 45000,
    'quantity': 0.1
})

输出效果:

{"timestamp": "2024-01-15 14:30:00", "level": "INFO", "logger": "trade", "message": "执行交易", "symbol": "BTC-USDT", "side": "buy", "price": 45000, "quantity": 0.1}

你看,这样查起来多方便。直接 grep 'BTC-USDT' trade.log | jq . 就能过滤出所有 BTC 的交易记录。

四、知识体系总览

下面这张图,是我自己总结的自动化运维框架。你照着搭,基本不会漏:

基差套利自动化运维框架 定时任务调度 异常监控报警 日志管理 APScheduler Cron / Interval 触发器 任务链与依赖管理 P0-P3 分级报警 微信 / 邮件 / 短信 心跳检测与进程守护 TimedRotatingFileHandler JSON 结构化日志 ELK / Splunk 集成 目标:7×24 小时无人值守运行 故障发现 < 5分钟 · 日志可追溯 30 天 · 任务零遗漏

五、整合到一起

最后,把这三块拼起来,就是一个完整的自动化运维方案:

import time
import logging
from apscheduler.schedulers.background import BackgroundScheduler

class BasisTradingSystem:
    def __init__(self):
        self.logger = setup_logger()
        self.alert = AlertManager()
        self.scheduler = BackgroundScheduler()
        self.heartbeat_count = 0
        
    def heartbeat_check(self):
        """心跳检测,每30秒执行一次"""
        self.heartbeat_count += 1
        self.logger.debug(f"心跳正常,第{self.heartbeat_count}次")
        
    def data_fetch_job(self):
        """数据拉取任务"""
        try:
            self.logger.info("开始拉取数据")
            # 数据拉取逻辑
            self.logger.info("数据拉取完成")
        except Exception as e:
            self.logger.error(f"数据拉取失败: {e}")
            self.alert.alert('P1', f"数据拉取异常: {e}")
            
    def basis_calc_job(self):
        """基差计算任务"""
        try:
            self.logger.info("开始计算基差")
            # 基差计算逻辑
            self.logger.info("基差计算完成")
        except Exception as e:
            self.logger.error(f"基差计算失败: {e}")
            self.alert.alert('P2', f"基差计算异常: {e}")
            
    def run(self):
        """启动系统"""
        self.logger.info("基差套利系统启动")
        
        # 注册定时任务
        self.scheduler.add_job(self.heartbeat_check, 'interval', seconds=30)
        self.scheduler.add_job(self.data_fetch_job, 'cron', hour=9, minute=0)
        self.scheduler.add_job(self.basis_calc_job, 'interval', minutes=5)
        
        # 启动调度器
        self.scheduler.start()
        
        try:
            # 保持主进程运行
            while True:
                time.sleep(60)
        except KeyboardInterrupt:
            self.logger.info("收到停止信号,正在关闭...")
            self.scheduler.shutdown()
            self.logger.info("系统已安全关闭")

if __name__ == "__main__":
    system = BasisTradingSystem()
    system.run()

最后说一句:自动化运维不是一蹴而就的。我建议你先从日志管理开始,然后加报警,最后再上定时调度。每一步都跑稳了,再往前走。别贪多,稳才是王道。


公众号:蓝海资料掘金营,微信deep3321