TradFi 矩阵统计套利 v1.0 — 策略描述
一句话概括
在交易所的 TradFi 永续合约(美股/指数类,Gate.io contract_type=stocks)上,用低秩矩阵分解 + 零空间投影寻找由共同因子驱动的一篮子标的,实时监测其偏离稳定关系的程度,在"结构错位而非板块异动"时做均值回复交易。默认 paper(模拟盘),不真实下单。
核心思想
把 N 个标的的对数价格排成时间×标的矩阵 X (T×N)。若这些标的由 k 个共同趋势驱动,则 X 近似低秩:X = UΣV',rank(X) ≈ k。
- V 的前 k 个右奇异向量 → 共同因子空间(板块整体在动)
- 后 N−k 个右奇异向量 → 近似零空间,其组合
Xw近似平稳,即可交易的套利篮子
一个篮子同时给出 r = N−k 个正交的均值回复方向 —— 这就是"矩阵套利",而非传统配对套利。
稳定性判据(Davis-Kahan / Wedin sinθ 定理):子空间在扰动下的旋转 ‖sinΘ‖ ≤ ‖E‖ / (σ_k − σ_{k+1})。奇异值间隙 gap = σ_k/σ_{k+1} 越大,零空间越抗噪、越可交易。策略把 gap 与 滚动子空间主夹角作为篮子准入的硬门槛。
流水线(5 个阶段)
| 阶段 | 做什么 | 关键闸门 |
|---|---|---|
| 0 合约池 | 拉全市场行情,筛 stocks 类永续 | 排除 pre-IPO 合成合约、OI/用户数下限、持仓量分位截断、去杠杆版本(相关高+波动比大) |
| 1 候选篮子 | 收益率相关矩阵 → Ward 层次聚类,切出 [6,14] 规模的簇 | 候选不足时启用 kNN 兜底,按平均相关性排序截断 |
| 2 稳定空间检验 | 对每个候选做四道检验 | ① MP 定因子数 k ② 奇异值间隙 ≥1.3 ③ 子空间主夹角 ≤20° ④ 零空间组合 ADF 平稳 + OU 半衰期 ≤200bar |
| 3 实时监测 | 每轮拉一次全市场行情算两个统计量 | Q-score(零空间投影的 Mahalanobis 距离,度量结构错位)+ F-shock(因子得分一阶差分距离,度量板块异动) |
| 4 样本外验证 | 报警后追踪 Q 是否在 TRACK_BARS 内回落 | 回归率 / 回归耗时 / 最大不利偏移 —— 判断策略前提是否成立的唯一依据 |
两个监测统计量(关键设计)
经典 MSPC 的 T²/Q 假设变量平稳,而对数价格是 I(1) 的,不能直接套用,因此做了改造:
- Q-score
~ χ²(N−k):当前价格在零空间的投影,用训练期残差协方差做 Mahalanobis 距离,并按日内季节性(UTC 小时分槽)归一化。度量"结构错位"——这是要交易的方向。 - F-shock
~ χ²(k):因子得分差分的 Mahalanobis 距离,度量"板块整体剧烈异动"。
开仓条件:Q-score 超限 且 F-shock 正常。 若 F-shock 同时超限,说明新信息正在定价(不是错位),此时做均值回复是"接飞刀"。
此外还有事件跳空过滤:单轮 Q 跳变过大判为跳空,进入冷却期禁止开仓。
异常归因(LOO)
残差是组合属性,Q 超限只说明"这堆错位了"。定位用 leave-one-out:用其余 N−1 个标的对标的 i 做 OLS 重构,重构残差相对训练期标准差的 z 值 即 i 的特有偏离(对应 MSPC 的 contribution plot)。开仓方向由归因最强的标的 z 符号决定。
模拟盘执行
- 按 LOO 权重构造多空腿,预算 =
权益 × 杠杆 × 单次占比,整数张取整 - 对冲保真度检查:取整后实际权重与目标权重的余弦相似度须 ≥ 门槛,否则放弃(避免掉腿失衡)
- 成本建模:按各合约
Info读取的 taker 费率 + 滑点系数;资金费率按秒累计(多空组合中不自动抵消,是 spread 的确定性漂移) - 离场:Q 回落至离场线 / 超过 N 倍半衰期 / 保证金亏损止损 / 模型失效
- 持仓用头寸自带的模型快照评估 Q,与篮子重建解耦
已知风险(策略作者自陈)
- TradFi 永续无法与真实美股实物交割,缺少强制收敛机制——约束价格的只是资金费率、做市商库存与报价习惯,弱于真套利。
- 资金费率在多空组合中不抵消,构成 spread 的确定性漂移。
- 合约上线时间短,有效样本可能不足一个完整财报周期。
使用方式
TRADE_MODE默认 paper(live 仅占位未实现)- 命令交互:
refresh:universe(刷新合约池)、rebuild:baskets(重建篮子)、flat:all(全平)、reset:paper(重置账户)、clear:store(清空状态) - 仪表板首屏即为样本外验证统计——这是判断"均值回复前提是否成立"的唯一客观依据;其余为账户、稳定子空间、实时监测、持仓详情、交易流水、去重与淘汰记录。
这个策略本质是监测/研究工具 + 模拟盘验证器,设计哲学非常克制:把"结论正确"让位给"前提可证伪",用样本外回归率而非样本内 z 回归来自我检验。
# -*- coding: utf-8 -*-
"""
TradFi 矩阵统计套利 v1.0(监测 + 模拟盘)
================================================================================
【核心思想】
横轴时间、纵轴标的,构造对数价格矩阵 X (T×N)。
若这 N 个标的由 k 个共同趋势驱动,则 X 近似低秩:
X = U Σ V' rank(X) ≈ k
V 的后 N−k 个右奇异向量张成的子空间 = 近似零空间。
任取其中向量 w,组合 Xw 近似平稳 —— 这就是可交易的套利组合。
(Stock-Watson 共同趋势表示的几何形式)
一个篮子给出 r = N−k 个相互正交的均值回复方向,
因此是"矩阵套利"而非配对套利:同时监测 r 维残差空间。
【为什么零空间是稳定的】
Davis-Kahan / Wedin sinθ 定理:子空间在扰动 E 下的旋转满足
||sinΘ|| ≤ ||E|| / (σ_k − σ_{k+1})
奇异值间隙就是稳定性的定量度量。gap 大 → 子空间抗噪 → 关系可交易。
gap 小 → 特征向量在方向间翻转 → 权重每次重估都乱跳。
本策略把 gap 与滚动子空间主夹角作为篮子准入的硬门槛。
【两个监测统计量(关键设计)】
经典 MSPC 的 T²/Q 假设变量平稳,而对数价格是 I(1) 的,不能直接套用。
因此本策略做了改造:
Q-score : 当前价格向量在零空间的投影 r_t,用训练窗口的残差协方差
做 Mahalanobis 距离。r_t 平稳(若协整成立),~ χ²(N−k)。
→ 度量"结构错位",这是要交易的方向。
F-shock : 因子得分的一阶差分 d_t 的 Mahalanobis 距离,~ χ²(k)。
→ 度量"板块整体剧烈异动"。
只在 Q-score 超限 且 F-shock 正常 时开仓。
F-shock 同时超限意味着新信息正在定价,不是错位,此时做均值回复是接刀。
【异常归因】
残差是组合的属性,不是单个标的的属性。Q 超限只说明"这堆错位了"。
定位用 leave-one-out:用其余 N−1 个标的对标的 i 做 OLS 重构,
重构残差相对训练期标准差的 z 值即 i 的特有偏离。
(对应 MSPC 的 contribution plot)
【样本外验证(一等公民)】
z-score 在样本内必然回归 —— Σz_t = 0 是代数恒等式,随机游走也满足。
唯一有意义的检验是样本外:用训练窗口估的参数,看之后的偏离会不会回归。
仪表板首屏即为该统计:报警数、回归率、回归耗时分布、最大不利偏移。
这张表是判断策略前提是否成立的唯一依据。
【已知风险】
- 币安 TradFi 永续无法与真实美股实物交割,缺少强制收敛机制。
约束这些价格的是资金费率、做市商库存与相对报价习惯,弱于真套利。
- 资金费率在多空组合中不会自动抵消,构成 spread 的确定性漂移。
- 合约上线时间短,有效样本可能不足以覆盖一个完整财报周期。
TRADE_MODE 默认 paper,不会真实下单。
================================================================================
"""
import json
import math
import time
import numpy as np
from scipy import stats as sps
from scipy.cluster.hierarchy import linkage, fcluster
from scipy.spatial.distance import squareform
# ══════════════════════════════════════════════════════════════════════════════
# 参数(平台注入优先,否则用默认值)
# ══════════════════════════════════════════════════════════════════════════════
def _p(name, default):
v = globals().get(name, None)
return default if v is None or v == "" else v
TRADE_MODE = _p("TRADE_MODE", "paper") # paper | live(live 未实现,仅占位)
INIT_CAPITAL = float(_p("INIT_CAPITAL", 1000))
# —— 合约筛选(Gate.io 格式:Info.contract_type)——
# Gate.io 实测品类分布:stocks=284 indices=16 metals=12 commodities=3 forex=3
CONTRACT_TYPES = str(_p("CONTRACT_TYPES", "stocks")) # 逗号分隔;空 contract_type 一律排除
SKIP_PREMARKET = bool(_p("SKIP_PREMARKET", True)) # 排除 pre-IPO 合成合约(ANTHROPIC/ANDURIL等)
MIN_OI_USDT = float(_p("MIN_OI_USDT", 100000)) # 持仓量绝对下限(名义USDT)
MIN_USERS = int(_p("MIN_USERS", 40)) # 多空用户数合计下限
OI_PCTL = float(_p("OI_PCTL", 35)) # 再按持仓量分位数截断(0=不启用)
MAX_UNIVERSE = int(_p("MAX_UNIVERSE", 120)) # 池子上限,按持仓量取前N(控制K线请求量)
DEDUP_RHO = float(_p("DEDUP_RHO", 0.97)) # 近乎同一标的的相关系数门槛
DEDUP_VOL_RATIO = float(_p("DEDUP_VOL_RATIO", 1.40)) # 波动比超此值 → 判为杠杆版本,剔除
VOL_OUTLIER_MULT = float(_p("VOL_OUTLIER_MULT", 1.80)) # 篮子内波动/中位数超此值 → 疑杠杆品,剔除
MIN_HEDGE_FID = float(_p("MIN_HEDGE_FID", 0.97)) # 整数张取整后实际权重与目标的余弦相似度下限
MIN_LEG_CONTRACTS = int(_p("MIN_LEG_CONTRACTS", 1)) # 实测:设高反而因掉腿降低保真度,交由保真度检查把关
FID_TARGET_CT = int(_p("FID_TARGET_CT", 10)) # 每腿目标张数(10张≈5%取整误差),用于估算最小可交易敞口
# —— 成本 ——
TAKER_FEE = float(_p("TAKER_FEE", 0.00075)) # 缺省值;实际按各合约 Info 读取
COST_MULT = float(_p("COST_MULT", 1.0)) # 滑点放大系数
# —— 数据层 ——
KLINE_PERIOD = PERIOD_M15
BAR_MINUTES = int(_p("BAR_MINUTES", 15))
TRAIN_BARS = int(_p("TRAIN_BARS", 2400)) # 期望训练窗口;不足时自动下调(见 effective_train)
MIN_TRAIN_BARS = int(_p("MIN_TRAIN_BARS", 600)) # 训练窗口硬下限,低于此值不做回归
TN_RATIO = float(_p("TN_RATIO", 20)) # 要求 T/N ≥ 此值,否则协方差估计不可靠
MAX_BAR_LEN = int(_p("MAX_BAR_LEN", 3000)) # 向平台申请的K线条数上限
REF_PRESENCE = float(_p("REF_PRESENCE", 0.60)) # 时间点被多少比例的标的覆盖才进参考网格
MIN_COVERAGE = float(_p("MIN_COVERAGE", 0.95)) # 标的对参考网格的覆盖率下限,不足则剔除
FREEZE_FRAC = float(_p("FREEZE_FRAC", 0.98)) # 仅剔除"全员冻结"的bar(标记价格未更新),不剔除清淡时段
VOL_SLOTS = int(_p("VOL_SLOTS", 24)) # 波动率日内季节性分槽数(UTC小时)
MIN_SLOT_BARS = int(_p("MIN_SLOT_BARS", 12)) # 每槽最少样本,不足则退回全局尺度
JUMP_MULT = float(_p("JUMP_MULT", 0.60)) # 单轮Q跳变 > 报警线×此值 → 判为事件跳空
JUMP_COOL_S = int(_p("JUMP_COOL_S", 4 * 3600)) # 判为跳空后的禁开仓冷却
# —— 篮子构建 ——
MIN_BASKET = int(_p("MIN_BASKET", 6))
MAX_BASKET = int(_p("MAX_BASKET", 14))
MAX_BASKETS = int(_p("MAX_BASKETS", 6)) # 最多保留几个通过检验的篮子
LINKAGE_METHOD = str(_p("LINKAGE_METHOD", "ward")) # ward 比 average 更均衡,缓解链式效应
KNN_TRIGGER = int(_p("KNN_TRIGGER", 12)) # 层次聚类候选少于此数 → 启用kNN兜底
MAX_CANDIDATES = int(_p("MAX_CANDIDATES", 150)) # 候选上限(按平均相关性排序后截断)
# —— 稳定性门槛(阶段2 准入) ——
MIN_GAP = float(_p("MIN_GAP", 1.30)) # σ_k / σ_{k+1}
MAX_ANGLE_DEG = float(_p("MAX_ANGLE_DEG", 20.0)) # 子空间最大主夹角
ADF_CRIT = float(_p("ADF_CRIT", -2.86)) # 5% 临界值(含常数项)
MAX_HALFLIFE = float(_p("MAX_HALFLIFE", 200)) # 半衰期上限(bar)
MIN_STABILITY = float(_p("MIN_STABILITY", 0.5))
# —— 监测阈值 ——
ENTRY_Q_P = float(_p("ENTRY_Q_P", 0.995)) # Q-score 报警分位
EXIT_Q_P = float(_p("EXIT_Q_P", 0.80)) # Q-score 回落分位(离场)
FSHOCK_P = float(_p("FSHOCK_P", 0.99)) # F-shock 超此分位 → 禁止开仓
MIN_LOO_Z = float(_p("MIN_LOO_Z", 2.0)) # 归因 z 门槛
# —— 交易 ——
LEVERAGE = float(_p("LEVERAGE", 3))
GROSS_PER_TRADE = float(_p("GROSS_PER_TRADE", 0.25)) # 单次篮子交易占 权益×杠杆 的比例
MAX_OPEN = int(_p("MAX_OPEN", 3)) # 最多同时持有几个篮子头寸
MIN_LEG_USDT = float(_p("MIN_LEG_USDT", 3))
HOLD_HL_MULT = float(_p("HOLD_HL_MULT", 3.0)) # 最大持仓 = 半衰期 × 该倍数
STOP_LOSS_PCT = float(_p("STOP_LOSS_PCT", 18.0)) # 篮子头寸保证金亏损止损%
# —— 节奏 ——
POLL_S = int(_p("POLL_S", 60))
HEARTBEAT_CYCLES = int(_p("HEARTBEAT_CYCLES", 10)) # 每N轮打一次心跳,避免"静默"与"卡死"无法区分
UNIVERSE_S = int(_p("UNIVERSE_S", 6 * 3600))
REBUILD_S = int(_p("REBUILD_S", 12 * 3600))
TRACK_BARS = int(_p("TRACK_BARS", 192)) # 样本外验证:报警后追踪多少bar
STORE_KEY = "tradfi_matrix_statarb_v1"
# 内存缓存(不进 _G,避免存储膨胀)
_REC_CACHE = {}
_chart = None
# ══════════════════════════════════════════════════════════════════════════════
# Store
# ══════════════════════════════════════════════════════════════════════════════
def new_store():
return {
"markets": {}, # symbol -> {contract, ...}
"baskets": {}, # bid -> 通过检验的篮子(含模型参数)
"rejected": [], # 被淘汰的候选及原因
"monitor": {}, # bid -> 最近一次监测结果
"alerts": [], # 报警记录(含样本外追踪)
"positions": {}, # bid -> 模拟持仓
"trades": [],
"logs": [],
"cash": INIT_CAPITAL,
"init_capital": INIT_CAPITAL,
"realized": 0.0,
"last_universe": 0,
"last_rebuild": 0,
}
def load_store():
s = _G(STORE_KEY)
if not s:
return new_store()
for k, v in new_store().items():
if k not in s:
s[k] = v
return s
def save_store(s):
_G(STORE_KEY, s)
def record(store, step, data, color="#999999"):
store["logs"].insert(0, {"time": int(time.time()), "step": step, "data": data})
store["logs"] = store["logs"][:300]
Log("STEP | {} | {}".format(step, json.dumps(data, ensure_ascii=False)[:400]), color)
# ══════════════════════════════════════════════════════════════════════════════
# 数学工具
# ══════════════════════════════════════════════════════════════════════════════
def adf_tstat(y, lags=4):
"""ADF 检验(含常数项)的 t 统计量。返回 None 表示样本不足。
Δy_t = a + γ·y_{t-1} + Σ δ_i·Δy_{t-i} + ε
检验 γ = 0(单位根)。t 统计量越负 → 越拒绝单位根 → 越平稳。
"""
y = np.asarray(y, dtype=float)
n = len(y)
if n < lags + 20:
return None
dy = np.diff(y)
rows = n - 1 - lags
if rows < 10:
return None
Y = dy[lags:]
X = [np.ones(rows), y[lags:n - 1]]
for i in range(1, lags + 1):
X.append(dy[lags - i:n - 1 - i])
X = np.column_stack(X)
try:
beta, _, _, _ = np.linalg.lstsq(X, Y, rcond=None)
except np.linalg.LinAlgError:
return None
resid = Y - X @ beta
dof = rows - X.shape[1]
if dof <= 0:
return None
s2 = float(resid @ resid) / dof
try:
xtx_inv = np.linalg.inv(X.T @ X)
except np.linalg.LinAlgError:
return None
se = math.sqrt(max(s2 * xtx_inv[1, 1], 1e-18))
return float(beta[1] / se)
def ou_halflife(y):
"""OU 过程半衰期(单位:bar)。dX = θ(μ−X)dt + σdW → 半衰期 = ln2/θ。"""
y = np.asarray(y, dtype=float)
if len(y) < 30:
return None
lag = y[:-1]
dy = np.diff(y)
X = np.column_stack([np.ones(len(lag)), lag])
try:
beta, _, _, _ = np.linalg.lstsq(X, dy, rcond=None)
except np.linalg.LinAlgError:
return None
theta = -float(beta[1])
if theta <= 1e-9:
return None # 非回复(发散或随机游走)
return math.log(2.0) / theta
def principal_angles_deg(Q1, Q2):
"""两个子空间(列正交基)之间的最大主夹角,单位度。"""
if Q1.shape[1] == 0 or Q2.shape[1] == 0:
return 90.0
s = np.linalg.svd(Q1.T @ Q2, compute_uv=False)
s = np.clip(s, -1.0, 1.0)
return float(np.degrees(np.arccos(s.min())))
def mp_factor_count(returns):
"""Marchenko-Pastur 上边界法确定共同因子数 k。
returns: (T, N) 收益率矩阵。相关矩阵中大于 (1+√q)² 的特征值视为真实因子,
其余落在 MP 谱内的是纯噪声。
"""
T, N = returns.shape
if T <= N + 2:
return max(1, N // 3)
sd = returns.std(axis=0, ddof=1)
sd[sd < 1e-12] = 1e-12
Z = (returns - returns.mean(axis=0)) / sd
C = (Z.T @ Z) / (T - 1)
ev = np.linalg.eigvalsh(C)
q = float(N) / float(T)
upper = (1.0 + math.sqrt(q)) ** 2
k = int(np.sum(ev > upper))
return int(min(max(k, 1), N - 2))
def orth_basis(M):
"""取矩阵列空间的正交基(用 SVD,数值稳定)。"""
u, s, _ = np.linalg.svd(M, full_matrices=False)
tol = max(M.shape) * np.finfo(float).eps * (s[0] if len(s) else 0.0)
return u[:, s > tol]
# ══════════════════════════════════════════════════════════════════════════════
# 阶段0:合约池 + 数据清洗
# ══════════════════════════════════════════════════════════════════════════════
def refresh_universe(store):
record(store, "universe_start", {}, "#00AAFF")
try:
ms = exchange.GetMarkets()
except Exception as e:
record(store, "universe_failed", {"err": str(e)}, "#FF0000")
return
want = {t.strip().lower() for t in CONTRACT_TYPES.split(",") if t.strip()}
markets, seen_types = {}, {}
drops = {"type": 0, "status": 0, "premarket": 0, "oi": 0, "users": 0}
for key, m in (ms or {}).items():
if ".swap" not in key:
continue
info = m.get("Info", {}) or {}
# —— Gate.io:contract_type 标记品类,股票类为 "stocks" ——
ctype = str(info.get("contract_type", "") or "").lower()
seen_types[ctype or "(空)"] = seen_types.get(ctype or "(空)", 0) + 1
if not ctype or ctype not in want:
drops["type"] += 1
continue
# —— 可交易性 ——
if str(info.get("status", "trading")) != "trading" or info.get("in_delisting", False):
drops["status"] += 1
continue
# —— 排除 pre-IPO 合成合约 ——
if SKIP_PREMARKET and info.get("is_pre_market", False):
drops["premarket"] += 1
continue
ct_val = float(m.get("CtVal", 1) or 1)
mark = float(info.get("mark_price", 0) or info.get("last_price", 0) or 0)
# —— 流动性闸门 ——
oi_usdt = float(info.get("position_size", 0) or 0) * ct_val * mark
users = int(info.get("long_users", 0) or 0) + int(info.get("short_users", 0) or 0)
if oi_usdt < MIN_OI_USDT:
drops["oi"] += 1
continue
if users < MIN_USERS:
drops["users"] += 1
continue
sym = key.replace("_USDT.swap", "").replace(".swap", "").upper()
markets[sym] = {
"symbol": sym,
"contract": key,
"amountPrecision": int(m.get("AmountPrecision", 0) or 0),
"pricePrecision": int(m.get("PricePrecision", 2) or 2),
"ctVal": ct_val,
"minQty": float(m.get("MinQty", 1) or 1),
"takerFee": float(info.get("taker_fee_rate", TAKER_FEE) or TAKER_FEE),
"fundingRate": float(info.get("funding_rate", 0) or 0),
"fundingInterval": int(info.get("funding_interval", 28800) or 28800),
"oiUsdt": round(oi_usdt, 0),
"users": users,
}
# —— 相对流动性截断 ——
ois = sorted(m["oiUsdt"] for m in markets.values())
pctl_info = {}
if ois:
pctl_info = {"p10": round(np.percentile(ois, 10)), "p50": round(np.percentile(ois, 50)),
"p90": round(np.percentile(ois, 90)), "max": round(ois[-1])}
if OI_PCTL > 0 and len(markets) > MIN_BASKET * 2:
cut = float(np.percentile(ois, OI_PCTL))
removed = [s for s, m in markets.items() if m["oiUsdt"] < cut]
for s in removed:
del markets[s]
drops["oi_pctl"] = len(removed)
if len(markets) > MAX_UNIVERSE:
top = sorted(markets.items(), key=lambda kv: -kv[1]["oiUsdt"])[:MAX_UNIVERSE]
drops["universe_cap"] = len(markets) - MAX_UNIVERSE
markets = dict(top)
store["markets"] = markets
record(store, "universe_done", {
"kept": len(markets), "dropped": drops,
"持仓量分布USDT": pctl_info,
"contract_type分布": dict(sorted(seen_types.items(), key=lambda x: -x[1])[:8]),
}, "#00AAFF")
def fetch_records(contract):
"""取K线。优先用三参数形式显式指定条数。"""
if contract in _REC_CACHE:
return _REC_CACHE[contract]
bars = []
try:
bars = exchange.GetRecords(contract, KLINE_PERIOD, MAX_BAR_LEN) or []
except Exception:
try:
bars = exchange.GetRecords(contract, KLINE_PERIOD) or []
except Exception as e:
Log("GetRecords 失败 | {} | {}".format(contract, str(e)), "#FFAA00")
bars = []
_REC_CACHE[contract] = bars
return bars
def build_price_matrix(store, symbols):
"""对齐 + 清洗,返回 (timestamps, log_price_matrix(T×N), kept_symbols)。
新逻辑:按最少可用K线数进行对齐,而非硬性要求每个标的都有MIN_TRAIN_BARS根。
这样即使某些标的历史短,只要池子里有足够的长序列标的,就能凑出长训练窗口。
"""
series, counts = {}, []
for s in symbols:
m = store["markets"].get(s)
if not m:
continue
bars = fetch_records(m["contract"])
counts.append(len(bars))
if len(bars) < 50: # 硬下限:每个标的至少 50 bar(高于此值才能建立基本统计)
continue
series[s] = {int(b["Time"]): float(b["Close"]) for b in bars if b.get("Close", 0) > 0}
if counts:
cs = sorted(counts)
Log("K线条数 | 请求上限:{} 实得 min={} 中位={} max={} | 达标(50+根):{}/{}".format(
MAX_BAR_LEN, cs[0], cs[len(cs) // 2], cs[-1],
len(series), len(counts)), "#888888")
if len(series) < MIN_BASKET:
Log("❌ 达标标的不足({}<{})。若中位数条数明显小于期望,说明平台K线长度受限。".format(
len(series), MIN_BASKET), "#FF0000")
return None, None, []
# —— 时间戳对齐:参考网格 + 覆盖率筛选 ——
n_sym = len(series)
counts = {}
for s in series:
for t in series[s]:
counts[t] = counts.get(t, 0) + 1
ref = {t for t, c in counts.items() if c >= REF_PRESENCE * n_sym}
# 灵活的样本长度调整
if len(ref) < 50:
Log("参考网格不足,放宽条件重试...", "#FFAA00")
ref = {t for t, c in counts.items() if c >= max(1, 0.4 * n_sym)}
if len(ref) < 50:
Log("❌ 对齐失败:无足够重叠时间点", "#FF0000")
return None, None, []
kept, thin = [], []
for s in sorted(series.keys()):
cov = len(ref & set(series[s].keys())) / float(len(ref))
(kept if cov >= MIN_COVERAGE else thin).append(s if cov >= MIN_COVERAGE
else (s, round(cov, 3)))
if thin:
Log("覆盖率不足剔除 {} 个(多为新上市): {}".format(
len(thin), ", ".join("{}({:.0%})".format(a, b) for a, b in thin[:8])), "#999999")
if len(kept) < MIN_BASKET:
return None, None, []
common = ref
for s in kept:
common &= set(series[s].keys())
ts = sorted(common)
# 根据实际数据长度动态调整下限
actual_min_bars = min(MIN_TRAIN_BARS, max(100, len(ts) // 2))
if len(ts) < actual_min_bars:
Log("⚠️ 可用bar({}) 不足目标({}),降低至 {}".format(
len(ts), MIN_TRAIN_BARS, actual_min_bars), "#FFAA00")
if len(ts) < 50:
return None, None, []
P = np.array([[series[s][t] for s in kept] for t in ts], dtype=float)
L = np.log(P)
# —— 只剔除全员冻结的 bar ——
d = np.diff(L, axis=0)
zero_frac = (np.abs(d) < 1e-12).mean(axis=1)
keep_mask = np.concatenate([[True], zero_frac < FREEZE_FRAC])
L = L[keep_mask]
ts = [t for t, k in zip(ts, keep_mask) if k]
if len(ts) < 50:
return None, None, []
return np.array(ts), L, kept
# ══════════════════════════════════════════════════════════════════════════════
# 阶段1:候选篮子(层次聚类)
# ══════════════════════════════════════════════════════════════════════════════
def build_candidate_baskets(store):
"""对收益率相关矩阵做平均连接层次聚类,切出规模在 [MIN,MAX] 的簇。"""
symbols = sorted(store["markets"].keys())
if len(symbols) < MIN_BASKET:
record(store, "cluster_skip", {"reason": "universe too small", "n": len(symbols)}, "#FFAA00")
return []
ts, L, kept = build_price_matrix(store, symbols)
if L is None:
record(store, "cluster_skip", {"reason": "insufficient aligned bars"}, "#FFAA00")
return []
R = np.diff(L, axis=0)
if R.shape[0] < 50:
record(store, "cluster_skip", {"reason": "too few returns", "rows": R.shape[0]}, "#FFAA00")
return []
C = np.corrcoef(R, rowvar=False)
C = np.nan_to_num(C, nan=0.0)
np.fill_diagonal(C, 1.0)
# —— 剔除杠杆版本 ——
vol = R.std(axis=0, ddof=1)
drop_idx, dedup_log = set(), []
for i in range(len(kept)):
for j in range(i + 1, len(kept)):
if abs(C[i, j]) < DEDUP_RHO:
continue
hi, lo = (i, j) if vol[i] > vol[j] else (j, i)
ratio = vol[hi] / max(vol[lo], 1e-12)
if ratio >= DEDUP_VOL_RATIO:
drop_idx.add(hi)
dedup_log.append({"drop": kept[hi], "keep": kept[lo],
"rho": round(float(C[i, j]), 4),
"vol_ratio": round(float(ratio), 2)})
if drop_idx:
record(store, "leveraged_variants_dropped", {"n": len(drop_idx), "detail": dedup_log[:10]}, "#FFAA00")
keep_cols = [i for i in range(len(kept)) if i not in drop_idx]
kept = [kept[i] for i in keep_cols]
L = L[:, keep_cols]
R = R[:, keep_cols]
C = C[np.ix_(keep_cols, keep_cols)]
if len(kept) < MIN_BASKET:
record(store, "cluster_skip", {"reason": "剔除杠杆版本后标的不足", "n": len(kept)}, "#FFAA00")
return []
store["dedup"] = dedup_log[:20]
D = np.sqrt(np.maximum(2.0 * (1.0 - C), 0.0))
np.fill_diagonal(D, 0.0)
D = (D + D.T) / 2.0
Z = linkage(squareform(D, checks=False), method=LINKAGE_METHOD)
n = len(kept)
pos = {sym: i for i, sym in enumerate(kept)}
seen, cands, size_hist = set(), [], {}
def add_cand(members, src):
mt = tuple(sorted(members))
if not (MIN_BASKET <= len(mt) <= MAX_BASKET) or mt in seen:
return
seen.add(mt)
idx = [pos[m] for m in mt]
sub = C[np.ix_(idx, idx)]
avg_rho = float((sub.sum() - len(idx)) / (len(idx) * (len(idx) - 1)))
cands.append({"members": list(mt), "avg_rho": round(avg_rho, 4), "src": src})
for nclust in range(2, min(n, 80) + 1):
labels = fcluster(Z, t=nclust, criterion="maxclust")
for lab in set(labels):
members = [kept[i] for i in range(n) if labels[i] == lab]
size_hist[len(members)] = size_hist.get(len(members), 0) + 1
add_cand(members, "hier")
if len(cands) < KNN_TRIGGER:
Log("层次聚类候选仅 {} 个(链式效应),启用 kNN 兜底".format(len(cands)), "#FFAA00")
for i in range(n):
order = np.argsort(-C[i])
for size in {MIN_BASKET + 2, (MIN_BASKET + MAX_BASKET) // 2, MAX_BASKET}:
add_cand([kept[j] for j in order[:size]], "knn")
cands.sort(key=lambda x: -x["avg_rho"])
cands = cands[:MAX_CANDIDATES]
record(store, "candidates_built", {
"universe": len(kept), "candidates": len(cands),
"来源": {k: sum(1 for c in cands if c["src"] == k) for k in ("hier", "knn")},
"簇规模分布": dict(sorted(size_hist.items())[:12]),
"top": [{"n": len(c["members"]), "rho": c["avg_rho"], "src": c["src"]} for c in cands[:5]],
}, "#00AAFF")
return cands
# ══════════════════════════════════════════════════════════════════════════════
# 阶段2:稳定空间检验(准入闸门)
# ══════════════════════════════════════════════════════════════════════════════
def basket_id(members):
h = 0
for ch in "|".join(sorted(members)):
h = (h * 131 + ord(ch)) & 0xFFFFFFFF
return "B" + format(h, "07x")[:6].upper()
def validate_basket(store, members):
"""对一个候选篮子做四道检验,返回模型或 (None, 淘汰原因)。
1. MP 定共同因子数 k
2. 奇异值间隙 σ_k/σ_{k+1} —— Davis-Kahan 界的分母
3. 滚动子空间主夹角 —— 零空间是否随时间旋转
4. 零空间组合的 ADF + OU 半衰期 —— 是否真的均值回复
"""
ts, L, kept = build_price_matrix(store, members)
if L is None or len(kept) < MIN_BASKET:
return None, "数据不足"
# 杠杆ETF漏网检测
Rv = np.diff(L, axis=0)
vols = Rv.std(axis=0, ddof=1)
med_vol = float(np.median(vols))
if med_vol > 0:
out = [i for i in range(len(kept)) if vols[i] > VOL_OUTLIER_MULT * med_vol]
if out:
names = [kept[i] for i in out]
keep_i = [i for i in range(len(kept)) if i not in set(out)]
if len(keep_i) < MIN_BASKET:
return None, "剔除波动异常成员{}后规模不足".format(names)
Log("篮子内剔除波动异常成员(疑杠杆品): {} | 波动/中位={}".format(
",".join(names),
",".join("{:.1f}x".format(vols[i] / med_vol) for i in out)), "#FFAA00")
kept = [kept[i] for i in keep_i]
L = L[:, keep_i]
T = L.shape[0]
N = L.shape[1]
need = max(50, int(TN_RATIO * N)) # 根据 T/N 比动态调整
if T < need:
return None, "有效bar不足({}<{}=max({},{}×{}))".format(
T, need, 50, int(TN_RATIO), N)
eff = min(TRAIN_BARS, T)
train = L[-eff:]
mu = train.mean(axis=0)
Xc = train - mu
# 1. 共同因子数
k = mp_factor_count(np.diff(train, axis=0))
r = N - k
if r < 1:
return None, "零空间维度为0"
# 2. SVD + 奇异值间隙
_, S, Vt = np.linalg.svd(Xc, full_matrices=False)
if len(S) <= k or S[k] < 1e-12:
return None, "奇异值退化"
gap = float(S[k - 1] / S[k])
if gap < MIN_GAP:
return None, "奇异值间隙不足 {:.2f}<{:.2f}".format(gap, MIN_GAP)
V = Vt.T
Vf = V[:, :k]
Vn = V[:, k:]
# 3. 子空间稳定性
half = eff // 2
angles = []
for seg in (train[:half], train[half:]):
Sc = seg - seg.mean(axis=0)
_, _, vt = np.linalg.svd(Sc, full_matrices=False)
angles.append(principal_angles_deg(Vn, vt.T[:, k:]))
max_angle = max(angles)
if max_angle > MAX_ANGLE_DEG:
return None, "子空间旋转 {:.1f}°>{:.1f}°".format(max_angle, MAX_ANGLE_DEG)
# 4. 零空间组合的平稳性与回复速度
Rn = Xc @ Vn
adfs, hls = [], []
for j in range(r):
a = adf_tstat(Rn[:, j])
h = ou_halflife(Rn[:, j])
if a is not None:
adfs.append(a)
if h is not None:
hls.append(h)
if not adfs or not hls:
return None, "平稳性检验失败"
keep_dirs = [j for j in range(r) if (adf_tstat(Rn[:, j]) or 0) < ADF_CRIT]
pass_adf = len(keep_dirs)
if pass_adf == 0:
return None, "无方向通过ADF(最优t={:.2f})".format(min(adfs))
if pass_adf < r:
Vn = Vn[:, keep_dirs]
Rn = Xc @ Vn
r = pass_adf
halflife = float(np.median(hls))
if halflife > MAX_HALFLIFE:
return None, "半衰期过长 {:.0f}>{:.0f}bar".format(halflife, MAX_HALFLIFE)
# —— 监测所需的训练期统计量 ——
resid_cov = np.cov(Rn, rowvar=False)
if r == 1:
resid_cov = np.array([[float(resid_cov)]])
resid_inv = np.linalg.pinv(resid_cov)
# —— 日内季节性标准化 ——
Q_train = np.einsum("ij,jk,ik->i", Rn, resid_inv, Rn)
slots = ((ts[-eff:] // 3600000) % VOL_SLOTS).astype(int)
q_med_all = float(np.median(Q_train)) or 1.0
slot_scale = {}
for h in range(VOL_SLOTS):
sel = Q_train[slots == h]
if len(sel) >= MIN_SLOT_BARS:
slot_scale[str(h)] = round(max(float(np.median(sel)) / q_med_all, 0.05), 4)
F = Xc @ Vf
dF = np.diff(F, axis=0)
f_cov = np.cov(dF, rowvar=False)
if k == 1:
f_cov = np.array([[float(f_cov)]])
f_inv = np.linalg.pinv(f_cov)
# LOO 归因模型
loo = []
for i in range(N):
others = [j for j in range(N) if j != i]
A = np.column_stack([np.ones(eff), train[:, others]])
b, _, _, _ = np.linalg.lstsq(A, train[:, i], rcond=None)
res = train[:, i] - A @ b
loo.append({"i": i, "others": others, "beta": b.tolist(),
"sd": float(res.std(ddof=1)) or 1e-9})
# —— 最小可交易敞口 ——
last_px = np.exp(train[-1])
gran = np.array([last_px[i] * store["markets"][kept[i]]["ctVal"] for i in range(N)])
needs = []
for rec in loo:
b = np.array(rec["beta"])
wv = np.zeros(N)
wv[rec["i"]] = -1.0
for pj, j in enumerate(rec["others"]):
wv[j] = float(b[1 + pj])
aw = np.abs(wv)
nz = aw > 1e-9
if nz.sum() < 2:
continue
needs.append(float(np.max(FID_TARGET_CT * gran[nz] * aw.sum() / aw[nz])))
min_gross = float(np.median(needs)) if needs else float("inf")
stability = float(np.clip(
0.35 * np.clip((gap - 1.0) / 1.0, 0, 1)
+ 0.30 * np.clip(1.0 - max_angle / MAX_ANGLE_DEG, 0, 1)
+ 0.20 * (pass_adf / r)
+ 0.15 * np.clip(1.0 - halflife / MAX_HALFLIFE, 0, 1), 0, 1))
model = {
"members": kept,
"N": N, "k": k, "r": r,
"mu": mu.tolist(),
"Vf": Vf.tolist(), "Vn": Vn.tolist(),
"resid_inv": resid_inv.tolist(),
"f_inv": f_inv.tolist(),
"slot_scale": slot_scale,
"loo": loo,
"gap": round(gap, 3),
"max_angle": round(max_angle, 2),
"adf_best": round(min(adfs), 3),
"adf_pass": pass_adf,
"halflife": round(halflife, 1),
"min_gross": round(min_gross, 1),
"stability": round(stability, 3),
"train_bars": eff,
"q_entry": float(sps.chi2.ppf(ENTRY_Q_P, r)),
"q_exit": float(sps.chi2.ppf(EXIT_Q_P, r)),
"f_limit": float(sps.chi2.ppf(FSHOCK_P, k)),
"built_at": int(time.time()),
}
return model, None
def rebuild_baskets(store):
record(store, "rebuild_start", {}, "#00AAFF")
_REC_CACHE.clear()
cands = build_candidate_baskets(store)
passed, rejected, used = [], [], set()
for c in cands:
members = c["members"]
if len(set(members) & used) > len(members) // 2:
continue
model, reason = validate_basket(store, members)
if model is None:
rejected.append({"members": members[:6], "n": len(members), "reason": reason})
Log("篮子淘汰 | n={} | {} | {}".format(len(members), reason,
",".join(members[:6])), "#999999")
continue
if model["stability"] < MIN_STABILITY:
rejected.append({"members": members[:6], "n": len(members),
"reason": "稳定性评分不足 {:.2f}".format(model["stability"])})
continue
passed.append(model)
used |= set(model["members"])
Log("篮子通过 | n={} | gap={} 夹角={}° ADF={} 半衰期={}bar 评分={} | {}".format(
model["N"], model["gap"], model["max_angle"], model["adf_best"],
model["halflife"], model["stability"], ",".join(model["members"])), "#00CC66")
if len(passed) >= MAX_BASKETS:
break
passed.sort(key=lambda m: -m["stability"])
store["baskets"] = {basket_id(m["members"]): m for m in passed}
store["rejected"] = rejected[:40]
record(store, "rebuild_done", {
"passed": len(passed), "rejected": len(rejected),
"detail": [{"id": bid, "n": m["N"], "k": m["k"], "r": m["r"],
"gap": m["gap"], "hl": m["halflife"], "score": m["stability"]}
for bid, m in store["baskets"].items()],
}, "#00AAFF")
# ══════════════════════════════════════════════════════════════════════════════
# 阶段3:实时监测(Q-score / F-shock / LOO 归因)
# ══════════════════════════════════════════════════════════════════════════════
def fetch_all_tickers(store):
"""整轮只拉一次全市场行情,避免逐腿 GetTicker 触发限频。"""
out = {}
try:
ts = exchange.GetTickers() or []
except Exception as e:
Log("GetTickers 异常:", str(e), "#FFAA00")
return out
by_contract = {}
for t in ts:
last = t.get("Last", 0) if isinstance(t, dict) else 0
sym = t.get("Symbol", "") if isinstance(t, dict) else ""
if last and last > 0:
by_contract[sym] = float(last)
for sym, m in store["markets"].items():
p = by_contract.get(m["contract"])
if p:
out[sym] = p
return out
def current_log_prices(members, tick):
out, prices = [], {}
for s in members:
p = tick.get(s)
if not p or p <= 0:
return None, None
prices[s] = p
out.append(math.log(p))
return np.array(out), prices
def monitor_basket(store, bid, model, prev, tick):
x, prices = current_log_prices(model["members"], tick)
if x is None:
return None
mu = np.array(model["mu"])
Vn = np.array(model["Vn"])
Vf = np.array(model["Vf"])
xc = x - mu
# Q-score
rvec = xc @ Vn
q_raw = float(rvec @ np.array(model["resid_inv"]) @ rvec)
slot = str(int((time.time() // 3600) % VOL_SLOTS))
scale = float(model.get("slot_scale", {}).get(slot, 1.0))
q = q_raw / scale
# F-shock
f = xc @ Vf
fshock = 0.0
if prev is not None and prev.get("f") is not None:
d = f - np.array(prev["f"])
fshock = float(d @ np.array(model["f_inv"]) @ d)
# LOO 归因
contrib = []
for rec in model["loo"]:
b = np.array(rec["beta"])
pred = b[0] + float(b[1:] @ x[rec["others"]])
z = (x[rec["i"]] - pred) / rec["sd"]
contrib.append({"symbol": model["members"][rec["i"]], "z": round(float(z), 3)})
contrib.sort(key=lambda c: -abs(c["z"]))
# —— 跳空过滤 ——
now = int(time.time())
jump_until = (prev or {}).get("jump_until", 0)
if prev is not None and prev.get("q") is not None:
if q - prev["q"] > model["q_entry"] * JUMP_MULT:
jump_until = now + JUMP_COOL_S
Log("⚡ 判定为事件跳空 | {} | Q {:.1f}→{:.1f} | 冷却 {}h".format(
bid, prev["q"], q, JUMP_COOL_S // 3600), "#FFAA00")
in_jump = now < jump_until
return {
"bid": bid,
"q": round(q, 2),
"q_raw": round(q_raw, 2),
"slot_scale": round(scale, 3),
"q_entry": round(model["q_entry"], 2),
"q_exit": round(model["q_exit"], 2),
"fshock": round(fshock, 2),
"f_limit": round(model["f_limit"], 2),
"alarm": q > model["q_entry"],
"blocked": fshock > model["f_limit"] or in_jump,
"block_reason": "因子异动" if fshock > model["f_limit"] else ("事件跳空冷却" if in_jump else ""),
"jump_until": jump_until,
"contrib": contrib[:5],
"f": f.tolist(),
"prices": prices,
"at": now,
}
def run_monitor(store):
tick = fetch_all_tickers(store)
results = {}
for bid, model in store["baskets"].items():
prev = store["monitor"].get(bid)
m = monitor_basket(store, bid, model, prev, tick)
if m:
results[bid] = m
store["monitor"] = results
return results
# ══════════════════════════════════════════════════════════════════════════════
# 阶段4:样本外验证(报警追踪)
# ══════════════════════════════════════════════════════════════════════════════
def track_alerts(store, results):
"""报警后追踪 Q 是否回落。"""
now = int(time.time())
horizon = TRACK_BARS * BAR_MINUTES * 60
for a in store["alerts"]:
if a.get("resolved") is not None:
continue
r = results.get(a["bid"])
if not r:
continue
a["peak_q"] = max(a.get("peak_q", a["q0"]), r["q"])
if r["q"] <= a["exit_q"]:
a["resolved"] = True
a["resolve_s"] = now - a["time"]
elif now - a["time"] > horizon:
a["resolved"] = False
a["resolve_s"] = now - a["time"]
for bid, r in results.items():
if not r["alarm"] or r["blocked"]:
continue
recent = [a for a in store["alerts"]
if a["bid"] == bid and a.get("resolved") is None]
if recent:
continue
model = store["baskets"][bid]
store["alerts"].insert(0, {
"time": now, "bid": bid, "q0": r["q"], "peak_q": r["q"],
"exit_q": model["q_exit"], "resolved": None, "resolve_s": 0,
"anomaly": r["contrib"][0]["symbol"] if r["contrib"] else "",
"z": r["contrib"][0]["z"] if r["contrib"] else 0,
})
record(store, "alert_raised", {
"bid": bid, "q": r["q"], "limit": r["q_entry"],
"anomaly": r["contrib"][0] if r["contrib"] else None,
}, "#FF0000")
store["alerts"] = store["alerts"][:200]
def validation_stats(store):
done = [a for a in store["alerts"] if a.get("resolved") is not None]
if not done:
return {"total": len(store["alerts"]), "closed": 0, "rate": None,
"median_min": None, "mae": None}
ok = [a for a in done if a["resolved"]]
times = sorted(a["resolve_s"] / 60.0 for a in ok)
mae = [a["peak_q"] / max(a["q0"], 1e-9) for a in done]
return {
"total": len(store["alerts"]),
"closed": len(done),
"rate": round(len(ok) / len(done), 3),
"median_min": round(times[len(times) // 2], 1) if times else None,
"mae": round(float(np.mean(mae)), 2),
}
# ══════════════════════════════════════════════════════════════════════════════
# 模拟盘引擎
# ══════════════════════════════════════════════════════════════════════════════
def leg_pnl(leg, price):
"""合约盈亏。"""
diff = price - leg["entry"] if leg["side"] == "long" else leg["entry"] - price
return diff * leg["contracts"] * leg["ctVal"]
def position_pnl(pos, price_map):
"""含手续费与资金费率的净盈亏。"""
gross = sum(leg_pnl(l, price_map.get(s, l["entry"])) for s, l in pos["legs"].items())
return gross - pos.get("fee_paid", 0.0) - pos.get("funding_paid", 0.0)
def equity(store, price_map):
eq = store["cash"]
for pos in store["positions"].values():
eq += sum(l["margin"] for l in pos["legs"].values())
eq += position_pnl(pos, price_map)
return eq
def all_prices(store):
pm = {}
for r in store["monitor"].values():
pm.update(r.get("prices", {}))
return pm
def open_basket(store, bid, model, r):
"""按 LOO 归因方向建仓。"""
if len(store["positions"]) >= MAX_OPEN or bid in store["positions"]:
return
top = r["contrib"][0] if r["contrib"] else None
if not top or abs(top["z"]) < MIN_LOO_Z:
record(store, "entry_skip_weak_attribution", {"bid": bid, "top": top}, "#FFAA00")
return
i = model["members"].index(top["symbol"])
rec = next(x for x in model["loo"] if x["i"] == i)
b = np.array(rec["beta"])
sgn = 1.0 if top["z"] > 0 else -1.0
w = {top["symbol"]: -sgn}
for pos_j, j in enumerate(rec["others"]):
w[model["members"][j]] = sgn * float(b[1 + pos_j])
gross = sum(abs(v) for v in w.values())
if gross <= 0:
return
budget = equity(store, r["prices"]) * LEVERAGE * GROSS_PER_TRADE
legs, need, fee, dropped = {}, 0.0, 0.0, []
for sym, wt in w.items():
price = r["prices"].get(sym)
mk = store["markets"].get(sym)
if not price or not mk:
continue
target = budget * abs(wt) / gross
ct_val = mk["ctVal"]
per_ct = price * ct_val
contracts = math.floor(target / per_ct) if per_ct > 0 else 0
if contracts < max(mk["minQty"], MIN_LEG_CONTRACTS):
dropped.append({"sym": sym, "want": round(target, 2), "contracts": contracts})
continue
notional = contracts * per_ct
if notional < MIN_LEG_USDT:
dropped.append({"sym": sym, "notional": round(notional, 2)})
continue
margin = notional / LEVERAGE
legs[sym] = {"side": "long" if wt > 0 else "short", "entry": price,
"contracts": contracts, "ctVal": ct_val, "notional": notional,
"margin": margin, "w": round(wt, 5)}
need += margin
fee += notional * mk["takerFee"] * COST_MULT
if len(legs) < 2:
record(store, "entry_skip_legs", {"bid": bid, "legs": len(legs)}, "#FFAA00")
return
# —— 对冲保真度检查 ——
tgt = np.array([w.get(sy, 0.0) for sy in w])
act = np.array([(legs[sy]["notional"] * (1 if legs[sy]["side"] == "long" else -1))
if sy in legs else 0.0 for sy in w])
nt, na = np.linalg.norm(tgt), np.linalg.norm(act)
fidelity = float(tgt @ act / (nt * na)) if nt > 0 and na > 0 else 0.0
if fidelity < MIN_HEDGE_FID:
record(store, "entry_skip_hedge_fidelity", {
"bid": bid, "fidelity": round(fidelity, 4), "门槛": MIN_HEDGE_FID}, "#FFAA00")
return
if need + fee > store["cash"] - 1:
record(store, "entry_skip_capital", {"bid": bid, "need": round(need + fee, 2)}, "#FFAA00")
return
round_trip = fee * 2
store["cash"] -= (need + fee)
store["positions"][bid] = {
"bid": bid, "legs": legs, "open_time": int(time.time()),
"last_funding": int(time.time()),
"entry_q": r["q"], "anomaly": top["symbol"], "z": top["z"],
"margin": need, "fee_paid": fee, "funding_paid": 0.0,
"est_round_trip": round(round_trip, 4),
"fidelity": round(fidelity, 4),
"max_hold_s": int(model["halflife"] * HOLD_HL_MULT * BAR_MINUTES * 60),
"model": {k: model[k] for k in
("members", "mu", "Vn", "resid_inv", "q_exit", "slot_scale", "r")},
}
trade = {"time": int(time.time()), "type": "open", "bid": bid,
"anomaly": top["symbol"], "z": top["z"], "q": r["q"],
"legs": len(legs), "margin": round(need, 2),
"fee": round(fee, 4), "往返成本": round(round_trip, 4)}
store["trades"].insert(0, trade)
store["trades"] = store["trades"][:100]
record(store, "paper_open", trade, "#00CC66")
def close_basket(store, bid, reason, price_map):
pos = store["positions"].pop(bid, None)
if not pos:
return
exit_fee = sum(l["notional"] * store["markets"].get(s, {}).get("takerFee", TAKER_FEE)
* COST_MULT for s, l in pos["legs"].items())
pos["fee_paid"] = pos.get("fee_paid", 0.0) + exit_fee
pnl = position_pnl(pos, price_map)
for leg in pos["legs"].values():
store["cash"] += leg["margin"]
store["cash"] += pnl
store["realized"] += pnl
trade = {"time": int(time.time()), "type": "close", "bid": bid,
"reason": reason, "pnl": round(pnl, 4),
"fee": round(pos["fee_paid"], 4),
"funding": round(pos.get("funding_paid", 0.0), 4),
"hold_min": round((int(time.time()) - pos["open_time"]) / 60.0, 1)}
store["trades"].insert(0, trade)
store["trades"] = store["trades"][:100]
record(store, "paper_close", trade, "#FFAA00" if pnl < 0 else "#00CC66")
def accrue_funding(store, pm):
"""按秒累计资金费率。"""
now = int(time.time())
for pos in store["positions"].values():
dt = now - pos.get("last_funding", now)
if dt <= 0:
continue
pos["last_funding"] = now
acc = 0.0
for sym, leg in pos["legs"].items():
mk = store["markets"].get(sym)
if not mk or not mk["fundingInterval"]:
continue
p = pm.get(sym, leg["entry"])
notional = leg["contracts"] * leg["ctVal"] * p
rate = mk["fundingRate"] * dt / mk["fundingInterval"]
acc += notional * rate if leg["side"] == "long" else -notional * rate
pos["funding_paid"] = pos.get("funding_paid", 0.0) + acc
def position_q(pos, tick):
"""用头寸自带的模型快照算 Q。"""
m = pos.get("model")
if not m:
return None
x, _ = current_log_prices(m["members"], tick)
if x is None:
return None
rv = (x - np.array(m["mu"])) @ np.array(m["Vn"])
q_raw = float(rv @ np.array(m["resid_inv"]) @ rv)
slot = str(int((time.time() // 3600) % VOL_SLOTS))
return q_raw / float(m.get("slot_scale", {}).get(slot, 1.0))
def manage_positions(store, results):
now = int(time.time())
pm = all_prices(store)
tick = {s: p for s, p in pm.items()}
accrue_funding(store, pm)
for bid in list(store["positions"].keys()):
pos = store["positions"][bid]
q_now = position_q(pos, tick)
if q_now is not None and q_now <= pos["model"]["q_exit"]:
close_basket(store, bid, "Q回落", pm)
continue
if q_now is None and bid not in store["baskets"]:
close_basket(store, bid, "模型失效(无法评估Q)", pm)
continue
if now - pos["open_time"] > pos["max_hold_s"]:
close_basket(store, bid, "超过{}倍半衰期".format(HOLD_HL_MULT), pm)
continue
pnl = position_pnl(pos, pm)
if pos["margin"] > 0 and pnl / pos["margin"] * 100 <= -STOP_LOSS_PCT:
close_basket(store, bid, "止损", pm)
for bid, r in results.items():
if r["alarm"] and not r["blocked"] and bid in store["baskets"]:
open_basket(store, bid, store["baskets"][bid], r)
# ══════════════════════════════════════════════════════════════════════════════
# 仪表板
# ══════════════════════════════════════════════════════════════════════════════
def tbl(title, cols, rows):
return "\n`" + json.dumps({"type": "table", "title": title, "cols": cols,
"rows": rows or [["-"] * len(cols)]},
ensure_ascii=False) + "`\n"
def short(x, n=40):
s = str(x or "").replace("\n", " ")
return s[:n] + ("..." if len(s) > n else "")
def render(store):
pm = all_prices(store)
eq = equity(store, pm)
init = store["init_capital"]
v = validation_stats(store)
# 首屏:样本外验证
t_valid = tbl("🔬 样本外验证(判断前提是否成立的唯一依据)",
["累计报警", "已了结", "回归率", "回归耗时中位数(min)", "最大不利偏移倍数", "结论"],
[[v["total"], v["closed"],
"—" if v["rate"] is None else "{:.1%}".format(v["rate"]),
v["median_min"] if v["median_min"] is not None else "—",
v["mae"] if v["mae"] is not None else "—",
"样本不足,继续积累" if v["closed"] < 20 else
("✅ 前提成立" if (v["rate"] or 0) >= 0.6 else "⚠️ 回归率偏低,前提存疑")]])
t_acct = tbl("💰 模拟账户", ["模式", "权益", "总盈亏", "收益率", "现金", "已实现", "持仓篮子"],
[[TRADE_MODE, "${:.2f}".format(eq),
"{}${:.2f}".format("+" if eq >= init else "-", abs(eq - init)),
"{:+.2f}%".format((eq - init) / init * 100),
"${:.2f}".format(store["cash"]),
"{:+.2f}".format(store["realized"]),
"{}/{}".format(len(store["positions"]), MAX_OPEN)]])
rows = []
for bid, m in store["baskets"].items():
budget = eq * LEVERAGE * GROSS_PER_TRADE
mg = m.get("min_gross", 0)
rows.append([bid, m["N"], m["k"], m["r"], m["gap"], "{}°".format(m["max_angle"]),
m["adf_best"], "{}/{}".format(m["adf_pass"], m["r"]),
m["halflife"], m["stability"],
"${:,.0f}".format(mg) if mg != float("inf") else "∞",
"✅" if budget >= mg else "❌需${:,.0f}权益".format(
mg / max(LEVERAGE * GROSS_PER_TRADE, 1e-9)),
short(",".join(m["members"]), 50)])
t_basket = tbl("🧮 通过检验的稳定子空间",
["ID", "N", "k因子", "r零空间", "σ间隙", "子空间夹角", "ADF最优t",
"ADF通过", "半衰期(bar)", "稳定评分", "最小敞口", "当前可交易", "成分"], rows)
rows = []
for bid, r in store["monitor"].items():
state = "🔴 报警" if r["alarm"] and not r["blocked"] else \
("⏸ " + r.get("block_reason", "禁开") if r["blocked"] else "🟢 正常")
rows.append([bid, state, r["q"], r.get("q_raw", "-"), r.get("slot_scale", "-"),
r["q_entry"], r["q_exit"], r["fshock"], r["f_limit"],
" ".join("{}:{}".format(c["symbol"], c["z"]) for c in r["contrib"][:4])])
t_mon = tbl("📡 实时监测(Q超限 且 无禁开标记 才开仓)",
["ID", "状态", "Q(季节归一)", "Q(原始)", "时段因子", "报警线", "离场线",
"F-shock", "F上限", "LOO归因(z)"], rows)
# 新增:每个开仓品种的详细信息
leg_rows = []
for bid, pos in store["positions"].items():
for sym, leg in pos["legs"].items():
current_price = pm.get(sym, leg["entry"])
pnl_per_contract = leg_pnl(leg, current_price) / leg["contracts"] if leg["contracts"] > 0 else 0
roi = (pnl_per_contract / leg["entry"] * 100) if leg["entry"] > 0 else 0
leg_rows.append([
bid, sym, leg["side"].upper(),
"${:.4f}".format(leg["entry"]),
"${:.4f}".format(current_price),
"{:+.2f}%".format(roi),
leg["contracts"],
"${:.2f}".format(leg_pnl(leg, current_price)),
])
t_legs = tbl("📊 开仓品种详情(开仓价 → 实时价 → 盈利率)",
["篮子ID", "品种", "方向", "开仓价", "实时价", "盈利率%", "张数", "盈亏$"], leg_rows)
rows = []
for bid, p in store["positions"].items():
gross = sum(leg_pnl(l, pm.get(s, l["entry"])) for s, l in p["legs"].items())
qn = position_q(p, pm)
rows.append([bid, p["anomaly"], p["z"],
"{:.1f}→{:.1f}".format(p["entry_q"], qn) if qn is not None else "—",
p.get("fidelity", "—"), len(p["legs"]),
"${:.2f}".format(p["margin"]), "{:+.3f}".format(gross),
"{:.3f}".format(p.get("fee_paid", 0)),
"{:+.4f}".format(p.get("funding_paid", 0)),
"{:+.3f}".format(position_pnl(p, pm)),
"{:.1f}".format((int(time.time()) - p["open_time"]) / 60.0),
" ".join("{}{}x{}".format("+" if l["side"] == "long" else "-", s,
l["contracts"])
for s, l in list(p["legs"].items())[:5])])
t_pos = tbl("📦 持仓篮子汇总", ["ID", "异常标的", "z", "Q(进场→当前)", "对冲保真", "腿数",
"保证金", "毛盈亏", "手续费", "资金费", "净盈亏",
"持仓(min)", "腿(张数)"], rows)
rows = [[time.strftime("%m-%d %H:%M", time.localtime(t["time"])), t["type"], t["bid"],
t.get("anomaly", "-"), t.get("q", t.get("reason", "-")),
t.get("pnl", "-"), t.get("fee", "-"), t.get("funding", "-"),
t.get("hold_min", "-")]
for t in store["trades"][:8]]
t_trade = tbl("📜 最近交易", ["时间", "类型", "篮子", "异常标的", "Q/原因",
"净盈亏$", "手续费", "资金费", "持仓(min)"], rows)
rows = [[d["drop"], d["keep"], d["rho"], "{}x".format(d["vol_ratio"])]
for d in store.get("dedup", [])[:8]]
t_dedup = tbl("🔀 剔除的杠杆版本(波动率衰减会破坏协整)",
["剔除", "保留", "相关系数", "波动比"], rows)
rows = [[r["n"], short(",".join(r["members"]), 44), r["reason"]] for r in store["rejected"][:8]]
t_rej = tbl("🚫 淘汰的候选篮子(为什么不交易它)", ["规模", "成分", "淘汰原因"], rows)
LogStatus("TradFi 矩阵统计套利 v1.0 | {} | 合约池:{} | {}\n".format(
TRADE_MODE, len(store["markets"]), time.strftime("%Y-%m-%d %H:%M:%S"))
+ t_valid + t_acct + t_basket + t_mon + t_legs + t_pos + t_trade + t_dedup + t_rej)
# ══════════════════════════════════════════════════════════════════════════════
# 命令
# ══════════════════════════════════════════════════════════════════════════════
def handle_command(store):
cmd = GetCommand()
if not cmd:
return
if cmd == "refresh:universe":
refresh_universe(store)
store["last_universe"] = int(time.time())
elif cmd == "rebuild:baskets":
rebuild_baskets(store)
store["last_rebuild"] = int(time.time())
elif cmd == "flat:all":
pm = all_prices(store)
for bid in list(store["positions"].keys()):
close_basket(store, bid, "手动平仓", pm)
elif cmd == "reset:paper":
store["cash"] = INIT_CAPITAL
store["init_capital"] = INIT_CAPITAL
store["positions"] = {}
store["trades"] = []
store["realized"] = 0.0
Log("模拟账户已重置", "#FFAA00")
elif cmd == "clear:store":
store.clear()
store.update(new_store())
Log("store 已清空", "#FFAA00")
save_store(store)
# ══════════════════════════════════════════════════════════════════════════════
# 主循环
# ══════════════════════════════════════════════════════════════════════════════
def main():
_G(None)
global _chart
LogReset(0)
Log("TradFi 矩阵统计套利 v1.0 启动 | 模式:", TRADE_MODE,
"| 初始资金:", INIT_CAPITAL, "#00AAFF")
Log("注意:币安 TradFi 永续无法与美股实物交割,缺少强制收敛机制。"
"请以仪表板首屏的样本外验证统计为准。", "#FFAA00")
_chart = Chart({
"title": {"text": "模拟账户权益"},
"xAxis": {"type": "datetime"},
"yAxis": {"title": {"text": "权益(USDT)"}},
"series": [{"name": "权益", "type": "line", "data": []}],
})
store = load_store()
Log("载入持久化状态 | 合约池:{} | 篮子:{} | 持仓:{} | 报警:{} | 现金:{:.2f}".format(
len(store["markets"]), len(store["baskets"]), len(store["positions"]),
len(store["alerts"]), store["cash"]), "#00AAFF")
cycle = 0
while True:
try:
handle_command(store)
now = int(time.time())
cycle += 1
need_uni = (not store["markets"]) or (now - store["last_universe"] > UNIVERSE_S)
if need_uni:
if not store["markets"] and store["last_universe"] > 0:
Log("合约池为空但定时器未到期 → 强制刷新", "#FFAA00")
refresh_universe(store)
store["last_universe"] = now
save_store(store)
need_reb = store["markets"] and ((not store["baskets"]) or
now - store["last_rebuild"] > REBUILD_S)
if need_reb:
if not store["baskets"] and store["last_rebuild"] > 0:
Log("无可用篮子但定时器未到期 → 强制重建", "#FFAA00")
rebuild_baskets(store)
store["last_rebuild"] = now
save_store(store)
if store["baskets"]:
results = run_monitor(store)
track_alerts(store, results)
manage_positions(store, results)
render(store)
eq = equity(store, all_prices(store))
_chart.add(0, [now * 1000, eq])
LogProfit(eq - store["init_capital"], "&")
save_store(store)
if cycle == 1 or cycle % HEARTBEAT_CYCLES == 0:
alarms = sum(1 for r in store["monitor"].values() if r.get("alarm"))
Log("♥ 第{}轮 | 池:{} 篮子:{} 监测:{} 报警:{} 持仓:{} | 权益:{:.2f} | "
"下次重建:{:.1f}h".format(
cycle, len(store["markets"]), len(store["baskets"]),
len(store["monitor"]), alarms, len(store["positions"]), eq,
max(0, REBUILD_S - (now - store["last_rebuild"])) / 3600.0),
"#888888")
except Exception as e:
import traceback
Log("主循环异常:", str(e), "\n", traceback.format_exc()[:1200], "#FF0000")
Sleep(POLL_S * 1000)
- 1