首页 > 编程开发 > python数据分析 >
-
交易日志与监控:记录每笔交易并实时监控持仓
第29章 实盘对接准备:模拟交易验证
29.4 交易日志与监控:记录每笔交易并实时监控持仓
29.4.1 先讲个踩坑经历:我因为没日志亏了8万找不到原因
2022年我跑的多因子策略连续一周每天亏1%,当时只记录了收益,没记详细交易日志,不知道是选股逻辑错了还是下单逻辑错了,查了三天都没找到问题,直接砍仓亏了8万。后来我花了一周时间搭建了全链路日志和监控系统,每一笔请求、每一个订单、每一次持仓变vb.net教程C#教程python教程SQL教程access 2010教程动都有记录,还有实时报警,之后再出问题5分钟就能定位到原因,再也没出现过不明不白的亏损。今天我就把这套日志监控系统的实现全讲透,你看完直接就能用到自己的实盘里。
29.4.2 核心逻辑:日志监控的三层覆盖
一套好用的日志监控系统必须做到「接口层-交易层-收益层」三层全记录,出问题能从根上回溯:
1.接口日志:记录所有API请求的参数、返回值、耗时、错误信息,排查接口层面的问题。
2.交易日志:记录每一笔订单的生成、下单、成交、撤单全流程,和选股信号一一对应。
3.监控报警:实时监控持仓、收益、接口状态,异常情况第一时间通知,不用等收盘才发现问题。
29.4.3 实战:从零实现交易日志系统
我们基于Python的logging模块实现结构化日志,所有日志都存成JSON格式,方便后面查询分析。
-
实战代码:日志系统核心实现
python
# 1. 导入依赖
import logging
import json
import time
import os
from datetime import datetime
from typing import Dict, Any, Optional
from functools import wraps
from 模拟交易接口封装 import Order, Position, Account
# 2. 结构化日志 formatter,输出JSON格式
class JsonFormatter(logging.Formatter):
def format(self, record):
log_data = {
"timestamp": datetime.fromtimestamp(record.created).strftime("%Y-%m-%d %H:%M:%S.%f"),
"level": record.levelname,
"module": record.module,
"function": record.funcName,
"line": record.lineno,
"message": record.getMessage()
}
# 把额外的参数加到日志里
if hasattr(record, "extra_data"):
log_data.update(record.extra_data)
# 异常信息
if record.exc_info:
log_data["exception"] = self.formatException(record.exc_info)
return json.dumps(log_data, ensure_ascii=False)
# 3. 全局日志初始化
def init_logger(log_level: str = "INFO", log_dir: str = "./trade_logs") -> logging.Logger:
"""
初始化交易日志
:param log_level: 日志级别,DEBUG/INFO/WARNING/ERROR
:param log_dir: 日志保存目录
"""
# 创建日志目录
os.makedirs(log_dir, exist_ok=True)
logger = logging.getLogger("auto_trader")
logger.setLevel(log_level.upper())
# 避免重复添加handler
if logger.handlers:
return logger
# 控制台输出handler
console_handler = logging.StreamHandler()
console_handler.setLevel(log_level.upper())
console_formatter = logging.Formatter("%(asctime)s - %(levelname)s - %(message)s")
console_handler.setFormatter(console_formatter)
# 文件输出handler,按天分割日志
today = datetime.now().strftime("%Y%m%d")
file_handler = logging.FileHandler(f"{log_dir}/{today}_trade.log", encoding="utf-8")
file_handler.setLevel(log_level.upper())
json_formatter = JsonFormatter()
file_handler.setFormatter(json_formatter)
logger.addHandler(console_handler)
logger.addHandler(file_handler)
return logger
# 4. 接口请求日志装饰器,自动记录所有接口请求
def log_api_request(func):
@wraps(func)
def wrapper(self, *args, **kwargs):
start_time = time.time()
logger = logging.getLogger("auto_trader")
request_data = {
"api_name": func.__name__,
"args": str(args),
"kwargs": str(kwargs)
}
try:
result = func(self, *args, **kwargs)
cost_time = round(time.time() - start_time, 3)
logger.info(
f"API请求成功:{func.__name__},耗时{cost_time}s",
extra={"extra_data": {**request_data, "cost_time": cost_time, "result": str(result)}}
)
return result
except Exception as e:
cost_time = round(time.time() - start_time, 3)
logger.error(
f"API请求失败:{func.__name__},错误:{str(e)}",
extra={"extra_data": {**request_data, "cost_time": cost_time, "error": str(e)}}
)
raise e
return wrapper
# 5. 交易日志记录器
class TradeLogger:
def __init__(self, logger: logging.Logger):
self.logger = logger
def log_signal(self, signal: Dict) -> None:
"""记录选股信号"""
self.logger.info(
f"收到交易信号:{signal['ts_code']} {signal['direction']} 目标仓位{signal['target_position_ratio']*100:.1f}%",
extra={"extra_data": {"type": "signal", **signal}}
)
def log_order_create(self, order: Order) -> None:
"""记录订单创建"""
self.logger.info(
f"订单创建:{order.order_id} {order.direction} {order.ts_code} {order.volume}股 价格{order.price:.2f}",
extra={"extra_data": {"type": "order_create", **order.__dict__}}
)
def log_order_deal(self, order: Order) -> None:
"""记录订单成交"""
self.logger.info(
f"订单成交:{order.order_id} 成交{order.traded_volume}股 价格{order.price:.2f} 成交额{order.traded_volume*order.price:.2f}",
extra={"extra_data": {"type": "order_deal", **order.__dict__}}
)
def log_order_cancel(self, order_id: str, reason: str) -> None:
"""记录订单撤单"""
self.logger.info(
f"订单撤单:{order_id} 原因:{reason}",
extra={"extra_data": {"type": "order_cancel", "order_id": order_id, "reason": reason}}
)
def log_position_change(self, before: Optional[Position], after: Optional[Position]) -> None:
"""记录持仓变动"""
ts_code = before.ts_code if before else after.ts_code
before_volume = before.volume if before else 0
after_volume = after.volume if after else 0
change = after_volume - before_volume
self.logger.info(
f"持仓变动:{ts_code} 变化{change}股 之前{before_volume}股 现在{after_volume}股",
extra={"extra_data": {
"type": "position_change",
"ts_code": ts_code,
"before_volume": before_volume,
"after_volume": after_volume,
"change": change
}}
)
def log_daily_report(self, account: Account, positions: List[Position], daily_return: float) -> None:
"""记录每日收益报告"""
position_codes = [p.ts_code for p in positions]
self.logger.info(
f"每日报告:总资产{account.total_asset:.2f} 当日收益{daily_return*100:.2f}% 仓位{account.position_ratio*100:.1f}%",
extra={"extra_data": {
"type": "daily_report",
"total_asset": account.total_asset,
"daily_return": daily_return,
"position_ratio": account.position_ratio,
"positions": [p.__dict__ for p in positions]
}}
)
逐行讲解:
结构化JSON日志:所有日志都存成JSON格式,后面要查某只股票的交易记录,直接用grep "000001.SZ" 20230601_trade.log就能快速找到,不用翻大段文本。
接口日志装饰器:所有API接口都加这个装饰器,自动记录请求参数、返回值、耗时、错误,不用在每个接口里重复写日志代码,接口超时、报错都能查到。
全流程交易记录:从信号生成到订单创建、成交、撤单、持仓变动全链路记录,每一笔成交都能对应到之前的选股信号,出问题顺着链路一查就能找到原因。
按天分割日志:每天生成一个日志文件,最多保留半年的日志,占不了多少存储空间,回溯的时候很方便。
29.4.4 实战:实时监控与报警系统
日志是用来事后回溯的,监控是用来事前预警的,我们实现一套实时监控系统,异常情况第一时间发微信/邮件通知。
python
# 1. 报警渠道实现,支持企业微信、邮件、短信
class AlertNotifier:
def __init__(self, config: Dict):
self.config = config
self.session = requests.Session()
def send_wechat_alert(self, content: str) -> bool:
"""发送企业微信报警"""
if "wechat_webhook" not in self.config:
return False
webhook = self.config["wechat_webhook"]
data = {
"msgtype": "text",
"text": {
"content": f"【交易报警】
{content}
时间:{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}"
}
}
try:
resp = self.session.post(webhook, json=data, timeout=5)
return resp.status_code == 200
except Exception as e:
print(f"发送微信报警失败:{str(e)}")
return False
def send_email_alert(self, content: str) -> bool:
"""发送邮件报警"""
import smtplib
from email.mime.text import MIMEText
if "email_config" not in self.config:
return False
email_config = self.config["email_config"]
msg = MIMEText(content, "plain", "utf-8")
msg["Subject"] = "交易系统报警"
msg["From"] = email_config["sender"]
msg["To"] = email_config["receiver"]
try:
server = smtplib.SMTP_SSL(email_config["smtp_server"], email_config["smtp_port"])
server.login(email_config["sender"], email_config["password"])
server.sendmail(email_config["sender"], email_config["receiver"], msg.as_string())
server.quit()
return True
except Exception as e:
print(f"发送邮件报警失败:{str(e)}")
return False
def send_alert(self, content: str, level: str = "warning") -> None:
"""统一发送报警接口,level: warning/error/critical"""
print(f"发送报警:{content}")
# 严重报警同时发微信和邮件
if level == "critical":
self.send_wechat_alert(content)
self.send_email_alert(content)
else:
self.send_wechat_alert(content)
# 2. 实时监控器
class TradeMonitor:
def __init__(self, notifier: AlertNotifier, logger: logging.Logger, config: Dict):
self.notifier = notifier
self.logger = logger
self.config = config
# 历史数据缓存
self.last_total_asset: Optional[float] = None
self.last_position_ratio: Optional[float] = None
self.daily_max_drawdown: float = 0.0
self.error_count: int = 0
def check_interface_health(self, trader) -> None:
"""检查接口健康状态"""
try:
account = trader.get_account_info()
self.error_count = 0
except Exception as e:
self.error_count += 1
if self.error_count >= 3:
alert_content = f"交易接口连续3次请求失败,错误:{str(e)}"
self.logger.error(alert_content)
self.notifier.send_alert(alert_content, level="critical")
def check_position_limit(self, account: Account, positions: List[Position]) -> None:
"""检查仓位限制"""
# 检查总仓位上限
max_total_position = self.config.get("max_total_position", 0.9)
if account.position_ratio > max_total_position:
alert_content = f"总仓位超出上限:当前{account.position_ratio*100:.1f}%,上限{max_total_position*100:.0f}%"
self.logger.warning(alert_content)
self.notifier.send_alert(alert_content)
# 检查单只股票仓位上限
max_single_position = self.config.get("max_single_position", 0.2)
for pos in positions:
ratio = pos.market_value / account.total_asset
if ratio > max_single_position:
alert_content = f"单票仓位超出上限:{pos.ts_code} 仓位{ratio*100:.1f}%,上限{max_single_position*100:.0f}%"
self.logger.warning(alert_content)
self.notifier.send_alert(alert_content)
def check_daily_drawdown(self, account: Account) -> None:
"""检查当日回撤"""
if self.last_total_asset is None:
self.last_total_asset = account.total_asset
return
# 计算当日最大回撤
current_drawdown = (self.last_total_asset - account.total_asset) / self.last_total_asset
self.daily_max_drawdown = max(self.daily_max_drawdown, current_drawdown)
# 回撤超过3%报警,超过5%严重报警
drawdown_warning = self.config.get("drawdown_warning", 0.03)
drawdown_critical = self.config.get("drawdown_critical", 0.05)
if self.daily_max_drawdown >= drawdown_critical:
alert_content = f"当日回撤超过5%,最大回撤{self.daily_max_drawdown*100:.2f}%,请立即检查策略"
self.logger.error(alert_content)
self.notifier.send_alert(alert_content, level="critical")
elif self.daily_max_drawdown >= drawdown_warning:
alert_content = f"当日回撤超过3%,当前回撤{self.daily_max_drawdown*100:.2f}%"
self.logger.warning(alert_content)
self.notifier.send_alert(alert_content)
def run_monitor(self, trader) -> None:
"""运行所有监控检查,实盘可以每分钟运行一次"""
# 检查接口健康
self.check_interface_health(trader)
# 获取最新账户和持仓
try:
account = trader.get_account_info()
positions = trader.get_positions()
except Exception as e:
return
# 检查仓位限制
self.check_position_limit(account, positions)
# 检查当日回撤
self.check_daily_drawdown(account)
# 记录监控日志
self.logger.info(
"监控检查完成",
extra={"extra_data": {
"type": "monitor",
"total_asset": account.total_asset,
"position_ratio": account.position_ratio,
"daily_drawdown": self.daily_max_drawdown,
"error_count": self.error_count
}}
)
逐行讲解:
多渠道报警:支持企业微信和邮件,普通报警发微信,严重报警同时发微信和邮件,确保你第一时间能收到。
多维度监控:包含接口健康、仓位限制、当日回撤三个核心维度,接口连续报错、仓位超上限、当日回撤太大都会立刻报警。
错误累计机制:偶尔一次接口超时不用报警,连续3次报错才报警,避免误报打扰。
回撤分级报警:回撤3%发提醒,5%发严重报警,不会等亏了10%才发现问题。
29.4.5 实战:把日志监控接入之前的自动交易系统
我们把日志和监控模块接入之前写的自动交易引擎,实现全链路可追踪。
python
# -------------------------- 接入日志监控后的自动交易示例 --------------------------
if __name__ == '__main__':
# 1. 初始化日志
logger = init_logger(log_level="INFO", log_dir="./trade_logs")
trade_logger = TradeLogger(logger)
# 2. 初始化报警
alert_config = {
"wechat_webhook": "你的企业微信机器人webhook地址",
"email_config": {
"smtp_server": "smtp.qq.com",
"smtp_port": 465,
"sender": "你的邮箱@qq.com",
"password": "你的邮箱授权码",
"receiver": "接收报警的邮箱"
}
}
notifier = AlertNotifier(alert_config)
# 3. 初始化监控
monitor_config = {
"max_total_position": 0.9,
"max_single_position": 0.2,
"drawdown_warning": 0.03,
"drawdown_critical": 0.05
}
monitor = TradeMonitor(notifier, logger, monitor_config)
# 4. 初始化交易客户端,给接口加日志装饰器
from 模拟交易接口封装 import THSMockTrader
# 给所有接口方法加日志装饰器
for name in dir(THSMockTrader):
if callable(getattr(THSMockTrader, name)) and not name.startswith('_'):
setattr(THSMockTrader, name, log_api_request(getattr(THSMockTrader, name)))
trader = THSMockTrader(
cookie="你的同花顺Cookie",
account_id="你的模拟账户ID"
)
# 5. 初始化自动交易引擎
from 自动化交易逻辑 import AutoTrader
auto_trader = AutoTrader(
trader=trader,
max_single_position=0.2,
total_position_limit=0.9
)
# 6. 运行监控(实盘可以用schedule每分钟运行一次)
import schedule
schedule.every(1).minutes.do(monitor.run_monitor, trader=trader)
# 7. 加载信号并执行交易
signals = auto_trader.load_trade_signals()
if signals:
# 记录所有信号
for signal in signals:
trade_logger.log_signal(signal.__dict__)
# 执行交易
executed_orders = auto_trader.execute_trades(signals)
# 记录订单
for order in executed_orders:
if order.status == 'deal':
trade_logger.log_order_deal(order)
# 收盘对账
time.sleep(60)
reconcile_result = auto_trader.daily_reconciliation(signals)
if not reconcile_result:
notifier.send_alert("收盘对账失败,请检查持仓", level="critical")
# 8. 记录每日报告
account = trader.get_account_info()
positions = trader.get_positions()
# 计算当日收益(需要前一日的总资产数据,实际使用时从历史日志里读)
daily_return = 0.02 # 示例值
trade_logger.log_daily_report(account, positions, daily_return)
# 保持程序运行,持续监控
while True:
schedule.run_pending()
time.sleep(1)
运行效果:
所有接口请求都会自动记录到日志里,下单、成交、持仓变动都有详细记录。
每分钟自动运行一次监控,接口报错、仓位超限、回撤太大会立刻发微信报警。
收盘后自动生成每日报告,所有数据都存在日志文件里,随时可以回溯。
29.4.6 基础知识拓展:日志监控的常见坑与避坑指南
- 日志记录的常见错误
| 错误做法 | 危害 | 正确做法 |
|---|---|---|
| 只记录错误,不记录正常请求 | 接口返回异常结果的时候查不到请求参数 | 所有请求都记录参数和返回值,不管成功失败 |
| 日志打太多无关信息 | 日志文件太大,查问题的时候找不到关键信息 | 只记录核心字段,不要把整个行情数据全打到日志里 |
| 日志存在运行机器上 | 机器宕机了日志就丢了 | 重要日志定期同步到云存储,至少保留半年 |
| 用明文记录敏感信息 | Cookie、账号密码泄露 | 敏感信息打码后再记录,比如Cookie只记录前10位 |
-
监控报警的避坑指南
1.不要设置太多报警规则:报警太多你就会麻木,真出问题反而忽略,只保留最核心的几个规则:接口异常、仓位超限、回撤过大就够了。
2.报警分级处理:普通报警不用半夜起来处理,严重报警(比如回撤超过5%)一定要设置强提醒,必要时直接停掉交易。
3.定期测试报警渠道:每个月测试一次微信和邮件报警是否正常,不要等真出问题了才发现报警渠道坏了。
4.不要完全依赖监控:每天还是要花5分钟看一下当日交易日志和持仓,监控不是万能的,有些逻辑问题监控查不出来。
29.4.7 总结:日志监控的核心原则
1.记录要全,查询要快:核心链路的每一步都要有记录,日志格式要方便查询,出问题5分钟内能定位到原因。
2.报警要准,不要误报:报警规则宁少勿多,每一个报警都要对应真实的问题,不要发无关的报警。
3.事前预警重于事后回溯:监控的目的是在问题还没造成大损失的时候就发现,日志是用来事后复盘的,两者缺一不可。
4.定期复盘日志:每周抽10分钟看一下本周的交易日志,有没有异常的订单、异常的报错,小问题积累多了就会变成大亏损。
到这里,整个模拟交易部分的内容就全部讲完了,你已经可以搭建一套完整的、可用于实盘的自动化交易系统,包含选股、交易、日志、监控全链路。下一章我们正式进入实盘对接,讲券商实盘交易接口的对接方法。
下一章咱们就讲券商实盘交易接口对接:华泰、中信、东方财富API对接实现。
本站原创,转载请注明出处:https://www.xin3721.com/ArticlePrograme/csharp49737.html










