Files
Chan/research/live/shadow_signal.py
T
jackandCursor 6c843639d5 影子路径开 lean 模式,实测完整链路 1410 → 683ms
拉到新引擎后在服务器侧直接实测,没有沿用 3.4 倍那个换算——那是在旧代码上量的。

先做等价性验证。不能直接引用 step46 的对拍结论:它固化的是 bsp_list 那条链的
哈希,而影子路径走 find_fast_bsp3 + build_htf_zones + htf_fx_timeline +
attach_htf_context,两条链读的东西不一样。所以新建 verify_lean_parity.py 在
这条路径上逐根对拍 compute() 的每个返回字段。

其中一个坑:随机取窗口测不到信号分支。信号密度约 1/2000 根,头 12 个窗口命中
0 个,「一致」只覆盖了早退路径。改成一半窗口对齐到已知信号根,命中率才上来。
最终 180 窗口 / 两模式各 90 命中 / 零分歧。

生产实测(inner_ms,容器内同口径):

  compute_ms      646 → 132ms   4.89x
  inner_ms        612 → 128ms   4.80x
  lag_signal_ms  1410 → 683ms   完整链路,落回 800ms 线内

比 3.4 倍更好,因为是引擎 ~3.5x 叠 lean ~1.35x。

两点判读上的订正:

- lag_data_ms 那 576→490ms 是噪声,不要记在引擎账上。均值 721±36 vs
  740±127,重叠;而且引擎本来就影响不到交易所与网络那一段。
- 仍有 44% 的根超 800ms,但尾部现在完全由数据腿主导(lag_data P90 1482ms
  vs compute P90 237ms)。计算既不是瓶颈也不是尾部主因了,继续压计算换不到
  尾部改善。800ms 那道闸取的是最近 30 根的中位数,683ms 已满足。

start.sh 加 SHADOW_LEAN(默认 1),设 0 可退回 full 复量两模式差异。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 04:24:08 +08:00

