乐于分享
好东西不私藏

量化系统详细设计文档

量化系统详细设计文档

 对面向中低频多策略交易的端到端量化系统架构 各层模块的深入技术设计,涵盖数据结构、接口契约、核心算法与实现要点,面向研发实现。

目标读者:研发工程师

前置文档:架构设计文档

具体到量化策略调优操作,估计还得好几篇。。微信排版好头疼,先看看点赞多不多再看还写不写

目录 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 86400

1.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 importDictOptional@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[strfloat] = {}        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[strfloat]):        self.target_portfolio = weights

DESIGN 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        )

信号聚合策略

聚合方法

公式

适用场景

等权平均

score = mean(signal_i)

各信号源质量相近

IC 加权

score = Σ(ic_i × signal_i) / Σ|ic_i|

有历史 IC 统计

波动率倒数

score = Σ(signal_i / vol_i) / Σ(1/vol_i)

各信号波动差异大

投票制

score = sign(Σ sign(signal_i))

规则型信号离散

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[strfloat] = {}  # 各策略资金占比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[strfloat],                  current: dict[strfloat],                  prices: dict[strfloat],                  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 parent

4.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.4

4.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.status

RECOVERY PRINCIPLE

恢复的核心原则是以交易所为准。本地状态可能在崩溃时丢失或不一致,但交易所的订单状态是唯一真相。恢复时必须主动查询交易所并对齐,任何差异都记录日志。宁可少记成交(后续补记),不可多记(导致重复下单)。