1m 腿的滑点余量只有几个 bp,所以要测的必须是生产路径的滑点——换个运行时 测出来的数就不作数。框架因此从「滑点已知后再定」提前到测量阶段就定为 Hummingbot(Spot/Perp 连接器均 v2.0,Bitget 是 Foundation Partner)。 新增 research/live/。前置测量: - bench_compute.py 本机算力,1m 单币 0.318s、三币串行 1.38s - venue_parity.py Binance 与 Bitget 同根信号重合仅 14.6~42.6% - signal_sensitivity.py 0.25bp 扰动就换掉一半信号 - aggregate_robustness.py 但总体期望不降——脆的是信号身份,不是 alpha - bitget_baseline.py 因此改用 Bitget 原生基线定预算:余量 BTC -0.13bp、 ETH +4.02bp、SOL +2.92bp。BTC 本就为负,只作延迟测量的参照物 运行时选型: - parity_env.py 容器与本机信号逐一相同(下标、中枢数、checksum 全等), 容器内 0.26s/币反而更快。故 chanlun 直接挂载进容器,不必另起信号服务。 装进现有 .venv 那条路走不通:Hummingbot 要 numba>=0.61.2 与 aiohttp<3.14,与本机 Python 3.14 冲突 - latency_ccxt.py / latency_hummingbot.py / latency_compare.py 初测显示 Hummingbot 比 ccxt.pro 慢约 1030ms,90 根逐根配对里 80~97% 更慢 - probe_ws_action.py 否掉「丢弃 snapshot」的猜测:换根首条就是 update - probe_hb_vs_raw.py 与 latency_attribute.py 四路归因——容器网络 2~18ms、 Hummingbot 处理 -10~-30ms,1350~1480ms 全落在解析方式上 - probe_ws_payload.py 定位根因:Bitget 换根会推一条带两根的消息 [上一根, 新一根],而上游取 data["data"][0] 拿到的是上一根,新一根要等 下一条单元素消息 修复: - patched_candles.py 处理消息里的全部元素。不能简单改成 [-1]——那样上一根 的收盘价会永远停在换根前约 1 秒的那次推送上,而信号对 0.25bp 都敏感 - verify_patch.py 60 根配对验证:拿回 1060~1090ms,与原始 WS 只差 5~14ms 已贴理论下限,19 根已收盘 K 线 OHLCV 逐根未变。折算 ETH 省 0.54bp、 SOL 省 0.42bp。此 bug 值得向上游反馈 影子交易器: - shadow_hb.py 不下单,读连接器真实盘口按仓位吃单深度算成交价,与次根开盘价 (回测 entry_delay=1 的口径)相减,分解成延迟漂移、盘口价差、深度冲击。 盘口 10Hz 滚动缓冲 30 秒,把延迟变成自变量:每个信号记 0.5/1/2/5s 与实际 算完时刻各一个滑点值,本机算得慢也不影响能读出的曲线 - shadow_signal.py 信号计算隔离到子进程。0.26s 是纯 CPU 且 chanlun 受 GIL 限制,放进 asyncio 循环会把行情处理一起卡住 - shadow_report.py 首日延迟门槛与滑点曲线报表 不用 paper trade 测滑点:它的成交由 Hummingbot 自己的撮合模型模拟, 测出来是模型行为而非市场行为。 Co-authored-by: Cursor <cursoragent@cursor.com>
205 lines
8.2 KiB
Python
205 lines
8.2 KiB
Python
"""前置测量一:本机算力——影子交易器的计算延迟下限。
|
||
|
||
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()
|