原对拍每根都只追加 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>
223 lines
8.6 KiB
Python
223 lines
8.6 KiB
Python
"""逐根对拍「全量重算」与「增量追加」,并量提速。
|
||
|
||
## 为什么必须逐根对拍,不能引用 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()
|