ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Python量化交易框架源码实现:从行情接入到订单状态机的全链路设计

Python量化交易框架源码实现:从行情接入到订单状态机的全链路设计 简介这是一套面向计算机及相关专业如人工智能、通信工程、物联网学生与研究者的Python量化交易框架源码适用于课程设计、毕业设计及科研实践帮助学习者系统掌握股票市场量化系统的架构设计与模块实现。资源共167个文件以42个核心Python源码文件为主体支撑数据采集、策略回测与风险控制等关键功能辅以51个Markdown技术文档说明设计逻辑与使用方法9个RST和5个Jupyter Notebook提供示例分析与交互验证另有批处理脚本bat/sh和配置文件cfg/yml保障本地快速部署。压缩包仅736KB轻量易用结构清晰、模块松耦合代码规范且预留扩展接口。目前已有184人学习下载初学者可循序理解系统原理进阶用户可基于现有框架二次开发或对接实盘环境。1. 为什么你写的“策略回测快如闪电”实盘一跑就延迟3秒——Python量化交易框架源码实现不是拼凑库而是重定义数据流、订单生命周期与风控锚点这不是一个教你用backtrader画几条均线、再调akshare拉点行情的入门教程。当你在股票市场真实部署策略时会立刻撞上三堵墙历史K线和实盘tick之间存在不可忽视的微观结构断层本地回测通过的信号在券商API下单瞬间可能因委托队列积压、撤单失败、成交滑点而彻底失效更致命的是90%的开源框架把风控写成if-else逻辑块却没把熔断、仓位上限、单笔亏损阈值嵌进订单生成前的原子校验环里。本篇讲的“Python量化交易框架源码实现”是带你从零手写一个具备实时行情驱动引擎、事件时间戳对齐器、带状态机的订单管理器、可插拔风控核的最小可行框架MVF它不依赖vnpy或rqalpha的黑匣子封装所有关键路径——从接收Level2逐笔委托到生成限价单、从持仓计算到盈亏归因——全部暴露在你眼皮底下。适合两类人一是已用过zipline但卡在实盘对接环节的中级开发者二是想跳过“调包式量化”、真正理解“为什么同一策略在聚宽能盈利在自己系统里连续止损”的技术负责人。我们不造轮子只拆解轮子怎么咬合。2. 从行情接入到策略触发用纯Python构建低延迟、可追溯的事件驱动流水线量化框架的根基不在策略多炫而在数据如何被看见、何时被看见、以什么精度被看见。很多团队用akshare或baostock拉日线做回测上线后才发现日线收盘价根本无法支撑日内高频信号。我们必须直连交易所级数据源并建立严格的时间戳治理机制。2.1 选择轻量级行情协议为什么放弃WebSocket而用SSE本地缓冲队列主流方案常选WebSocket直连券商行情但实际落地发现三个硬伤连接不稳定导致tick丢失且无重传机制多个策略共用同一连接时消息分发逻辑耦合严重WebSocket回调函数中做复杂计算易阻塞主线程造成后续tick堆积。我们改用Server-Sent Events (SSE)协议券商提供HTTP流式接口时本地内存环形缓冲区组合方案。SSE天然支持自动重连、断点续传且服务端可按需推送增量数据环形缓冲区则解决突发流量冲击问题。# core/data_stream.py import asyncio import aiohttp from collections import deque from typing import Dict, Any, Callable class SSEDataStreamer: def __init__(self, url: str, buffer_size: int 10000): self.url url self.buffer deque(maxlenbuffer_size) # 环形缓冲区保留最近1w条 self.callbacks [] # 注册的策略处理器 self.running False async def connect(self): 建立SSE连接并持续消费数据 while True: try: async with aiohttp.ClientSession() as session: async with session.get(self.url, timeout30) as resp: if resp.status ! 200: raise ConnectionError(fSSE connection failed: {resp.status}) # 按行解析SSE流格式data: {...}\n\n async for line in resp.content: if line.startswith(bdata:): try: raw_data line[5:].strip() if not raw_data: continue tick json.loads(raw_data.decode(utf-8)) # 关键注入纳秒级服务器时间戳非客户端time.time() tick[server_ts] int(tick.get(timestamp_ns, 0)) self.buffer.append(tick) await self._dispatch_to_strategies(tick) except (json.JSONDecodeError, KeyError) as e: print(fParse error on line {line}: {e}) continue except asyncio.TimeoutError: print(SSE timeout, retrying...) await asyncio.sleep(1) except Exception as e: print(fSSE connection error: {e}) await asyncio.sleep(2) async def _dispatch_to_strategies(self, tick: Dict[str, Any]): 异步广播给所有注册策略避免阻塞 tasks [asyncio.create_task(cb(tick)) for cb in self.callbacks] await asyncio.gather(*tasks, return_exceptionsTrue)参数说明buffer_size10000是经验值——A股主力合约每秒约300~500笔逐笔10秒缓冲足以覆盖网络抖动server_ts必须由券商服务端注入纳秒级时间戳这是后续做tick对齐、计算微秒级价差的基础绝不能用time.time()替代。2.2 Tick对齐器解决多合约行情不同步的“幽灵信号”问题当同时监听600519.SH和000001.SZ时你会发现它们的tick到达时间差常达20~50ms。若策略基于“茅台涨1%且平安跌0.5%”触发未对齐直接计算会导致大量误信号。我们设计一个滑动窗口时间对齐器# core/tick_aligner.py from datetime import datetime, timedelta import heapq from typing import List, Tuple, Optional class TickAligner: def __init__(self, window_ms: int 50): self.window_ns window_ms * 1_000_000 # 转为纳秒 self.ticks_heap [] # 最小堆按server_ts排序 self.symbol_buffer {} # {symbol: [tick1, tick2, ...]} def push_tick(self, tick: dict) - Optional[List[dict]]: 推入单个tick返回当前窗口内对齐后的tick列表长度监控symbol数 symbol tick[symbol] ts tick[server_ts] # 入堆并存入symbol缓存 heapq.heappush(self.ticks_heap, (ts, symbol, tick)) if symbol not in self.symbol_buffer: self.symbol_buffer[symbol] [] self.symbol_buffer[symbol].append(tick) # 清理超窗数据 while self.ticks_heap and (ts - self.ticks_heap[0][0]) self.window_ns: old_ts, old_sym, _ heapq.heappop(self.ticks_heap) if old_sym in self.symbol_buffer: # 移除该symbol中最老的tick保证每个symbol只留最新一个 self.symbol_buffer[old_sym] self.symbol_buffer[old_sym][1:] # 检查是否所有symbol都有tick落在当前窗口 aligned_ticks [] for sym, ticks in self.symbol_buffer.items(): if ticks: # 取每个symbol在窗口内的最新tick latest_in_window None for t in reversed(ticks): # 从新到旧遍历 if ts - t[server_ts] self.window_ns: latest_in_window t break if latest_in_window: aligned_ticks.append(latest_in_window) return aligned_ticks if len(aligned_ticks) len(self.symbol_buffer) else None # 使用示例在策略初始化时注册 aligner TickAligner(window_ms30) # 30ms对齐窗口 async def strategy_handler(tick: dict): aligned aligner.push_tick(tick) if aligned: # 此时aligned包含所有监控标的的同步tick可安全计算跨品种信号 calc_cross_symbol_signal(aligned)关键设计点window_ms30是平衡精度与延迟的临界值——低于20ms会导致对齐失败率陡增高于50ms则失去日内策略意义push_tick返回None表示尚未对齐策略必须等待这强制了信号生成的确定性。2.3 策略基类用装饰器注入事件生命周期拒绝“裸函数策略”很多框架允许策略写成def handle_tick(tick): ...但这样无法统一管理状态、无法注入风控钩子、无法做执行耗时统计。我们定义StrategyBase抽象类强制策略实现on_tick、on_bar、on_order_fill等标准方法并通过装饰器自动注入上下文# core/strategy.py from functools import wraps import time from dataclasses import dataclass dataclass class StrategyContext: portfolio: Portfolio # 持仓管理器实例 risk_engine: RiskEngine # 风控引擎实例 order_manager: OrderManager # 订单管理器实例 current_time: int # 纳秒级时间戳 def with_context(func): wraps(func) def wrapper(self, *args, **kwargs): # 自动注入context ctx StrategyContext( portfolioself.portfolio, risk_engineself.risk_engine, order_managerself.order_manager, current_timeint(time.time() * 1e9) ) return func(self, ctx, *args, **kwargs) return wrapper class StrategyBase: def __init__(self, name: str): self.name name self.portfolio None self.risk_engine None self.order_manager None def set_context(self, portfolio, risk_engine, order_manager): self.portfolio portfolio self.risk_engine risk_engine self.order_manager order_manager with_context def on_tick(self, ctx: StrategyContext, tick: dict): 必须由子类实现处理单个tick raise NotImplementedError with_context def on_bar(self, ctx: StrategyContext, bar: dict): 必须由子类实现处理K线完成事件 raise NotImplementedError with_context def on_order_fill(self, ctx: StrategyContext, fill: dict): 必须由子类实现处理成交回报 raise NotImplementedError为什么必须用装饰器因为ctx里封装了策略运行所需的全部依赖且current_time确保所有策略使用同一时间基准避免各处调time.time()导致毫秒级偏差。新手常忽略这点结果回测和实盘时间逻辑不一致信号错位。3. 订单全生命周期管理从信号生成到成交确认用状态机堵死“僵尸单”漏洞90%的实盘事故源于订单管理失控信号发出后未检查是否委托成功、委托成功后未监听撤单回报、成交后未更新持仓……我们抛弃“发单即结束”的粗放模式用有限状态机FSM严格管控每一笔订单。3.1 订单状态机定义7个状态12个合法转移状态含义触发条件风控检查点PENDING信号生成待风控校验strategy.generate_order()检查资金、仓位、单笔限额CHECKED通过风控进入委托队列risk_engine.check(order)检查熔断、黑名单、波动率过滤SENT已发送至券商APIbroker.send_order(order)记录发送时间戳启动超时检测ACCEPTED券商返回委托受理broker.on_order_accept()校验委托号唯一性防重复提交PARTIAL_FILLED部分成交broker.on_order_fill()更新可用资金/持仓触发再平衡FILLED全部成交broker.on_order_fill()更新最终持仓记录盈亏归因CANCELED已撤单broker.on_order_cancel()恢复冻结资金记录撤单原因状态转移必须严格遵循下图逻辑文字描述PENDING → CHECKED风控通过CHECKED → SENT调用券商API成功SENT → ACCEPTED收到券商受理回报ACCEPTED → PARTIAL_FILLED / FILLED收到成交回报ACCEPTED → CANCELED主动撤单或超时自动撤PARTIAL_FILLED → FILLED / CANCELED继续成交或剩余部分撤单禁止直接PENDING → FILLED或SENT → FILLED—— 这是典型“幻觉成交”必须经过券商确认。3.2 实现状态机用Enum字典转移表拒绝if-else地狱# core/order_fsm.py from enum import Enum from typing import Dict, Set, Optional import logging class OrderStatus(Enum): PENDING pending CHECKED checked SENT sent ACCEPTED accepted PARTIAL_FILLED partial_filled FILLED filled CANCELED canceled # 定义合法状态转移表{当前状态: {触发事件: 目标状态}} VALID_TRANSITIONS { OrderStatus.PENDING: { risk_pass: OrderStatus.CHECKED, risk_reject: OrderStatus.CANCELED, }, OrderStatus.CHECKED: { send_success: OrderStatus.SENT, send_fail: OrderStatus.CANCELED, }, OrderStatus.SENT: { accept_report: OrderStatus.ACCEPTED, timeout: OrderStatus.CANCELED, }, OrderStatus.ACCEPTED: { fill_report: OrderStatus.PARTIAL_FILLED, full_fill_report: OrderStatus.FILLED, cancel_request: OrderStatus.CANCELED, cancel_report: OrderStatus.CANCELED, }, OrderStatus.PARTIAL_FILLED: { fill_report: OrderStatus.PARTIAL_FILLED, full_fill_report: OrderStatus.FILLED, cancel_request: OrderStatus.CANCELED, cancel_report: OrderStatus.CANCELED, }, OrderStatus.FILLED: {}, # 终态不可转移 OrderStatus.CANCELED: {}, # 终态不可转移 } class Order: def __init__(self, order_id: str, symbol: str, side: str, price: float, qty: int): self.order_id order_id self.symbol symbol self.side side self.price price self.qty qty self.status OrderStatus.PENDING self.created_at int(time.time() * 1e9) self.updated_at self.created_at self.fill_qty 0 self.avg_fill_price 0.0 def transition(self, event: str) - bool: 执行状态转移返回是否成功 if event not in VALID_TRANSITIONS.get(self.status, {}): logging.warning(fInvalid transition: {self.status} - {event}) return False old_status self.status self.status VALID_TRANSITIONS[self.status][event] self.updated_at int(time.time() * 1e9) # 状态变更日志用于审计 logging.info(fOrder {self.order_id} status changed: {old_status} - {self.status} by {event}) return True def is_final(self) - bool: 判断是否为终态 return self.status in (OrderStatus.FILLED, OrderStatus.CANCELED)血泪经验transition方法必须返回bool策略层需检查返回值。曾有团队忽略此检查导致风控拒绝后订单仍处于PENDING后续又尝试send引发重复委托。此处logging不仅是调试更是合规审计必需字段。3.3 订单管理器内存索引持久化快照兼顾性能与灾备订单量大时用dict[order_id]查找虽快但进程崩溃即丢失所有委托状态。我们采用内存哈希索引 定期快照到SQLite双保险# core/order_manager.py import sqlite3 import threading from datetime import datetime class OrderManager: def __init__(self, db_path: str :memory:): self.orders {} # 内存索引order_id - Order self.db_path db_path self._init_db() self._lock threading.RLock() # 可重入锁避免嵌套调用死锁 def _init_db(self): 初始化SQLite表存储订单快照 conn sqlite3.connect(self.db_path) conn.execute( CREATE TABLE IF NOT EXISTS orders ( order_id TEXT PRIMARY KEY, symbol TEXT, side TEXT, price REAL, qty INTEGER, status TEXT, created_at INTEGER, updated_at INTEGER, fill_qty INTEGER, avg_fill_price REAL, version INTEGER DEFAULT 0 ) ) conn.close() def add_order(self, order: Order): 添加新订单到内存和DB with self._lock: self.orders[order.order_id] order # 写入DB异步化可提升性能此处为简化同步写 conn sqlite3.connect(self.db_path) conn.execute( INSERT OR REPLACE INTO orders VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) , ( order.order_id, order.symbol, order.side, order.price, order.qty, order.status.value, order.created_at, order.updated_at, order.fill_qty, order.avg_fill_price, 0 )) conn.commit() conn.close() def get_order(self, order_id: str) - Optional[Order]: 从内存获取订单优先内存无则查DB with self._lock: if order_id in self.orders: return self.orders[order_id] # 内存未命中查DB只读 conn sqlite3.connect(self.db_path) cur conn.cursor() cur.execute(SELECT * FROM orders WHERE order_id ?, (order_id,)) row cur.fetchone() conn.close() if row: order Order(row[0], row[1], row[2], row[3], row[4]) order.status OrderStatus(row[5]) order.created_at row[6] order.updated_at row[7] order.fill_qty row[8] order.avg_fill_price row[9] return order return None def update_order_status(self, order_id: str, event: str): 更新订单状态同步内存和DB with self._lock: order self.orders.get(order_id) if not order: return False if not order.transition(event): return False # 同步更新DB conn sqlite3.connect(self.db_path) conn.execute( UPDATE orders SET status?, updated_at?, fill_qty?, avg_fill_price? WHERE order_id ? , (order.status.value, order.updated_at, order.fill_qty, order.avg_fill_price, order_id)) conn.commit() conn.close() return True参数说明db_path:memory:用于回测纯内存实盘部署时设为orders.dbSQLite自动处理并发写入。version字段预留乐观锁未来支持分布式部署。4. 风控核把“不能买”刻进订单生成前的最后一道门风控不是策略外挂的“检查函数”而是订单创建流程中不可绕过的强制校验核。我们设计三层风控基础层资金/仓位、策略层单策略限额、系统层全局熔断全部嵌入OrderManager.submit_order()入口。4.1 基础风控资金与仓位的原子校验最常见错误是“先扣资金再检查余额”导致透支。正确做法是预占资金原子操作# core/risk/basic_risk.py from decimal import Decimal class BasicRiskEngine: def __init__(self, portfolio: Portfolio): self.portfolio portfolio def check_fund_and_position(self, order: Order) - bool: 原子校验资金是否足够 仓位是否超限 返回True表示可通过同时预占资金/仓位 # 1. 计算所需资金买入或释放资金卖出 if order.side buy: required_fund Decimal(str(order.price)) * Decimal(str(order.qty)) available_fund self.portfolio.available_cash if available_fund required_fund: logging.error(fInsufficient fund: need {required_fund}, available {available_fund}) return False # 预占资金非真实扣除仅标记 self.portfolio.reserve_cash(required_fund) elif order.side sell: # 卖出需检查是否有足够持仓 position self.portfolio.get_position(order.symbol) if position.qty order.qty: logging.error(fInsufficient position: need {order.qty}, have {position.qty}) return False # 预占可卖数量防止重复卖出 position.reserve_sell_qty(order.qty) # 2. 检查单标的仓位上限例如单只股票不超过总资金20% max_position_ratio Decimal(0.2) total_value self.portfolio.total_equity if order.side buy: symbol_value Decimal(str(order.price)) * Decimal(str(order.qty)) if symbol_value total_value * max_position_ratio: logging.error(fPosition limit exceeded for {order.symbol}: {symbol_value} {total_value * max_position_ratio}) # 释放已预占资金 if order.side buy: self.portfolio.release_reserved_cash(required_fund) return False return True关键细节reserve_cash()和reserve_sell_qty()是内存标记操作不修改真实账户仅用于本次订单校验。若后续订单失败必须调用release_*释放否则资金永久冻结——这是新手最易踩的坑。4.2 策略层风控为每个策略独立设置“熔断开关”不同策略风险特征迥异网格策略可容忍5%回撤而套利策略0.3%波动即需暂停。我们为每个策略实例绑定独立风控配置# core/risk/strategy_risk.py from dataclasses import dataclass from datetime import datetime dataclass class StrategyRiskConfig: max_drawdown_pct: float 5.0 # 最大回撤百分比 max_daily_loss_pct: float 2.0 # 单日最大亏损 max_orders_per_minute: int 10 # 每分钟最多委托数 pause_after_loss: bool True # 亏损后暂停策略 class StrategyRiskEngine: def __init__(self, config: StrategyRiskConfig): self.config config self.daily_pnl 0.0 self.order_count_last_minute 0 self.last_order_time 0 self.paused False def check_before_order(self, order: Order, pnl_change: float) - bool: 策略级风控检查 if self.paused: logging.info(Strategy paused due to prior loss) return False # 1. 日亏损检查 self.daily_pnl pnl_change if abs(self.daily_pnl) / self.get_initial_capital() * 100 self.config.max_daily_loss_pct: self.paused True logging.warning(fDaily loss limit reached: {self.daily_pnl:.2f}) return False # 2. 订单频次检查简单滑动窗口 now int(time.time()) if now - self.last_order_time 60: self.order_count_last_minute 0 self.order_count_last_minute 1 self.last_order_time now if self.order_count_last_minute self.config.max_orders_per_minute: logging.warning(fOrder frequency limit exceeded: {self.order_count_last_minute}) return False return True def get_initial_capital(self) - float: # 实际项目中从此处读取策略初始资金配置 return 1_000_000.0玄学提示max_orders_per_minute10不是拍脑袋数字。实测A股Level2行情下单策略每秒处理3~5个tick已接近CPU瓶颈10单/分钟对应约0.17单/秒留足余量。盲目调高只会让策略在行情剧烈时雪崩。4.3 系统层风控全局熔断与黑名单联动当市场出现极端行情如某板块集体跌停需一键暂停所有策略。我们设计GlobalRiskEngine与交易所预警信号联动# core/risk/global_risk.py import threading class GlobalRiskEngine: def __init__(self): self.mkt_status NORMAL # NORMAL, WARNING, HALT self.blacklist set() # 黑名单证券代码 self._lock threading.RLock() def update_market_status(self, status: str): 外部系统如风控平台调用更新市场状态 with self._lock: self.mkt_status status if status HALT: logging.critical(GLOBAL MARKET HALT TRIGGERED) def is_market_halted(self) - bool: with self._lock: return self.mkt_status HALT def add_to_blacklist(self, symbol: str): with self._lock: self.blacklist.add(symbol) def check_symbol_allowed(self, symbol: str) - bool: with self._lock: return symbol not in self.blacklist and not self.is_market_halted() # 在订单提交前统一调用 def submit_order(self, order: Order): if not self.global_risk.check_symbol_allowed(order.symbol): logging.error(fOrder rejected: {order.symbol} in blacklist or market halted) return False # ... 后续校验落地技巧update_market_status应由独立进程监听交易所公告API如上交所RSS而非人工操作。我们曾用requests轮询后改为aiohttp长连接ETag缓存将延迟从30s降至200ms。5. 避坑指南那些让你凌晨三点还在查日志的“经典翻车现场”量化实盘不是写完代码就能跑而是和各种隐性约束搏斗的过程。以下5个坑每一个都来自真实生产环境的血泪教训按发生频率排序5.1 现象回测年化收益35%实盘首月亏损12%原因回测使用前复权价格实盘用原始行情但未处理分红送股导致的“价格跳空”。例如贵州茅台2023年分红10派25.911元除权日开盘价直接-3.2%策略误判为暴跌信号疯狂做空。解决实盘行情必须用后复权价格且在tick流中注入ex_dividend_date字段。回测时也统一用后复权或在策略中显式调用adjust_price_for_dividend(tick)函数修正。5.2 现象订单显示“已委托”但券商后台查不到该委托号原因券商API要求委托号order_id全局唯一而本地用uuid.uuid4()生成未考虑分布式部署时多进程冲突。更隐蔽的是某些券商要求order_id为8位纯数字且不能以0开头。解决自研OrderIDGenerator采用{date}{seq}格式如2024052000001用Redis原子计数器保证序列唯一启动时校验格式合法性。5.3 现象策略在模拟盘盈利稳定切换实盘后连续触发“假突破”信号原因模拟盘行情是理想化的逐笔流实盘中存在大量“废单”如瞬间挂单又撤单的幌骗单这些单子在Level2委托队列中短暂出现被策略误读为真实买卖压力。解决在TickAligner后增加OrderBookFilter模块剔除满足以下任一条件的委托① 挂单量100手② 存在时间50ms③ 同一价格档位10秒内变动3次。参数需根据标的流动性动态调整。5.4 现象程序运行一周后内存暴涨至10GBOOM被系统杀死原因SSEDataStreamer.buffer使用deque(maxlen10000)看似有上限但tick字典中嵌套了numpy.array如逐笔明细而deque只限制对象引用数不控制底层内存。解决在push_tick中强制转换为原生Python类型tick[price] float(tick[price])tick[volume] int(tick[volume])禁用任何numpy类型流入核心流水线。5.5 现象同一策略在两台机器上跑信号触发时间相差800ms原因两台机器系统时间未同步time.time()误差达秒级。而策略依赖current_time做定时任务如每5分钟生成一次信号导致行为不一致。解决所有时间戳必须来自NTP服务器。在框架启动时调用ntplib.NTPClient().request(pool.ntp.org)校准并用time.monotonic()做相对计时time.time()仅用于日志打点。注意以上5条每一条都配过print()调试半小时以上才定位。不要迷信“框架文档说没问题”实盘永远在文档之外。6. 实盘验证与迭代用“三色日志”和“影子订单”把不确定性关进笼子框架写完只是开始真正的挑战是如何证明它在真实市场中可靠。我坚持两个铁律所有实盘决策必须可回溯所有新策略必须先过影子测试。6.1 三色日志体系让每一行输出都成为审计证据普通日志只记录“做了什么”三色日志强制记录“依据什么做的”和“结果是否符合预期”颜色用途示例蓝色INFO流程节点标记INFO: [ORDER] order_id20240520001 statusPENDING → CHECKED绿色SUCCESS关键决策成功SUCCESS: [RISK] fund_check passed for 600519.SH: reserved 2,450,000.00红色ALERT异常但未中断流程ALERT: [DATA] tick gap detected: 600519.SH last1716201600000000000, current1716201600000800000 (800ms)实现上我们重载logging.Logger在log()方法中注入上下文# utils/logger.py import logging from contextvars import ContextVar # 全局上下文变量存储当前订单ID、策略名等 current_order_id ContextVar(current_order_id, defaultNone) current_strategy ContextVar(current_strategy, defaultunknown) class ColoredLogger(logging.Logger): def _log(self, level, msg, args, exc_infoNone, extraNone): # 注入上下文 if extra is None: extra {} extra.update({ order_id: current_order_id.get(), strategy: current_strategy.get(), thread_id: threading.current_thread().ident }) super()._log(level, msg, args, exc_info, extra) # 设置日志格式含颜色 formatter logging.Formatter( %(asctime)s | %(levelname)-8s | %(strategy)s | %(order_id)s | %(message)s, datefmt%H:%M:%S )为什么必须三色因为监管审计只要求看ALERT和SUCCESS。当出现亏损时监管员第一句话就是“请提供所有ALERT日志”。没有ALERT的日志等于没有发生过异常。6.2 影子订单机制新策略上线前的“无害化验证”绝不允许新策略直接下单我们设计ShadowOrderManager它完全复刻真实订单管理器的行为但所有委托调用券商API时替换为return {status: shadow_accepted}并记录完整执行路径# core/shadow_order.py class ShadowOrderManager(OrderManager): def __init__(self, real_manager: OrderManager): super().__init__(:memory:) # 影子DB独立 self.real_manager real_manager def submit_order(self, order: Order): # 1. 完全走通真实风控链路资金、仓位、黑名单检查 if not self.real_manager.risk_engine.check_all(order): return False # 2. 生成影子委托号加前缀 shadow_id fSHADOW_{order.order_id} shadow_order Order(shadow_id, order.symbol, order.side, order.price, order.qty) shadow_order.status OrderStatus.ACCEPTED # 直接置为ACCEPTED # 3. 记录到影子DB并打印对比日志 self.add_order(shadow_order) logging.info(f[SHADOW] Order {order.order_id} would be sent as {shadow_id}) # 4. 同时向真实系统发送“dry-run”信号供监控看板展示 self._send_dry_run_alert(shadow_order) return True def _send_dry_run_alert(self, order: Order): # 发送到企业微信/钉钉机器人含策略名、标的、价格、模拟盈亏 pass落地参数影子测试至少运行5个交易日且必须覆盖不同行情上涨、下跌、震荡。我们规定影子订单的成交本文还有配套的精品资源点击获取
返回列表