Files
jackandCursor 0a0fd2f682 影子信号走增量:清空十币 560ms → 247ms
compute() 每根新建 TF_DF 换成按 (symbol, timeframe) 缓存的流式对象。worker
进程被复用,所以缓存跨根存活;2 worker 轮流拿 10 币,每个 worker 最终缓存
全部 20 条流,实测内存开销落在噪声里(597→598MiB)。

## 前提先验,否则整个改动建立在沙子上

init_stream/append_bar **没有 trim**,dataframe 靠 pd.concat 无界增长。所以
增量必然让窗口每根 +1,只能周期性重建拉回,两次重建之间窗口是 [W, W+500]
而非恒定 W。于是必须先证明 compute() 输出对窗口长度不敏感——否则增量等于
静默换掉一批信号,不报错不崩。

verify_window_sens.py:三币 75 个信号窗口,+200/+500/+1000 三档全部逐字段
一致。step39 说的是「命中率在 2000 根饱和」,饱和不等于不变,这是两回事。

## 对拍

verify_incr_parity.py:三币 1,800 根、21 个命中、各跨 1 次重建边界,逐字段
零分歧。不能引用 HANDOFF §5.5——那验的是 bsp_list 那条链的整体哈希,而这里
是 find_fast_bsp3 那条链,且流式对象跨根复用,状态污染只会让信号悄悄换一批。

对拍顺带定论一件读代码定不了的事:cal_bi_list **不依赖** klc.trend。
init_stream/append_bar 从不调 cal_trend(它只在 get_klc_list 里),所以追加
出来的 klc 其 trend 恒为 UNKNOWN,而批量构建的有值;两者结果逐字段相同。
HANDOFF §5.5 那句「bi.py:221 读 klc.trend,笔的计算依赖它」不成立——221 行
在 cal_trend 自己的循环里,读的是它自身的序列状态。

## 重建不走 init_stream

init_stream 是逐行 dataframe.iloc[idx],正是引擎提速刚修掉的反模式:2001 根
要 238.5ms,而批量 lean 只 74.3ms,慢 3.2 倍。第一版用它重建,10 个币启动时
各来一次,清空反而涨到 1686ms。改用 TF_DF(df, lean=True) 重建,append_bar
靠 _ensure_stream_state 就能接上。

## 实测

append_bar 21.8ms vs 批量 lean 重建 77.7ms = 3.56x,与研究侧测的 3.7x 一致。
拆解:add_indicators 全表 7.5ms(34%,为加一根重算 2001 行)+ cal_bi_list
整表重扫 11.1ms(51%)+ concat 1.4ms。这两项都在引擎侧,值得反馈。

十币 / 2 核:清空 560→247ms,排队 92→10ms,纯计算 219→108ms。判定从
「加 worker 无用,唯一出路是增量」变成「宽裕,无需优化」。

注意 inner 108ms 里 chan 构建只占约 22ms,其余是 build_htf_zones /
find_fast_bsp3 / attach_htf_context。**瓶颈已不在 chan 构建**,再压增量收益
有限。

stream_bars 落到 latency CSV:恒等于 2001 说明缺口判定在每根都回退重建、
增量静默失效,这一点从耗时上看不出是哪一环。实测窗口稳定长大。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 06:17:47 +08:00

255 lines
12 KiB
Python
Raw Permalink 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")
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