"""前置测量一:本机算力——影子交易器的计算延迟下限。 step39 在 Mac 上量到 2000 根窗口单次 0.20s。这台机器只有 2 vCPU 且单核更慢, 而计算延迟直接吃 1m 腿仅 3.9bp 的滑点预算,所以必须在写影子交易器之前实测。 量四件事: 1m 侧 TF_DF(2000根) + build_htf_zones + find_fast_bsp3 5m 侧 TF_DF(800根) + 分型 + 时间线(只在 5m 收盘那根变,可摊薄到 1/5) 串行 3 币 最后一个币要等多久 —— 这就是它的下单延迟 并行 3 币 2 核跑 3 进程会互相抢核,未必比串行快 输出 out/bench_compute.csv,判据是与 3.9bp 预算对应的漂移:本机 1m 波动 BTC 6.5bp/分钟,按 √t 折算,1 秒延迟约 0.8bp、5 秒约 1.9bp。 """ from __future__ import annotations import argparse import os import sys import time import warnings from concurrent.futures import ProcessPoolExecutor 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().parent RESEARCH = HERE.parent sys.path.insert(0, str(RESEARCH)) sys.path.insert(0, str(RESEARCH.parent)) pd.set_option("display.width", 240) LTF, HTF = "1m", "5m" WIN_LTF = 2000 WIN_HTF = max(WIN_LTF // 5, 800) SYMS = ("BTC", "ETH", "SOL") def _signal_once(sl_ltf: pd.DataFrame) -> tuple[int, float]: """1m 侧:切片重建结构与信号。与 step39 的 _pit_once 同一套调用。""" from chanlun import TF_DF from lib.fast_bsp3 import find_fast_bsp3 from lib.nested_level import build_htf_zones t0 = time.perf_counter() chan = TF_DF(sl_ltf, 1, LTF) cdf = chan.dataframe zones = build_htf_zones(cdf, LTF, chan=chan) n_sig = 0 if not zones.empty: sig = find_fast_bsp3(cdf, zones.reset_index(drop=True)) n_sig = len(sig) return n_sig, time.perf_counter() - t0 def _agree_once(sl_htf: pd.DataFrame) -> tuple[int, float]: """5m 侧:切片重建分型时间线。与 step39 的 _pit_agree 同一套调用。""" from chanlun import TF_DF from lib.fx_signal import extract_fx_signals, signals_to_frame from lib.nested_bsp import htf_fx_timeline t0 = time.perf_counter() chan = TF_DF(sl_htf, 1, HTF) tl = htf_fx_timeline(signals_to_frame(extract_fx_signals(chan, chan.dataframe)), chan.dataframe) return len(tl), time.perf_counter() - t0 def _load(sym: str) -> tuple[pd.DataFrame, pd.DataFrame]: from lib.data import load_local pair = f"{sym}/USDT:USDT" df_l = load_local(pair, LTF) df_h = load_local(pair, HTF) if df_l is None or df_h is None: raise SystemExit(f"{sym} 本地数据缺失,本机只有 BTC/ETH/SOL") return df_l, df_h def _slices(sym: str, reps: int, seed: int) -> list[tuple[pd.DataFrame, pd.DataFrame]]: """预先切好窗口。读 feather 是本机 IO,实盘数据来自 WS 缓冲,不计入延迟。""" df_l, df_h = _load(sym) rng = np.random.default_rng(seed) # 从末段随机取窗口,避开数据头部(指标预热)与尾部(不足一窗) lo, hi = max(WIN_LTF, len(df_l) - 200_000), len(df_l) - 1 ltf_ts = df_l["timestamp"].to_numpy() htf_ts = df_h["timestamp"].to_numpy() out = [] for i in rng.integers(lo, hi, size=reps): i = int(i) sl_l = df_l.iloc[i - WIN_LTF + 1: i + 1].reset_index(drop=True) # 5m 只喂到不晚于该 1m 根的部分,与实盘一致 h_end = int(np.searchsorted(htf_ts, ltf_ts[i], side="right")) sl_h = (df_h.iloc[h_end - WIN_HTF: h_end].reset_index(drop=True) if h_end >= WIN_HTF else None) out.append((sl_l, sl_h)) return out def bench_pair(task: tuple) -> dict: """算一根:1m 信号 + 5m 时间线。在子进程里跑,import 都在函数内。""" import warnings as _w _w.filterwarnings("ignore") sys.path.insert(0, str(RESEARCH)) sys.path.insert(0, str(RESEARCH.parent)) sym, sl_l, sl_h = task _, dt_l = _signal_once(sl_l) dt_h = np.nan if sl_h is None else _agree_once(sl_h)[1] return {"sym": sym, "ltf_s": dt_l, "htf_s": dt_h} def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--reps", type=int, default=15, help="每币采样窗口数") ap.add_argument("--rounds", type=int, default=5, help="串行/并行各跑几轮") args = ap.parse_args() print(f"[本机算力] {os.cpu_count()} 逻辑核 · 1m 窗口 {WIN_LTF} 根 · " f"5m 窗口 {WIN_HTF} 根 · 每币 {args.reps} 次\n", flush=True) print("########## 1. 单币单次耗时 ##########", flush=True) prepared = {s: _slices(s, args.reps, 7) for s in SYMS} per = [] for s in SYMS: d = pd.DataFrame([bench_pair((s, l, h)) for l, h in prepared[s]]) per.append(d) print(f" {s}: 1m {d['ltf_s'].median():.3f}s (P95 {d['ltf_s'].quantile(.95):.3f}s)" f" · 5m {d['htf_s'].median():.3f}s", flush=True) allp = pd.concat(per, ignore_index=True) m_ltf = allp["ltf_s"].median() m_htf = allp["htf_s"].median() # 5m 侧只在每 5 根 1m 里变一次,其余 4 根可复用缓存 amort = m_ltf + m_htf / 5 print(f"\n 三币合并中位:1m {m_ltf:.3f}s · 5m {m_htf:.3f}s") print(f" 单币每根摊薄成本 {amort:.3f}s(5m 每 5 根才重算一次)") print(f" 对比 Mac 的 0.20s:本机慢 {m_ltf / 0.20:.1f} 倍") print("\n########## 2. 三币串行 vs 并行(最后一个币的下单延迟)##########", flush=True) print(" 只计算子耗时:读数据不计入,进程池常驻不计启动开销", flush=True) ser, par2, par3 = [], [], [] pools = {w: ProcessPoolExecutor(max_workers=w) for w in (2, 3)} try: # 先各跑一次把子进程的 import 预热掉,否则首轮全是模块加载时间 for w, ex in pools.items(): list(ex.map(bench_pair, [(s, *prepared[s][0]) for s in SYMS])) for k in range(args.rounds): batch = [(s, *prepared[s][k % args.reps]) for s in SYMS] t0 = time.perf_counter() for t in batch: bench_pair(t) ser.append(time.perf_counter() - t0) for w, bag in ((2, par2), (3, par3)): t0 = time.perf_counter() list(pools[w].map(bench_pair, batch)) bag.append(time.perf_counter() - t0) print(f" 轮 {k + 1}: 串行 {ser[-1]:.2f}s · 并行2 {par2[-1]:.2f}s · " f"并行3 {par3[-1]:.2f}s", flush=True) finally: for ex in pools.values(): ex.shutdown() print("\n########## 3. 延迟折算成 bp(3.9bp 预算的参照)##########") # 各币 1m 收益标准差,用 √t 把延迟折成价格漂移的一个标准差 vol = {} for s in SYMS: df_l, _ = _load(s) r = np.diff(np.log(df_l["close"].to_numpy(dtype=float)[-400_000:])) vol[s] = float(np.std(r) * 1e4) rows = [] for name, xs in (("串行", ser), ("并行2", par2), ("并行3", par3)): d = float(np.median(xs)) rows.append({"方案": name, "三币总耗时": f"{d:.2f}s", **{f"{s} 漂移1σ": f"{vol[s] * np.sqrt(d / 60):.2f}bp" for s in SYMS}}) tb = pd.DataFrame(rows) print(tb.to_string(index=False)) print(f"\n 1m 波动实测:" + " · ".join(f"{s} {vol[s]:.1f}bp/分钟" for s in SYMS)) print(" 漂移 1σ 是随机部分;入场时价格正朝信号方向跑,系统性追价另计。") out = RESEARCH / "out" / "bench_compute.csv" pd.DataFrame({ "指标": ["cpu核数", "1m窗口", "5m窗口", "1m中位s", "1m_P95s", "5m中位s", "单币摊薄s", "串行3币s", "并行2_3币s", "并行3_3币s"], "值": [os.cpu_count(), WIN_LTF, WIN_HTF, round(m_ltf, 3), round(allp["ltf_s"].quantile(.95), 3), round(m_htf, 3), round(amort, 3), round(float(np.median(ser)), 2), round(float(np.median(par2)), 2), round(float(np.median(par3)), 2)], }).to_csv(out, index=False) print(f"\n产物写入 {out}") if __name__ == "__main__": main()