193 lines
8.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""影子交易器的信号函数——在子进程里跑,不碰事件循环。
单次调用约 0.26s 的纯 CPU,且 chanlun 是纯 Python 受 GIL 限制,放进
Hummingbot 的 asyncio 循环里会把行情处理一起卡住,所以必须隔离到独立进程。
## 口径必须与预算同源,缺一项数就不可比
预算(`lib/shadow_budget.BUDGET_BP`)算在 step42 的这套滤网上,
本文件逐行对齐 `step42_exit_tp_1m.run_one`
同向 h1_agree == 1
中枢阶梯 多头要求当前中枢整体高于前一个(zd > 前 zg),空头反之
ATR 门控 atr_pct ≥ ATR_GATE_BP(当前 8bp
早先这里只有 h1_agree。缺阶梯与门控测的就不是我们要交易的那批信号,
而这一项改常数解决不了——必须改信号路径本身。
门控阈值是**费率的函数**不是市场常数(低 ATR 信号的毛质量反而最好,
断崖只在扣费后出现),换 VIP 档或换交易所要重扫,不要抄 8bp。
## 与回测的两点差别
其一,这里只关心**最后一根已收盘 K 线**上有没有信号——实盘只能在当下下单。
其二,`atr_pct` 的分母取次根开盘价,与 `exit_model.walk_exits` 一致,
所以调用方必须把次根开盘价传进来。
未通过滤网的信号也一并返回并打上标志:过滤后样本很稀(门控后 8 个币
合计约 38 笔/周),未过滤的可作提前读数。但**统计主口径只能用 pass_all**
在我们根本不会下单的根上测滑点会把判据算宽。
"""
from __future__ import annotations
import os
import time
import warnings
warnings.filterwarnings("ignore")
for _v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS"):
os.environ.setdefault(_v, "1")
LEAN = os.environ.get("SHADOW_LEAN", "1") not in ("0", "", "false")
def compute(df_l, df_h, entry_px: float | None = None,
lean: bool | None = None) -> dict:
"""在 df_l 的最后一根上找信号。df_l/df_h 都只含已收盘 K 线。
entry_px 是次根开盘价(回测 entry_delay=1 的成交价),用作 atr_pct 的
分母。取不到时退回用信号根收盘价,并在返回里标 atr_ref="close"。
lean=True 让 TF_DF 只构建到中枢,跳过线段/走势中枢/MACD 状态机。本路径
只读 chan.dataframe 与 chan.klc_list,不碰 bsp_list/seg_list/chanmacd
所以可以跳。但静态检查会漏间接依赖,等价性由 verify_lean_parity.py 在
这条路径上逐根实测,不套用 step46 那 5 个对拍用例——那些用例走的是
bsp_list,覆盖不到 fast_bsp3 + 嵌套上下文这条链。
返回 dict
last_idx 最后一根在 chanlun 处理后 dataframe 里的下标
n_bars 实际参与计算的根数
atr_pct 信号根 ATR / 次根开盘价
hits 命中列表,每项含方向与三个滤网标志、pass_all
error 出错时的说明,正常为 None
"""
import numpy as np
import pandas as pd
lean = LEAN if lean is None else lean
try:
from chanlun import TF_DF
from lib.fast_bsp3 import find_fast_bsp3
from lib.fx_signal import extract_fx_signals, signals_to_frame
from lib.nested_bsp import attach_htf_context, htf_fx_timeline
from lib.nested_level import build_htf_zones
from lib.shadow_budget import ATR_GATE_BP
chan_l = TF_DF(df_l, 1, "1m", lean=lean)
cdf = chan_l.dataframe
last = len(cdf) - 1
base = {"last_idx": last, "n_bars": int(len(df_l)), "hits": [],
"atr_pct": None, "atr_ref": None, "error": None}
# ATR 门控。分母与 exit_model.walk_exits 一致,取次根开盘价。
# 放在任何早退之前——无信号的根也要记,才能在线看到门控的真实刷除率
atr = float(cdf["atr"].to_numpy(dtype=float)[last]) \
if "atr" in cdf.columns else float("nan")
ref = entry_px if (entry_px and np.isfinite(entry_px)) else \
float(cdf["close"].to_numpy(dtype=float)[last])
atr_pct = atr / ref if (np.isfinite(atr) and ref) else float("nan")
gate_ok = bool(np.isfinite(atr_pct) and atr_pct * 1e4 >= ATR_GATE_BP)
base["atr_pct"] = None if not np.isfinite(atr_pct) else float(atr_pct)
base["atr_ref"] = "next_open" if (entry_px and np.isfinite(entry_px)) \
else "close"
zones = build_htf_zones(cdf, "1m", chan=chan_l).reset_index(drop=True)
if zones.empty:
return base
# 中枢阶梯:当前中枢是否整体脱离前一个。与 step42 同一算法
z = zones.copy()
prev_zg, prev_zd = z["zg"].shift(), z["zd"].shift()
z["z_above"], z["z_below"] = z["zd"] > prev_zg, z["zg"] < prev_zd
z["zone_i"] = np.arange(len(z))
sig = find_fast_bsp3(cdf, zones)
if sig is None or sig.empty:
return base
sig = sig.merge(z[["zone_i", "z_above", "z_below"]],
on="zone_i", how="left")
# 5m 同向。算不出时 h1_agree 记 0,该信号自然不会通过 pass_all
if df_h is not None and len(df_h) > 0:
chan_h = TF_DF(df_h, 1, "5m", lean=lean)
hdf = chan_h.dataframe
tl = htf_fx_timeline(
signals_to_frame(extract_fx_signals(chan_h, hdf)), hdf)
sig = attach_htf_context(sig, cdf, tl, "h1")
else:
sig["h1_agree"] = 0
cur = sig[sig["entry_idx"].astype(int) == last]
if cur.empty:
return base
hits = []
for _, r in cur.iterrows():
d = int(r["direction"])
push = r["z_above"] if d == 1 else r["z_below"]
ladder_ok = bool(pd.notna(push) and bool(push))
# attach_htf_context 在入场时刻之前没有大级别分型时写 NaN。
# 不能写成 `int(x or 0)`——NaN 是真值,会走到 int(nan) 抛异常,
# 整根的信号就此丢掉,只留一行报错
raw = r.get("h1_agree", 0)
agree = int(raw) if pd.notna(raw) else 0
hits.append({"direction": d, "h1_agree": agree,
"ladder_ok": int(ladder_ok), "gate_ok": int(gate_ok),
"pass_all": int(agree == 1 and ladder_ok and gate_ok)})
base["hits"] = hits
return base
except Exception as e: # 子进程里异常必须带回主进程,否则只见超时不见原因
import traceback
return {"last_idx": -1, "n_bars": int(len(df_l)) if df_l is not None else 0,
"hits": [], "atr_pct": None, "atr_ref": None,
"error": f"{type(e).__name__}: {e}",
"traceback": traceback.format_exc()}
NUM_COLS = ("timestamp", "open", "high", "low", "close", "volume")
def _rebuild(rows) -> "object":
"""只传数值列,date 在这里按 lib/data.py 的同一规则重建。
跨进程传 tz-aware 的 datetime 既慢又容易在字符串往返中丢时区,
而时区若与回测不一致,chanlun 的 K 线标签就会错位。
"""
import pandas as pd
df = pd.DataFrame(rows, columns=list(NUM_COLS))
df["timestamp"] = df["timestamp"].astype("int64")
date = pd.to_datetime(df["timestamp"], unit="ms", utc=True) \
.dt.tz_convert("Asia/Shanghai")
df.insert(1, "date", date)
return df
def compute_packed(payload: tuple) -> dict:
"""ProcessPoolExecutor 的入口:收 (l_rows, h_rows, entry_px[, t_submit])。
返回里带上 `queue_ms` 与 `inner_ms`,把父进程看到的墙钟时间拆开:
compute_ms(父进程测)= queue_ms + 反序列化 + inner_ms + 回传
这个拆分决定「加核有没有用」。排队占大头说明 worker 数不够(币数多于
worker 数时,同一秒收盘的币只能排队),加核直接见效;纯计算占大头说明
单核性能受限,加核帮不上,得从算法上改成增量更新。
两者混在一个数里就只能靠猜。
"""
t_start = time.time()
l_rows, h_rows, entry_px, *rest = payload
t_submit = rest[0] if rest else None
t0 = time.perf_counter()
df_l = _rebuild(l_rows)
df_h = _rebuild(h_rows) if h_rows else None
out = compute(df_l, df_h, entry_px)
out["inner_ms"] = int((time.perf_counter() - t0) * 1000)
# 同一台机器,父子进程时钟一致,可直接相减
out["queue_ms"] = int((t_start - t_submit) * 1000) \
if t_submit is not None else None
return out