对面向中低频多策略交易的端到端量化系统架构 各层模块的深入技术设计,涵盖数据结构、接口契约、核心算法与实现要点,面向研发实现。
目标读者:研发工程师
前置文档:架构设计文档
具体到量化策略调优操作,估计还得好几篇。。微信排版好头疼,先看看点赞多不多再看还写不写
目录 Contents
.数据层详细设计
.策略层详细设计
.风控层详细设计
.执行层详细设计
.跨层基础设施设计
01
数据层详细设计
数据层是系统的地基,由四个子模块构成:采集适配器、清洗流水线、存储引擎、特征计算引擎。数据自左向右流经四个子模块,最终向策略层提供统一的行情与因子查询接口。

1.1 采集适配器(Ingestion Adapters)
适配器模式设计
每个数据源对应一个 Adapter,继承统一的 DataSourceBase 基类,对外暴露标准化回调接口,对内处理协议差异。新增数据源只需新增 Adapter 类,不动核心采集逻辑。
from abc import ABC, abstractmethodfrom dataclasses import dataclassfrom typing import Callableimport time# ---- 标准化数据结构 ----@dataclassclassTick: symbol: str price: float volume: float bid_price: float ask_price: float timestamp: float# epoch seconds, UTC source: str@dataclassclassBar: symbol: str open: float high: float low: float close: float volume: float timestamp: float freq: str# "1m", "5m", "1d"# ---- 适配器基类 ----classDataSourceBase(ABC):def__init__(self, name: str): self.name = name self._on_tick: Callable = None self._on_bar: Callable = None self._connected = Falsedefset_callbacks(self, on_tick, on_bar): self._on_tick = on_tick self._on_bar = on_bar@abstractmethoddefconnect(self): ...@abstractmethoddefsubscribe(self, symbols: list[str]): ...@abstractmethoddefdisconnect(self): ...# 连接管理:断线重连def_reconnect_with_backoff(self, max_retries=5):for i inrange(max_retries):try: self.connect() self._connected = TruereturnexceptConnectionError: time.sleep(min(2**i, 30)) # 指数退避raiseConnectionError(f"{self.name} 重连失败")# ---- CTP 期货适配器示例 ----classCTPAdapter(DataSourceBase):def__init__(self, config):super().__init__("CTP") self.config = configdefconnect(self):# CTP 登录、认证、结算单确认 self._api = CtpApi(self.config) self._api.login()defsubscribe(self, symbols): self._api.subscribe_market_data(symbols)def_on_ctp_tick(self, raw_tick):# 将 CTP 原始格式转为标准 Tick tick = Tick( symbol=raw_tick.InstrumentID, price=raw_tick.LastPrice, volume=raw_tick.Volume, bid_price=raw_tick.BidPrice1, ask_price=raw_tick.AskPrice1, timestamp=raw_tick.UpdateTime, source="CTP" )if self._on_tick: self._on_tick(tick)连接健壮性设计
心跳检测
每 5 秒发送心跳包,3 次未响应判定断线,触发重连。
指数退避重连
重连间隔 1s → 2s → 4s → 8s → 30s,避免雪崩。
数据缺口检测
每分钟检查最近 tick 时间戳,超 10 秒无数据告警。
优雅降级
主源断线自动切备用源,恢复后切回,切源事件记录。
1.2 清洗流水线(Cleaning Pipeline)
清洗流水线采用责任链模式,原始数据依次经过校验、复权、对齐、去重四个处理节点。每个节点可独立开关与配置,支持热更新。

