第二十九节:自动化运维:定时任务调度、异常监控报警、日志管理
做量化交易,尤其是基差套利,最怕什么?
怕策略跑着跑着停了,怕半夜数据源断了没人知道,怕服务器重启后脚本没起来。我见过太多团队,策略逻辑写得漂漂亮亮,最后死在运维上。
说白了,自动化运维就是给你的策略上个“保险”。今天咱们就聊聊,怎么把这套东西搭起来。
一、定时任务调度:让机器替你干活
基差套利策略,通常需要定时拉取数据、计算价差、执行交易。手动操作?不存在的。我个人习惯用 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 的交易记录。
四、知识体系总览
下面这张图,是我自己总结的自动化运维框架。你照着搭,基本不会漏:
五、整合到一起
最后,把这三块拼起来,就是一个完整的自动化运维方案:
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