补验多 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>
This commit is contained in:
jack
2026-08-28 06:24:51 +08:00
co-authored by Cursor
parent 34a8f37b39
commit 545fcc7c77
+72 -4
View File
@@ -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: