模拟交易系统架构设计:从信号扫描到自动执行 模拟交易系统架构设计从信号扫描到自动执行在量化交易开发中模拟交易Paper Trading系统是连接策略研究与实盘的关键桥梁。一个设计良好的模拟交易系统不仅要能验证策略逻辑还要能暴露实盘中可能遇到的延迟、滑点、数据异常等问题。本文分享一个完整的模拟交易系统架构涵盖信号扫描、市场分析、交易执行和通知推送四个核心模块。系统架构总览整个系统采用模块化设计各模块通过消息队列或直接调用解耦便于独立测试和替换。数据流向为行情数据源 → 信号扫描模块 → 市场分析模块 → 模拟交易引擎 → 推送通知模块。---------------- ------------------ ------------------ | 行情数据源 | -- | 信号扫描模块 | -- | 市场分析模块 | | (WebSocket/API) | | (Strategy Engine)| | (Risk Checker) | ---------------- ------------------ ----------------- | v ---------------- ------------------ ------------------ | 推送通知模块 | -- | 模拟交易引擎 | -- | 订单生成器 | | (Telegram/邮件) | | (Paper Engine) | | (Order Creator) | ---------------- ------------------ ------------------模块职责划分1. 信号扫描模块该模块负责监听行情数据运行策略逻辑生成交易信号。核心设计要点是策略与数据解耦——策略只接收标准化的K线数据不关心数据来源。# signal_scanner.py from abc import ABC, abstractmethod import pandas as pd import numpy as np from typing import Dict, List, Optional from dataclasses import dataclass, field from datetime import datetime dataclass class Signal: symbol: str direction: str # long or short strength: float # 0.0 to 1.0 timestamp: datetime price: float metadata: Dict field(default_factorydict) class BaseStrategy(ABC): 策略基类所有策略必须实现generate_signal方法 abstractmethod def generate_signal(self, data: pd.DataFrame) - Optional[Signal]: 输入标准化K线数据输出交易信号 pass class MovingAverageCrossStrategy(BaseStrategy): 双均线交叉策略示例 def __init__(self, fast_period: int 5, slow_period: int 20): self.fast_period fast_period self.slow_period slow_period def generate_signal(self, data: pd.DataFrame) - Optional[Signal]: if len(data) self.slow_period 1: return None fast_ma data[close].rolling(self.fast_period).mean() slow_ma data[close].rolling(self.slow_period).mean() # 金叉 if fast_ma.iloc[-2] slow_ma.iloc[-2] and fast_ma.iloc[-1] slow_ma.iloc[-1]: strength min(abs(fast_ma.iloc[-1] - slow_ma.iloc[-1]) / slow_ma.iloc[-1] * 100, 1.0) return Signal( symboldata[symbol].iloc[-1], directionlong, strengthstrength, timestampdatetime.now(), pricedata[close].iloc[-1] ) # 死叉 elif fast_ma.iloc[-2] slow_ma.iloc[-2] and fast_ma.iloc[-1] slow_ma.iloc[-1]: strength min(abs(fast_ma.iloc[-1] - slow_ma.iloc[-1]) / slow_ma.iloc[-1] * 100, 1.0) return Signal( symboldata[symbol].iloc[-1], directionshort, strengthstrength, timestampdatetime.now(), pricedata[close].iloc[-1] ) return None class SignalScanner: 信号扫描器管理多个策略和数据源 def __init__(self, strategies: List[BaseStrategy]): self.strategies strategies self.data_buffer {} # symbol - DataFrame def on_bar_update(self, symbol: str, bar: Dict): 接收新K线数据更新缓冲区并运行策略 if symbol not in self.data_buffer: self.data_buffer[symbol] pd.DataFrame(columns[open, high, low, close, volume]) # 添加新bar保留最近500根 new_row pd.DataFrame([bar]) self.data_buffer[symbol] pd.concat([self.data_buffer[symbol], new_row], ignore_indexTrue) self.data_buffer[symbol] self.data_buffer[symbol].tail(500) # 运行所有策略 signals [] for strategy in self.strategies: signal strategy.generate_signal(self.data_buffer[symbol]) if signal: signals.append(signal) return signals2. 市场分析模块该模块在信号生成后进行二次确认主要做风险过滤和仓位计算。我将其设计为可插拔的过滤器链每个过滤器独立负责一项检查。# market_analyzer.py from typing import List, Callable from dataclasses import dataclass import numpy as np dataclass class AnalysisResult: approved: bool position_size: float risk_score: float reason: str class MarketAnalyzer: 市场分析器执行一系列过滤器 def __init__(self, filters: List[Callable] None): self.filters filters or [] def add_filter(self, filter_func: Callable): self.filters.append(filter_func) def analyze(self, signal, market_data: dict) - AnalysisResult: 运行所有过滤器决定是否接受信号 for filter_func in self.filters: result filter_func(signal, market_data) if not result.approved: return result return AnalysisResult(approvedTrue, position_size1.0, risk_score0.0) # 过滤器示例 def volatility_filter(signal, market_data, max_volatility0.05): 波动率过滤拒绝高波动行情 symbol signal.symbol recent_bars market_data.get(symbol, []) if len(recent_bars) 20: return AnalysisResult(approvedFalse, position_size0, risk_score1.0, reason数据不足) returns np.diff([bar[close] for bar in recent_bars[-20:]]) volatility np.std(returns) if volatility max_volatility: return AnalysisResult(approvedFalse, position_size0, risk_scorevolatility, reasonf波动率过高: {volatility:.4f}) return AnalysisResult(approvedTrue, position_size1.0, risk_scorevolatility) def volume_filter(signal, market_data, min_volume1000): 成交量过滤拒绝流动性不足的品种 symbol signal.symbol recent_bars market_data.get(symbol, []) if not recent_bars: return AnalysisResult(approvedFalse, position_size0, risk_score1.0, reason无数据) avg_volume np.mean([bar[volume] for bar in recent_bars[-10:]]) if avg_volume min_volume: return AnalysisResult(approvedFalse, position_size0, risk_score1.0, reasonf成交量不足: {avg_volume:.0f}) return AnalysisResult(approvedTrue, position_size1.0, risk_score0.0)3. 模拟交易引擎这是系统的核心。模拟交易引擎需要处理订单状态机、持仓管理和资金核算。关键设计是使用撮合队列模拟真实撮合延迟而不是立即成交。# paper_trading_engine.py from enum import Enum from typing import Dict, List import asyncio from datetime import datetime from dataclasses import dataclass, field class OrderStatus(Enum): PENDING pending FILLED filled REJECTED rejected CANCELLED cancelled class OrderType(Enum): MARKET market LIMIT limit dataclass class Order: order_id: str symbol: str side: str # buy or sell quantity: float order_type: OrderType price: float 0.0 status: OrderStatus OrderStatus.PENDING created_at: datetime field(default_factorydatetime.now) filled_at: datetime None dataclass class Position: symbol: str quantity: float 0.0 avg_price: float 0.0 def update(self, side: str, price: float, quantity: float): 更新持仓 if side buy: total_cost self.avg_price * self.quantity price * quantity self.quantity quantity self.avg_price total_cost / self.quantity if self.quantity 0 else 0 else: self.quantity - quantity if self.quantity 0: self.quantity 0 self.avg_price 0 class PaperTradingEngine: 模拟交易引擎 def __init__(self, initial_capital: float 100000.0, fill_delay_ms: int 100): self.initial_capital initial_capital self.cash initial_capital self.positions: Dict[str, Position] {} self.orders: List[Order] [] self.fill_delay_ms fill_delay_ms async def place_order(self, symbol: str, side: str, quantity: float, order_type: OrderType OrderType.MARKET, price: float 0.0) - Order: 下单并模拟撮合延迟 order Order( order_idfORD{len(self.orders)1:06d}, symbolsymbol, sideside, quantityquantity, order_typeorder_type, priceprice ) self.orders.append(order) # 模拟撮合延迟 await asyncio.sleep(self.fill_delay_ms / 1000) # 市价单直接成交限价单检查价格 if order.order_type OrderType.MARKET: order.status OrderStatus.FILLED order.filled_at datetime.now() self._execute_fill(order, order.price if order.price 0 else self._get_market_price(symbol)) else: # 限价单需要外部触发价格检查 order.status OrderStatus.PENDING return order def _execute_fill(self, order: Order, fill_price: float): 执行成交更新账户 # 检查资金 cost fill_price * order.quantity if order.side buy and cost self.cash: order.status OrderStatus.REJECTED return # 更新持仓 if order.symbol not in self.positions: self.positions[order.symbol] Position(symbolorder.symbol) self.positions[order.symbol].update(order.side, fill_price, order.quantity) # 更新现金 if order.side buy: self.cash - cost else: self.cash cost order.status OrderStatus.FILLED order.filled_at datetime.now() def _get_market_price(self, symbol: str) - float: 获取市场价实际系统中从行情模块获取 # 简化实现 return 100.0 def get_portfolio_value(self, market_prices: Dict[str, float]) - float: 计算总资产 total self.cash for symbol, pos in self.positions.items(): if symbol in market_prices: total pos.quantity * market_prices[symbol] return total # 使用示例 async def main(): engine PaperTradingEngine(initial_capital50000) # 模拟下单 order await engine.place_order(BTCUSDT, buy, 0.1, OrderType.MARKET, price40000) print(fOrder {order.order_id}: {order.status}) order2 await engine.place_order(BTCUSDT, sell, 0.05, OrderType.MARKET, price41000) print(fOrder {order2.order_id}: {order2.status}) # 计算总资产 market_prices {BTCUSDT: 41000} total engine.get_portfolio_value(market_prices) print(fTotal portfolio value: {total:.2f}) if __name__ __main__: asyncio.run(main())4. 推送通知模块信号触发和订单成交后需要及时通知。这里设计一个支持多通道的通知管理器通过异步队列解耦通知发送与业务逻辑。# notification.py import asyncio from typing import Dict, List, Callable from dataclasses import dataclass from datetime import datetime dataclass class Notification: title: str message: str level: str # info, warning, critical timestamp: datetime datetime.now() class NotificationManager: 通知管理器支持多通道异步发送 def __init__(self): self.channels: List[Callable] [] self.queue asyncio.Queue() self._worker_task None def add_channel(self, channel_func: Callable): 添加通知通道如Telegram、邮件、Webhook self.channels.append(channel_func) async def start(self): 启动后台发送任务 self._worker_task asyncio.create_task(self._worker()) async def stop(self): 停止后台任务 if self._worker_task: self._worker_task.cancel() async def send(self, title: str, message: str, level: str info): 异步发送通知 notification Notification(titletitle, messagemessage, levellevel) await self.queue.put(notification) async def _worker(self): 后台消费队列并发送通知 while True: try: notification await self.queue.get() for channel in self.channels: try: await channel(notification) except Exception as e: print(fChannel error: {e}) except asyncio.CancelledError: break # 通道实现示例Telegram async def telegram_channel(notification: Notification, bot_token: str, chat_id: str): 发送Telegram通知 import httpx text f*{notification.title}*\n{notification.message}\n{notification.timestamp} url fhttps://api.telegram.org/bot{bot_token}/sendMessage async with httpx.AsyncClient() as client: await client.post(url, json{ chat_id: chat_id, text: text, parse_mode: Markdown }) # 通道实现示例Webhook async def webhook_channel(notification: Notification, webhook_url: str): 发送Webhook通知 import httpx payload { title: notification.title, message: notification.message, level: notification.level, timestamp: notification.timestamp.isoformat() } async with httpx.AsyncClient() as client: await client.post(webhook_url, jsonpayload) # 使用示例 async def notification_demo(): manager NotificationManager() # 配置通道 manager.add_channel(lambda n: telegram_channel(n, YOUR_BOT_TOKEN, YOUR_CHAT_ID)) manager.add_channel(lambda n: webhook_channel(n, http://localhost:3000/webhook)) await manager.start() # 发送通知 await manager.send( title交易信号触发, messageBTCUSDT 出现金叉信号建议做多, levelwarning ) # 等待发送完成 await asyncio.sleep(2) await manager.stop() if __name__ __main__: asyncio.run(notification_demo())数据流向与关键技术选型数据流时序行情采集使用WebSocket接收实时K线比REST轮询延迟更低信号生成每次新K线到达时触发策略计算风险评估信号产生后同步执行过滤器链订单执行通过异步队列提交订单模拟撮合延迟通知推送订单状态变化或信号产生时异步推送技术选型建议组件推荐方案理由异步框架asyncio aiohttp原生异步适合IO密集型任务数据存储SQLite开发/ TimescaleDB生产时序数据用专业数据库消息队列Redis Streams轻量级支持消费者组任务调度APScheduler支持cron表达式适合定时任务配置管理pydantic .env类型安全环境隔离关键设计要点异步边界所有IO操作网络请求、数据库必须异步化避免阻塞事件循环错误隔离每个模块独立try-except防止单点故障导致整个系统崩溃状态持久化定期将订单和持仓状态写入数据库支持系统重启恢复回放能力记录所有市场数据和交易决策便于