Files
Chan/research/live/verify_incr_parity.py
jackandCursor 545fcc7c77 补验多 worker 交错下的一次追加多根
原对拍每根都只追加 1 根,覆盖不到多 worker 交错的实际路径:
ProcessPoolExecutor 不保证同一个币落到同一个 worker,所以每个 worker 隔 nw
根才再见到这个币,一次要补 nw 根。「补 nw 根等价于连续追加 nw 次」是推理,
没实测过。

--interleave N 用 N 份独立缓存轮流接同一个币。BTC/SOL 各 300 根、2 份缓存,
逐字段零分歧,各缓存末窗 2299 根。

顺带记清亲和性的性质:缓存是 worker 进程内的 dict、键含 symbol,不存在
「worker A 的状态被 B 读到」或「拿到别的币的状态」。亲和性影响的是内存
(每个 worker 最终缓存全部币)与补根次数,不影响正确性。实测内存
597→598MiB,代价在噪声里。

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

223 lines
8.6 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.
"""逐根对拍「全量重算」与「增量追加」,并量提速。
## 为什么必须逐根对拍,不能引用 HANDOFF §5.5
§5.5 验的是 step46 那批用例(走 bsp_list 那条链),且是「追加 150~200 根 vs
全量重建」的整体哈希。影子路径不同:
- 走 find_fast_bsp3 + build_htf_zones + htf_fx_timeline + attach_htf_context
- 流式对象**跨根复用**,而 worker 轮流拿多个币,同一条流可能隔几根才被
再次追加。状态污染只会让信号悄悄换一批,不报错、不崩
而且代码阅读已经暴露一处偏差:`init_stream/append_bar` 从不调 `cal_trend`
(它只在 `get_klc_list` 里),所以增量路径下 `klc.trend` 恒为 UNKNOWN。
HANDOFF 说「笔的计算依赖 klc.trend」——若为真,增量的笔就和全量不同。
那句话所引的 bi.py:221 其实在 `cal_trend` 自己的循环里,不是 `cal_bi_list`
的依赖。**这条只能由对拍来定论**,不能靠读代码。
## 判据
逐字段相同,排除 last_idx/n_bars(随窗口长度必然变,见 verify_window_sens
与计时字段。数量相同而标志不同一样算失败。
模拟真实调用模式:连续推进,且每根都按「全量」和「增量」各算一次,增量那侧
复用同一条流。
python research/live/verify_incr_parity.py --syms BTC,ETH,SOL --n 300
"""
from __future__ import annotations
import argparse
import os
import sys
import time
import warnings
from pathlib import Path
import numpy as np
import pandas as pd
warnings.filterwarnings("ignore")
for _v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS"):
os.environ.setdefault(_v, "1")
HERE = Path(__file__).resolve()
sys.path.insert(0, str(HERE.parents[1]))
sys.path.insert(0, str(HERE.parents[2]))
sys.path.insert(0, str(HERE.parent))
BASE_L, BASE_H = 2001, 801
def run_one(sym: str, cache: Path, n: int, start_at: int | None) -> dict:
import shadow_signal as ss
from verify_lean_parity import SKIP, WINDOW_KEYS, canon, load, signal_bars
skip = SKIP + WINDOW_KEYS + ("stream_bars",)
l_all, h_all = load(sym, "1m", cache), load(sym, "5m", cache)
l_ts = l_all["timestamp"].to_numpy("int64")
h_ts = h_all["timestamp"].to_numpy("int64")
# 从最后一个信号根往前 n 根开始,保证这段里一定有信号分支被执行
if start_at is None:
try:
sb = signal_bars(sym, cache)
sb = sb[(sb > BASE_L + n) & (sb < len(l_all) - 1)]
start_at = int(sb[-1]) - n + 5 if len(sb) else BASE_L + 10
except Exception:
start_at = BASE_L + 10
ends = [e for e in range(start_at, start_at + n) if e < len(l_all) - 1]
if not ends:
raise RuntimeError("窗口不足")
ss._STREAMS.clear()
same = diff = 0
t_full = t_incr = 0.0
n_hits = 0
first = None
rebuilds = 0
prev_grown = 0
for e in ends:
hi = int(np.searchsorted(h_ts, l_ts[e], side="right"))
df_l = l_all.iloc[e - BASE_L + 1:e + 1]
df_h = h_all.iloc[max(0, hi - BASE_H):hi]
entry = float(l_all["open"].to_numpy(float)[e + 1])
t0 = time.perf_counter()
rf = ss.compute(df_l.copy(), df_h.copy(), entry, incr=False)
t_full += time.perf_counter() - t0
t0 = time.perf_counter()
ri = ss.compute(df_l.copy(), df_h.copy(), entry, sym=sym, incr=True)
t_incr += time.perf_counter() - t0
grown = len(ss._STREAMS[(sym, "1m")][0].dataframe)
if grown <= prev_grown:
rebuilds += 1
prev_grown = grown
n_hits += len(rf.get("hits") or [])
if canon(rf, skip) == canon(ri, skip):
same += 1
else:
diff += 1
if first is None:
first = (e, canon(rf, skip), canon(ri, skip))
k = len(ends)
print(f" {k} 根 · 一致 {same} · 不一致 {diff} · 命中 {n_hits} 个 · "
f"重建 {rebuilds} 次 · 末窗 {prev_grown} 根")
print(f" 单根 全量 {t_full / k * 1000:.1f}ms → "
f"增量 {t_incr / k * 1000:.1f}ms "
f"{t_full / max(t_incr, 1e-9):.2f}x")
if first:
e, a, b = first
print(f" ⚠ 首个分歧 idx={e}\n 全量: {a[:300]}\n 增量: {b[:300]}")
ss._STREAMS.clear()
del l_all, h_all
return {"sym": sym, "n": k, "same": same, "diff": diff, "hits": n_hits,
"full_ms": t_full / k * 1000, "incr_ms": t_incr / k * 1000}
def interleave(sym: str, cache: Path, n: int, nw: int) -> dict:
"""模拟多 worker 交错:nw 份独立缓存轮流接同一个币。
这是单进程对拍覆盖不到的路径。`ProcessPoolExecutor` 不保证同一个币落到
同一个 worker,所以每个 worker 只能隔 nw 根才再见到这个币,一次要补 nw
根。补根走的是 `for row in df[ts > last_ts]` 那个循环——逻辑上等价于连续
追加 nw 次,但「等价」是推理,没实测过。
缓存是 worker 进程内的 dict、键含 symbol,所以不存在「worker A 的状态被
worker B 读到」或「拿到别的币的状态」。亲和性影响的是内存(每个 worker
最终缓存全部币)与补根次数,不影响正确性——本函数就是来证这一点的。
"""
import shadow_signal as ss
from verify_lean_parity import SKIP, WINDOW_KEYS, canon, load, signal_bars
skip = SKIP + WINDOW_KEYS + ("stream_bars",)
l_all, h_all = load(sym, "1m", cache), load(sym, "5m", cache)
l_ts = l_all["timestamp"].to_numpy("int64")
h_ts = h_all["timestamp"].to_numpy("int64")
try:
sb = signal_bars(sym, cache)
sb = sb[(sb > BASE_L + n) & (sb < len(l_all) - 1)]
start = int(sb[-1]) - n + 5 if len(sb) else BASE_L + 10
except Exception:
start = BASE_L + 10
ends = [e for e in range(start, start + n) if e < len(l_all) - 1]
caches: list[dict] = [{} for _ in range(nw)]
same = diff = n_hits = 0
first = None
for j, e in enumerate(ends):
hi = int(np.searchsorted(h_ts, l_ts[e], side="right"))
df_l = l_all.iloc[e - BASE_L + 1:e + 1]
df_h = h_all.iloc[max(0, hi - BASE_H):hi]
entry = float(l_all["open"].to_numpy(float)[e + 1])
rf = ss.compute(df_l.copy(), df_h.copy(), entry, incr=False)
# 轮流换缓存 = 轮流换 worker
ss._STREAMS = caches[j % nw]
ri = ss.compute(df_l.copy(), df_h.copy(), entry, sym=sym, incr=True)
n_hits += len(rf.get("hits") or [])
if canon(rf, skip) == canon(ri, skip):
same += 1
else:
diff += 1
if first is None:
first = (e, canon(rf, skip), canon(ri, skip))
grown = [len(c[(sym, "1m")][0].dataframe) for c in caches
if (sym, "1m") in c]
print(f" {len(ends)} 根 · {nw} 份缓存轮流 · 一致 {same} · 不一致 {diff}"
f" · 命中 {n_hits} 个 · 各缓存末窗 {grown}")
if first:
e, a, b = first
print(f" ⚠ 首个分歧 idx={e}\n 全量: {a[:300]}\n 增量: {b[:300]}")
ss._STREAMS = {}
del l_all, h_all
return {"sym": sym, "n": len(ends), "same": same, "diff": diff,
"hits": n_hits, "full_ms": 0.0, "incr_ms": 0.0}
def main() -> None:
ap = argparse.ArgumentParser()
ap.add_argument("--interleave", type=int, default=0,
help="模拟这么多个 worker 轮流接同一个币(一次补多根)")
ap.add_argument("--syms", default="BTC,ETH,SOL")
ap.add_argument("--cache", default="research/live/cache")
ap.add_argument("--n", type=int, default=300)
ap.add_argument("--start", type=int, default=None)
a = ap.parse_args()
rows = []
for sym in a.syms.split(","):
print(f"\n{'=' * 70}\n{sym}")
try:
if a.interleave:
rows.append(interleave(sym, Path(a.cache), a.n, a.interleave))
else:
rows.append(run_one(sym, Path(a.cache), a.n, a.start))
except Exception as e:
print(f" 跳过:{e!r}")
if not rows:
return
d = pd.DataFrame(rows)
print(f"\n\n{'=' * 70}\n汇总\n")
print(f" 对拍 {int(d['n'].sum()):,} 根 · 不一致 {int(d['diff'].sum())} · "
f"命中 {int(d['hits'].sum())} 个")
if d["incr_ms"].sum() > 0:
print(f" 单根 全量 {d['full_ms'].mean():.1f}ms → "
f"增量 {d['incr_ms'].mean():.1f}ms "
f"{d['full_ms'].sum() / max(d['incr_ms'].sum(), 1e-9):.2f}x")
if int(d["diff"].sum()) == 0:
print("\n 逐字段一致,增量可以上线。")
else:
print("\n ⛔ 有分歧,不要上线。增量流的状态与全量重建不等价。")
if __name__ == "__main__":
main()