classCleaningPipeline:def__init__(self): self.stages: list[PipelineStage] = [Validator(), # 数据校验Adjuster(), # 复权处理Aligner(), # 时间对齐Deduper(), # 去重 ]defprocess(self, raw_data) -> dict | None: data = raw_datafor stage in self.stages: data = stage.process(data)if data isNone:# 当前节点拒绝,进入隔离表 self._quarantine(raw_data, stage.name)returnNonereturn dataclassValidator(PipelineStage):"""数据校验:价格/量合法性、涨跌停检测"""defprocess(self, tick):if tick.price <= 0 or tick.volume < 0:returnNone# 涨跌停检测(需当日参考价)if self._is_price_limit(tick):returnNonereturn tickclassAdjuster(PipelineStage):"""复权处理:按因子表调整历史价格"""defprocess(self, tick): factor = self._get_adjust_factor(tick.symbol, tick.timestamp)returnreplace(tick, price=tick.price * factor)关键清洗规则
节点 | 规则 | 异常处理 |
Validator | 价格 > 0、量 ≥ 0、时间戳递增、涨跌停检测 | 拒绝 → 隔离表 |
Adjuster | 前复权(回测)/ 后复权(实盘),因子按日期查找 | 无因子 → 标记未复权 |
Aligner | 统一为 UTC 时间戳,跨市场按交易日历对齐 | 非交易日数据 → 丢弃 |
Deduper | 同一 symbol+timestamp 取最新,量增量 > 0 才更新 | 完全重复 → 丢弃 |
1.3 存储引擎设计
时序库表结构(TDengine)
TDengine 按"超级表 + 子表"模型组织:每个标的类型(如股票、期货)一张超级表,每个具体标的一张子表。这种模型既支持按标的精准查询,又支持跨标的聚合。
-- 超级表:股票 Tick 数据CREATE STABLE stock_ticks ( ts TIMESTAMP, price FLOAT, volume BIGINT, bid_price FLOAT, ask_price FLOAT, turnover DOUBLE) TAGS ( symbol BINARY(16), exchange BINARY(8), sector BINARY(32));-- 子表自动建表(按 symbol)INSERTINTO s600000 USING stock_ticks TAGS ('600000', 'SSE', '银行')VALUES ('2026-07-28 09:30:00.000', 12.45, 1000, 12.44, 12.46, 12450);-- 超级表:K线数据(按频率分表)CREATE STABLE bars_1m ( ts TIMESTAMP, open FLOAT, high FLOAT, low FLOAT, close FLOAT, volume BIGINT, turnover DOUBLE) TAGS (symbol BINARY(16), exchange BINARY(8));冷热分离与分区策略
·热数据:最近 3 个月 tick 数据留在 TDengine 内存缓存,毫秒级查询。
·温数据:3 个月至 2 年数据在 TDengine 磁盘层,10ms 级查询。
·冷数据:2 年以上数据归档为 Parquet 列式文件,按月分目录,批量分析时加载。
·分区:按月自动建分区,旧分区设为只读,定期迁移至冷存储。
Redis 热缓存设计
# 最新行情快照 —— Hash 结构,O(1) 读取HSET tick:snapshot:600000 price 12.45 volume 1000 ts 1722148200 bid 12.44 ask 12.46# 持仓快照 —— 策略层高频读取HSET position:strategy_001 600000 1000 600001 -500 cash 1000000# 因子缓存 —— 最近 N 天的因子值,TTL 自动过期SET factor:momentum_20:600000:20260728 0.0235 EX 864001.4 特征计算引擎
特征引擎将原始数据转化为因子值,核心要求是回测与实盘共享同一套计算逻辑。引擎支持批量计算(回测)与增量计算(实盘)两种模式,但因子定义代码完全相同。
from abc import ABC, abstractmethodimport numpy as npclassFactorBase(ABC):"""因子基类:回测与实盘共享""" name: str# 因子名称,如 "momentum_20" lookback: int# 回看周期,如 20 天 freq: str = "1d"# 计算频率@abstractmethoddefcompute(self, data: np.ndarray) -> float:"""输入历史数据数组,输出因子值""" ...classMomentumFactor(FactorBase):"""20 日动量因子:过去 20 日收益率""" name = "momentum_20" lookback = 20defcompute(self, closes: np.ndarray) -> float:iflen(closes) < self.lookback + 1:returnnp.nanreturn (closes[-1] - closes[-self.lookback]) / closes[-self.lookback]classFeatureEngine:def__init__(self, factors: list[FactorBase]): self.factors = factors self._cache = {} # 增量计算的中间状态defcompute_batch(self, symbol, hist_data):"""批量模式(回测):一次性计算全部因子"""return {f.name: f.compute(hist_data[f.freq])for f in self.factors}defcompute_incremental(self, symbol, new_bar):"""增量模式(实盘):仅用新数据更新因子""" window = self._update_cache(symbol, new_bar)return {f.name: f.compute(window[f.freq])for f in self.factors}CONSISTENCY GUARANTEE
FactorBase.compute() 是回测与实盘的唯一因子定义入口。批量模式将全部历史数据传入,增量模式维护滚动窗口后传入相同数据切片。无论哪种模式,同一时刻同一数据输入必定产出相同因子值。这是回测可信的技术保障。
1.5 对外接口定义
数据查询接口
classDataClient:"""策略层统一数据访问接口"""defget_bar(self, symbol, freq, count) -> list[Bar]:"""获取最近 N 根 K线"""defget_tick(self, symbol) -> Tick:"""获取最新 tick"""defget_history(self, symbol, start, end, freq) -> pd.DataFrame:"""获取历史 K线序列(回测用)"""classFeatureClient:"""因子查询接口,底层透明路由至缓存/库/文件"""defget(self, symbol, factor_name, date) -> float:"""查询单个因子值"""defget_panel(self, factor_name, date, universe) -> pd.Series:"""查询截面因子值(全标的)"""defget_history(self, symbol, factor_name, start, end) -> pd.Series:"""查询因子历史序列"""02
策略层详细设计
策略层由五个子模块构成:策略基类与生命周期、信号生成框架、Alpha 模型训练流水线、组合优化器、多策略调度引擎。数据从特征引擎流入,决策以"目标持仓"形式流出。

2.1 策略基类与生命周期
所有策略继承 StrategyBase,通过回调方法接收事件。回测引擎与实盘引擎均通过调用这些回调驱动策略,策略代码零改动切换环境。
from abc import ABCfrom dataclasses import dataclassfrom typing importDict, Optional@dataclassclassContext:"""策略运行上下文,引擎注入""" mode: str# "backtest" | "live" data_client: DataClient feature_client: FeatureClient portfolio: Portfolio# 当前持仓快照 clock: Clock# 当前时间classStrategyBase(ABC):"""策略基类 —— 回测与实盘的统一接口"""def__init__(self, name: str, params: dict): self.name = name self.params = params self.target_portfolio: Dict[str, float] = {} self.ctx: Optional[Context] = Nonedefon_init(self, ctx: Context):"""初始化:引擎启动时调用一次""" self.ctx = ctxdefon_bar(self, bar: Bar):"""K线回调:策略主逻辑入口 回测:引擎逐条回放历史 bar 实盘:实时行情聚合后触发"""raiseNotImplementeddefon_tick(self, tick: Tick):"""Tick回调:高频策略使用"""passdefon_fill(self, fill: Fill):"""成交回报回调:更新内部状态"""passdefon_stop(self):"""停止回调:清理资源"""pass# 策略不直接下单,仅设置目标持仓defset_target(self, weights: Dict[str, float]): self.target_portfolio = weightsDESIGN RULE
策略类只产出目标持仓,不创建订单。订单由执行层根据"目标 - 当前"差异生成。这种分离使策略无需关心当前持仓细节和交易所协议,也使回测可以直接对比目标持仓与理论持仓,隔离执行误差。
2.2 信号生成框架
信号生成将因子值转化为多空信号。系统支持三类信号源并行运行,通过统一的 SignalGenerator 接口聚合。
classSignalGenerator(ABC):"""信号生成器基类"""@abstractmethoddefgenerate(self, ctx: Context, universe: list[str]) -> pd.Series:"""返回每个标的的信号分数(正=多,负=空)""" ...classRuleSignal(SignalGenerator):"""规则型信号:均线交叉"""def__init__(self, fast=5, slow=20): self.fast, self.slow = fast, slowdefgenerate(self, ctx, universe): scores = {}for sym in universe: bars = ctx.data_client.get_bar(sym, "1d", self.slow + 1) ma_fast = np.mean([b.close for b in bars[-self.fast:]]) ma_slow = np.mean([b.close for b in bars[-self.slow:]]) scores[sym] = 1.0 if ma_fast > ma_slow else -1.0return pd.Series(scores)classMLSignal(SignalGenerator):"""机器学习信号:模型预测打分"""def__init__(self, model, feature_names): self.model = model self.feature_names = feature_namesdefgenerate(self, ctx, universe): features = self._build_features(ctx, universe)return pd.Series( self.model.predict(features[self.feature_names]), index=features.index )信号聚合策略
聚合方法 | 公式 | 适用场景 |
等权平均 |
| 各信号源质量相近 |
IC 加权 |
| 有历史 IC 统计 |
波动率倒数 |
| 各信号波动差异大 |
投票制 |
| 规则型信号离散 |
2.3 Alpha 模型训练流水线
Alpha 模型将信号转化为预期收益预测。训练流水线包含特征工程、数据分割、模型训练、评估、版本管理五个阶段。

classAlphaModel:def__init__(self, config): self.config = config self.model = None self.version = Nonedeftrain(self, factor_panel: pd.DataFrame, forward_returns: pd.Series):"""训练流程:严格时间分割,杜绝未来函数"""# 1. 时间序列分割(不可随机打乱) train_end = "2024-12-31" val_end = "2025-06-30" train_mask = factor_panel.index <= train_end val_mask = (factor_panel.index > train_end) & (factor_panel.index <= val_end) test_mask = factor_panel.index > val_end# 2. 训练 self.model = lightgbm.train( params=self.config.params, train_set=lightgbm.Dataset( factor_panel[train_mask], forward_returns[train_mask]), valid_sets=[lightgbm.Dataset( factor_panel[val_mask], forward_returns[val_mask])], num_boost_round=500, callbacks=[lightgbm.early_stopping(50)] )# 3. 评估 metrics = self._evaluate(factor_panel[test_mask], forward_returns[test_mask])assert metrics["ic"] > 0.03, "IC 过低,拒绝上线"# 4. 版本管理 self.version = ModelRegistry.register( model=self.model, train_data_hash=hash(factor_panel[train_mask].tobytes()), feature_names=list(factor_panel.columns), metrics=metrics, created_at=datetime.now() )defpredict(self, features: pd.DataFrame) -> pd.Series:"""线上推理:输出横截面预期收益排序分""" raw_pred = self.model.predict(features[self.config.feature_names])# 转为截面排序分(0~1),消除绝对值偏差return pd.Series(raw_pred, index=features.index).rank(pct=True)模型评估指标
指标 | 计算方式 | 合格阈值 | 含义 |
IC(信息系数) | corr(预测, 实际收益) | > 0.03 | 预测方向准确度 |
ICIR | mean(IC) / std(IC) | > 0.5 | IC 稳定性 |
Rank IC | spearman corr | > 0.05 | 排序预测力 |
换手率 | Σ|Δweight| / 2 | < 30%/日 | 过度交易检测 |
超额回撤 | max(peak - trough) | < 10% | 极端亏损控制 |
2.4 组合优化器设计
组合优化器将 Alpha 预测转化为目标持仓权重,在最大化预期收益的同时满足各类约束。核心是一个约束优化问题。
import cvxpy as cpclassPortfolioOptimizer:"""均值-方差优化器,支持行业中性、集中度约束"""def__init__(self, config): self.max_weight = config.get("max_weight", 0.05) self.max_turnover = config.get("max_turnover", 0.30) self.industry_neutral = config.get("industry_neutral", True)defoptimize(self, alpha: pd.Series, cov_matrix: pd.DataFrame, current_weights: pd.Series, industry_map: pd.Series) -> pd.Series:"""求解最优权重""" n = len(alpha) w = cp.Variable(n)# 目标函数:最大化 alpha - 风险惩罚 objective = cp.Maximize( alpha.values @ w - 0.5 * cp.quad_form(w, cov_matrix.values) ) constraints = [ cp.sum(w) == 1, # 满仓 w >= -self.max_weight, # 个股下限 w <= self.max_weight, # 个股上限 cp.abs(w - current_weights.values) <= self.max_turnover / 2, # 换手率约束 ]# 行业中性约束if self.industry_neutral:for ind in industry_map.unique(): mask = (industry_map == ind).values constraints.append(cp.sum(w[mask]) ==sum(current_weights[mask])) prob = cp.Problem(objective, constraints) prob.solve(solver=cp.OSQP)return pd.Series(w.value, index=alpha.index)OPTIMIZER NOTE
均值-方差优化对输入协方差矩阵极其敏感,容易过拟合。生产环境建议使用收缩估计(Shrinkage Estimator)降低协方差噪声,或直接用风险平价 / 等权作为稳健基线,仅在 Alpha 预测置信度高时启用优化器。
2.5 多策略调度引擎
调度引擎管理多个策略实例的生命周期,分配资金,聚合各策略的目标持仓为最终的组合目标。
classStrategyEngine:"""多策略调度引擎"""def__init__(self): self.strategies: Dict[str, StrategyBase] = {} self.allocations: Dict[str, float] = {} # 各策略资金占比defadd_strategy(self, strategy: StrategyBase, allocation: float): self.strategies[strategy.name] = strategy self.allocations[strategy.name] = allocationdefon_bar(self, bar: Bar):"""逐策略回调,聚合目标持仓""" final_portfolio = {}for name, strategy in self.strategies.items(): strategy.on_bar(bar) alloc = self.allocations[name]# 各策略目标持仓按资金占比加权for sym, weight in strategy.target_portfolio.items(): final_portfolio[sym] = final_portfolio.get(sym, 0) + weight * alloc# 输出最终目标持仓 → 风控层 self._emit_target_portfolio(final_portfolio)资金分配方法
固定分配
各策略按预设比例分配,简单但未考虑绩效差异。
波动率倒数
按策略历史波动率倒数分配,风险均衡。
凯利公式
按预期收益与方差最优分配,理论最优但对参数敏感。
动态调整
按滚动 Sharpe 动态调整,绩效好的加配、差的减配。
03
风控层详细设计
风控层采用规则引擎 + 三道防线架构。规则引擎以 DSL 定义风控规则,支持热更新与审计留痕;三道防线分别覆盖下单前、持仓校验、运行时监控,确保任何一笔订单在生命周期的每个阶段都被检查。

3.1 规则引擎设计
规则引擎以 YAML DSL 定义风控规则,非开发人员也可配置。规则由条件(Condition)和动作(Action)组成,引擎解析后编译为可执行函数。
规则 DSL 示例
# 风控规则定义文件 risk_rules.yamlrules:# 事前风控:单笔委托上限 - id: pre_order_max_qty stage: pre_trade condition: "order.qty > params.max_qty_per_order" action: reject message: "单笔委托量 {order.qty} 超过上限 {params.max_qty_per_order}"# 事前风控:涨跌停保护 - id: pre_price_limit stage: pre_trade condition: "order.price > tick.upper_limit or order.price < tick.lower_limit" action: reject message: "委托价超出涨跌停板"# 持仓限制:个股集中度 - id: pos_single_concentration stage: position_check condition: "portfolio.weight(symbol) > params.max_single_weight" action: reduce params: max_single_weight: 0.05 message: "个股 {symbol} 权重超限"# 事中风控:日内回撤熔断 - id: rt_drawdown_kill stage: real_time condition: "portfolio.drawdown > params.kill_drawdown" action: kill_switch params: kill_drawdown: 0.04 message: "日内回撤 {drawdown:.2%} 触发硬熔断"import yamlimport astclassRuleEngine:"""规则引擎:解析 DSL,编译为可执行规则"""def__init__(self, rules_file: str): self.rules = self._load_rules(rules_file)def_load_rules(self, path) -> list[Rule]:withopen(path) as f: raw = yaml.safe_load(f)return [self._compile(r) for r in raw["rules"]]def_compile(self, raw: dict) -> Rule:# 将条件字符串编译为 AST,安全求值(禁止函数调用) condition_ast = ast.parse(raw["condition"], mode="eval")returnRule( id=raw["id"], stage=raw["stage"], condition=condition_ast, action=raw["action"], params=raw.get("params", {}), message=raw["message"] )defcheck(self, stage: str, context: dict) -> RiskResult:"""执行指定阶段的全部规则"""for rule in self.rules:if rule.stage != stage:continueif self._eval_condition(rule.condition, context):returnRiskResult( passed=False, action=rule.action, rule_id=rule.id, message=rule.message.format(**context) )returnRiskResult(passed=True)@dataclassclassRiskResult: passed: bool action: str = ""# "reject" | "reduce" | "kill_switch" rule_id: str = "" message: str = ""
3.2 三道防线实现
事前风控(Pre-trade Check)
事前风控在订单创建后、发送交易所前执行。每笔订单必须通过全部事前规则才可放行,任何一条规则拒绝则订单进入拒单队列。
classPreTradeChecker:def__init__(self, rule_engine: RuleEngine, account: Account): self.engine = rule_engine self.account = accountdefcheck(self, order: Order) -> RiskResult: context = {"order": order,"account": self.account,"tick": self._get_latest_tick(order.symbol),"params": self._get_params(order.symbol), } result = self.engine.check("pre_trade", context)ifnot result.passed: self._log_rejection(order, result)return result持仓限制(Position Check)
持仓限制在目标持仓变化时执行,校验组合层面的敞口。它与事前风控的区别:事前风控针对单笔订单,持仓限制针对组合状态。
classPositionChecker:def__init__(self, rule_engine: RuleEngine, portfolio: Portfolio): self.engine = rule_engine self.portfolio = portfoliodefcheck(self, target_weights: dict) -> RiskResult:# 模拟目标持仓应用后的组合状态 projected = self.portfolio.project(target_weights) context = {"portfolio": projected,"params": self._get_params(), } result = self.engine.check("position_check", context)if result.action == "reduce":# 自动缩减超限标的的权重 adjusted = self._reduce_positions(target_weights, result)returnRiskResult(passed=True, action="adjusted", message=f"已自动减仓: {result.message}")return result事中风控(Real-time Monitor)
事中风控以独立线程持续运行,不依赖订单事件触发。它周期性检查组合状态,触发时执行硬熔断或软告警。
import threadingclassRealtimeMonitor:def__init__(self, rule_engine, portfolio, interval=1.0): self.engine = rule_engine self.portfolio = portfolio self.interval = interval # 检查周期(秒) self._running = False self._kill_active = Falsedefstart(self): self._running = True thread = threading.Thread(target=self._loop, daemon=True) thread.start()def_loop(self):while self._running:ifnot self._kill_active: self._check() time.sleep(self.interval)def_check(self): context = {"portfolio": self.portfolio,"params": self._get_params(),"drawdown": self.portfolio.intraday_drawdown(),"volatility": self.portfolio.realized_vol(), } result = self.engine.check("real_time", context)ifnot result.passed:if result.action == "kill_switch": self._execute_kill_switch(result)else: self._send_alert(result)def_execute_kill_switch(self, result):"""硬熔断:全部平仓 + 暂停交易""" self._kill_active = True self._send_alert(result, level="DANGER")# 向执行层发送全平指令 self.portfolio.flatten_all(reason="kill_switch")# 阻止策略层产出新信号 self.portfolio.lock(reason=result.message)
3.3 熔断与降级机制
软告警(WARNING)
触发条件较宽,仅通知人工。策略继续运行但暂停开新仓,仅允许减仓。
硬熔断(DANGER)
触发条件严格,自动全平+锁定。需人工解除后才能恢复交易。
策略降级
连续亏损时缩减该策略资金分配,而非直接停止。降级可逆,绩效恢复后自动升级。
市场档位
按市场波动状态切换风控档位:正常/波动/极端,极端档位自动收紧所有阈值。

SAFETY PRINCIPLE
风控层的硬熔断一旦触发,必须由人工显式解除,不可自动恢复。这是防止系统在异常状态下反复触发-恢复-触发导致更大损失的最后一道保险。解除操作需记录操作人、时间、原因。
04
执行层详细设计
执行层将目标持仓转化为实际订单并高效成交,由四个子模块构成:OMS 状态机、算法交易引擎、智能路由、交易所适配层。

4.1 OMS 订单状态机
每笔订单从创建到终态经历明确的状态流转。状态机确保订单在任何时刻都有一个确定的状态,防止漏单、重复下单。

from enum import Enumfrom dataclasses import dataclass, fieldfrom typing importOptionalimport timeclassOrderStatus(Enum): CREATED = "created" SENT = "sent" ACCEPTED = "accepted" PARTIALLY_FILLED = "partially_filled" FILLED = "filled" CANCELLING = "cancelling" CANCELLED = "cancelled" REJECTED = "rejected"classOrderSide(Enum): BUY = "buy" SELL = "sell"@dataclassclassOrder: id: str parent_id: Optional[str] # 母单ID symbol: str side: OrderSide qty: float price: Optional[float] # None=市价单 status: OrderStatus = OrderStatus.CREATED filled_qty: float = 0 filled_price: float = 0 created_at: float = field(default_factory=time.time) updated_at: float = field(default_factory=time.time)@propertydefis_active(self) -> bool:return self.status in ( OrderStatus.CREATED, OrderStatus.SENT, OrderStatus.ACCEPTED, OrderStatus.PARTIALLY_FILLED, OrderStatus.CANCELLING )@propertydefis_terminal(self) -> bool:return self.status in ( OrderStatus.FILLED, OrderStatus.CANCELLED, OrderStatus.REJECTED )classOrderManager:def__init__(self): self._orders: dict[str, Order] = {} self._active_orders: set[str] = set()defcreate_order(self, symbol, side, qty, price=None, parent_id=None) -> Order: order = Order( id=uuid4().hex, parent_id=parent_id, symbol=symbol, side=side, qty=qty, price=price ) self._orders[order.id] = order self._active_orders.add(order.id)return orderdefon_fill(self, fill: Fill): order = self._orders[fill.order_id] order.filled_qty += fill.qty order.filled_price = ( (order.filled_price * (order.filled_qty - fill.qty) + fill.price * fill.qty) / order.filled_qty ) order.updated_at = time.time()if order.filled_qty >= order.qty: order.status = OrderStatus.FILLED self._active_orders.discard(order.id)else: order.status = OrderStatus.PARTIALLY_FILLED
4.2 持仓差异与母单管理
执行层收到目标持仓后,先计算与当前持仓的差异,生成"需要买/卖多少"的订单需求,再由算法交易引擎拆分为子单执行。
classOrderGenerator:"""目标持仓 → 订单需求"""defgenerate(self, target: dict[str, float], current: dict[str, float], prices: dict[str, float], capital: float) -> list[OrderRequest]: requests = []for sym, target_weight in target.items(): current_weight = current.get(sym, 0) delta_weight = target_weight - current_weight delta_shares = int(delta_weight * capital / prices[sym])ifabs(delta_shares) < 100: # 最小交易单位continue side = OrderSide.BUY if delta_shares > 0 else OrderSide.SELL requests.append(OrderRequest( symbol=sym, side=side, qty=abs(delta_shares) ))return requestsclassParentOrderManager:"""母单管理:大单拆子单,跟踪执行进度"""defcreate_parent(self, request: OrderRequest, algo: str = "vwap", duration: int = 300) -> ParentOrder: parent = ParentOrder( id=uuid4().hex, request=request, algo=algo, duration=duration, total_qty=request.qty, filled_qty=0, status="active" )# 启动算法交易引擎执行 self._algo_engine.start(parent)return parent4.3 算法交易引擎
算法交易引擎按策略将母单拆分为子单,分批执行以降低市场冲击。所有算法实现统一接口。
classAlgoBase(ABC):"""算法交易基类"""@abstractmethoddefgenerate_child_orders(self, parent: ParentOrder, market_data: MarketData) -> list[Order]:"""根据市场数据生成下一批子单""" ...classVWAPAlgo(AlgoBase):"""VWAP算法:按历史成交量分布拆单"""def__init__(self): self._volume_profile = {} # 历史成交量分布defgenerate_child_orders(self, parent, market_data): remaining = parent.total_qty - parent.filled_qtyif remaining <= 0:return []# 获取当前时段的成交量占比 profile = self._volume_profile[parent.request.symbol] current_bucket = self._get_time_bucket() expected_frac = profile[current_bucket]# 本时段目标量 = 剩余量 × 时段占比 target_qty = min(remaining, remaining * expected_frac * 2)return [Order( id=uuid4().hex, parent_id=parent.id, symbol=parent.request.symbol, side=parent.request.side, qty=int(target_qty), price=None# 市价单 )]classIcebergAlgo(AlgoBase):"""冰山算法:只暴露小部分委托量"""def__init__(self, visible_fraction=0.1): self.visible_fraction = visible_fractiondefgenerate_child_orders(self, parent, market_data): remaining = parent.total_qty - parent.filled_qtyif remaining <= 0:return [] visible_qty = max(int(remaining * self.visible_fraction), 100)return [Order( id=uuid4().hex, parent_id=parent.id, symbol=parent.request.symbol, side=parent.request.side, qty=visible_qty, price=market_data.mid_price() # 限价单 )]
4.4 智能路由决策
多交易所场景下,路由器评估各 venue 的成本与成交概率,选择最优路径。决策每笔子单实时执行。
classSmartRouter:"""智能路由:选择最优执行 venue"""def__init__(self, venues: dict[str, ExchangeAdapter]): self.venues = venuesdefroute(self, order: Order) -> str:"""返回最优 venue 名称""" scores = {}for name, adapter in self.venues.items(): score = self._score_venue(order, adapter) scores[name] = scorereturnmax(scores, key=scores.get)def_score_venue(self, order, adapter) -> float:# 评分 = 流动性得分 - 成本得分 + 成交概率得分 liquidity = adapter.get_depth(order.symbol) fee = adapter.get_fee(order) fill_prob = adapter.estimate_fill_probability(order)return liquidity * 0.4 - fee * 1000 + fill_prob * 0.44.5 交易所适配层健壮性
自动重连
连接断开时指数退避重连,重连后主动查询所有活跃订单状态。
幂等下单
每个订单携带 client_order_id,交易所去重。重发不会产生重复订单。
心跳保活
每 5 秒发送心跳,3 次未响应判定断线。断线期间冻结订单发送。
状态恢复
系统重启后从持久化存储恢复所有活跃订单,主动查询交易所确认状态。
EXECUTION SAFETY
执行层最重要的安全原则是不漏单、不重复。每笔订单的 client_order_id 全局唯一且持久化,即使系统崩溃重启,也能通过该 ID 向交易所查询确切状态。宁可少成交(后续补单),不可重复下单。
05
跨层基础设施设计
跨层基础设施为四层业务提供通信、配置与状态管理能力,是系统运行的"神经系统"。
5.1 事件总线设计
事件总线是各层异步通信的核心。系统使用两类消息中间件:Kafka 承载高吞吐行情数据流,Redis Stream 传递低延迟内部事件。

事件类型定义
事件 | 通道 | 生产者 | 消费者 | 延迟要求 |
TickEvent | Kafka tick-topic | 数据采集 | 策略引擎 | < 50ms |
BarEvent | Kafka bar-topic | 数据采集 | 策略引擎 | < 100ms |
SignalEvent | Redis signal-stream | 策略引擎 | 风控层 | < 10ms |
OrderEvent | Redis order-stream | 风控层 | OMS | < 10ms |
FillEvent | Redis fill-stream | 交易所适配 | OMS/策略/监控 | < 10ms |
RiskEvent | Redis risk-stream | 风控层 | 监控/告警 | 立即 |
消费者组与断点续传
# Redis Stream 消费者组 —— 断点续传,不丢事件# 每个消费者有独立 ID,崩溃重启后从最后 ACK 位置继续classEventConsumer:def__init__(self, stream: str, group: str, consumer: str): self.stream = stream self.group = group self.consumer = consumer self._redis = redis.Redis() self._ensure_group()def_ensure_group(self):try: self._redis.xgroup_create(self.stream, self.group, id="$")exceptResponseError:pass# group 已存在defconsume(self, handler: Callable, block=1000, count=10):whileTrue:# 读取未 ACK 消息,崩溃恢复后自动重投 messages = self._redis.xreadgroup( self.group, self.consumer, {self.stream: ">"}, count=count, block=block )for msg_id, data in messages:try:handler(Event.from_dict(data)) self._redis.xack(self.stream, self.group, msg_id)exceptException:# 处理失败不 ACK,下次会重投 logger.exception(f"处理失败: {msg_id}")5.2 配置中心
配置中心集中管理所有可变参数:风控阈值、策略参数、数据源连接、交易所合约信息。支持热更新与版本审计。
classConfigCenter:"""配置中心:分层配置 + 热更新 + 审计"""def__init__(self, backend="etcd"): self.backend = backend self._cache = {} # 本地缓存 self._watchers = {} # 变更回调defget(self, key: str, default=None):# 三级查找:本地缓存 → 后端 → 默认值if key in self._cache:return self._cache[key] value = self._fetch_from_backend(key)if value isNone:return default self._cache[key] = valuereturn valuedefwatch(self, key: str, callback: Callable):"""监听配置变更,触发回调""" self._watchers[key] = callback self._backend.watch(key, lambda v: self._on_change(key, v))def_on_change(self, key, value): old = self._cache.get(key) self._cache[key] = value# 审计日志:谁、何时、改了什么 self._audit_log(key, old, value)# 触发回调if key in self._watchers: self._watchers[key](value)# 配置层级(后者覆盖前者)# global.yaml → market.yaml(A股) → strategy.yaml(动量策略) → runtime overrideconfig_hierarchy = ["config/global.yaml", # 全局默认"config/market_a_share.yaml", # A股市场"config/strategy_momentum.yaml", # 策略级]5.3 状态持久化与恢复
系统在运行过程中持续持久化关键状态,确保崩溃重启后能从断点恢复,不丢数据、不重复下单。

恢复优先级
步骤 | 动作 | 数据来源 | 耗时 |
1 | 加载持仓与资金快照 | PostgreSQL + Redis | < 1s |
2 | 重放未确认事件 | Redis Stream pending | < 2s |
3 | 向交易所查询活跃订单 | 交易所 API | < 5s |
4 | 对齐本地与交易所状态 | 差异修正 | < 1s |
5 | 恢复行情订阅与交易 | 数据采集 + 执行层 | < 2s |
classRecoveryManager:"""系统恢复管理器"""defrecover(self):# 1. 加载持仓快照 portfolio = self._load_portfolio_snapshot() account = self._load_account_snapshot()# 2. 重放未确认事件 pending = self._get_pending_events()for event in pending: self._replay_event(event, portfolio)# 3. 向交易所查询活跃订单状态 local_orders = self._get_local_active_orders()for order in local_orders: exchange_status = self._exchange.query_order(order) self._reconcile(order, exchange_status)# 4. 恢复交易 self._resume_trading()def_reconcile(self, local_order, exchange_status):"""对齐本地与交易所状态"""if exchange_status.filled_qty > local_order.filled_qty:# 交易所已成交但本地未记录 → 补记 delta = exchange_status.filled_qty - local_order.filled_qty self._record_fill(local_order, delta, exchange_status.price) logger.warning(f"恢复补单: {local_order.id} 补成交 {delta}")elif exchange_status.is_terminal andnot local_order.is_terminal:# 交易所已终结但本地仍活跃 → 更新状态 local_order.status = exchange_status.statusRECOVERY PRINCIPLE
恢复的核心原则是以交易所为准。本地状态可能在崩溃时丢失或不一致,但交易所的订单状态是唯一真相。恢复时必须主动查询交易所并对齐,任何差异都记录日志。宁可少记成交(后续补记),不可多记(导致重复下单)。

夜雨聆风