diff --git a/research/live/verify_incr_parity.py b/research/live/verify_incr_parity.py index 57ab500..94317c9 100644 --- a/research/live/verify_incr_parity.py +++ b/research/live/verify_incr_parity.py @@ -120,8 +120,72 @@ def run_one(sym: str, cache: Path, n: int, start_at: int | None) -> dict: "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) @@ -132,7 +196,10 @@ def main() -> None: for sym in a.syms.split(","): print(f"\n{'=' * 70}\n{sym}") try: - rows.append(run_one(sym, Path(a.cache), a.n, a.start)) + 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: @@ -141,9 +208,10 @@ def main() -> None: print(f"\n\n{'=' * 70}\n汇总\n") print(f" 对拍 {int(d['n'].sum()):,} 根 · 不一致 {int(d['diff'].sum())} · " f"命中 {int(d['hits'].sum())} 个") - 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 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: