影子信号走增量:清空十币 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>
This commit is contained in:
jack
2026-08-28 06:17:47 +08:00
co-authored by Cursor
parent c8a0062707
commit 0a0fd2f682
5 changed files with 349 additions and 9 deletions
+66 -4
View File
@@ -40,10 +40,64 @@ for _v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS"):
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) -> dict:
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 的
@@ -66,6 +120,8 @@ def compute(df_l, df_h, entry_px: float | None = None,
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
@@ -75,7 +131,8 @@ def compute(df_l, df_h, entry_px: float | None = None,
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)
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": [],
@@ -111,7 +168,8 @@ def compute(df_l, df_h, entry_px: float | None = None,
# 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)
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)
@@ -180,11 +238,15 @@ def compute_packed(payload: tuple) -> dict:
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)
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) \