"""影子交易器的信号函数——在子进程里跑,不碰事件循环。 单次调用约 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") INCR = os.environ.get("SHADOW_INCR", "1") not in ("0", "", "false") # 增量流缓存。worker 进程被复用,所以这个 dict 跨根存活。 # 键是 (symbol, timeframe)——2 个 worker 轮流拿 10 个币,每个 worker 最终会 # 缓存全部 10 个币,共 20 条流。 _STREAMS: dict[tuple, tuple] = {} # 两次重建之间允许窗口长多少根。 # # init_stream/append_bar **没有 trim**:dataframe 靠 pd.concat 无界增长。所以 # 增量必然让窗口每根 +1,只能周期性 init_stream 拉回。取 500 的两个理由: # 1. append_bar 里 rebuild_bi_zs 要整表重扫笔,是 O(n)。窗口涨 25% 成本也涨 # 约 25%,500/2001 正好把这个膨胀压在 25% 以内。 # 2. 重建约 51ms、追加约 14ms,摊到 500 根上重建只加 0.07ms/根。 # 前提「输出对窗口长度不敏感」由 verify_window_sens.py 验过(+200/+500/+1000 # 全部逐字段一致),否则这个方案等于静默换掉一批信号。 MAX_GROW = 500 def _chan_for(key: tuple, df, tf: str, lean: bool): """拿该窗口对应的 chan 对象,能增量就增量,否则重建。 三种情况必须回退到全量重建,否则会拿一个状态不对的流去出信号: 缓存没有 首次见到这个币 窗口已长过阈值 见 MAX_GROW 缓存末根不在新窗口 说明中间断了很多根(或时间戳回退),接不上 第三种是最要紧的。2 个 worker 轮流拿 10 个币,某个 worker 可能隔几根才再 看到同一个币,那几根要补齐;但若缺口大到超出窗口,就没法补,只能重建。 不检查而直接 append 会把不连续的 K 线接在一起,笔和中枢全错且不报错。 """ from chanlun import TF_DF ts = df["timestamp"].to_numpy("int64") st = _STREAMS.get(key) if st is not None: chan, last_ts, base_n = st if len(chan.dataframe) <= base_n + MAX_GROW and last_ts >= ts[0] \ and last_ts <= ts[-1] and (ts == last_ts).any(): for _, row in df[df["timestamp"] > last_ts].iterrows(): chan.append_bar(row) _STREAMS[key] = (chan, int(ts[-1]), base_n) return chan # 重建走**批量** init_TF_DF,不用 init_stream。init_stream 是逐行 # `dataframe.iloc[idx]`,正是引擎提速刚修掉的反模式:实测 2001 根要 # 238.5ms,而批量 lean 只要 74.3ms,慢 3.2 倍。 # append_bar 能接在批量构建的对象上——_ensure_stream_state 会补出 # _klc_feed_last_klu,其余列表 init_TF_DF 都建好了。 chan = TF_DF(df.copy(), 1, tf, lean=lean) _STREAMS[key] = (chan, int(ts[-1]), len(df)) return chan def compute(df_l, df_h, entry_px: float | None = None, lean: bool | None = None, sym: str | None = None, incr: 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 # 没有 sym 就无法给流分键,只能走全量——对拍脚本会用这条路径当基准 incr = (INCR if incr is None else incr) and sym is not None 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 = _chan_for((sym, "1m"), df_l, "1m", lean) if incr \ else 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 = _chan_for((sym, "5m"), df_h, "5m", lean) if incr \ else 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 sym = rest[1] if len(rest) > 1 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, sym=sym) # 落盘这两个数才能在线看出增量是否在生效:走了重建的根 grown 会等于窗口 st = _STREAMS.get((sym, "1m")) out["stream_bars"] = int(len(st[0].dataframe)) if st else None 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