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>
307 lines
12 KiB
Python
307 lines
12 KiB
Python
"""前置测量二:venue 对齐——Bitget 与 Binance 的 1m 是不是同一批信号。
|
||
|
||
研究数据全部来自 Binance,影子交易器却跑在 Bitget。若两家的 1m K 线有差异,
|
||
信号集就会不同,而这个差异会被误记到滑点账上——那样收集一两周也不可归因。
|
||
|
||
所以先把两家同期的 1m 拉齐,跑同一套管线,比三件事:
|
||
K 线层 时间戳缺口、close 价差(bp)、high/low 差异
|
||
信号层 原始 fast_bsp3 的重合率
|
||
过滤后 加 5m 同向过滤后的重合率(这才是实际要交易的那批)
|
||
|
||
重合率高 → 后面测到的滑点可以直接对照 Binance 回测的 3.9bp 预算。
|
||
重合率低 → 必须先补 Bitget 自己的回测基线,否则实验不可归因。
|
||
|
||
输出 out/venue_parity.csv。Bitget 数据缓存在 live/cache/,重跑不必再拉。
|
||
"""
|
||
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().parent
|
||
RESEARCH = HERE.parent
|
||
sys.path.insert(0, str(RESEARCH))
|
||
sys.path.insert(0, str(RESEARCH.parent))
|
||
pd.set_option("display.width", 240)
|
||
|
||
CACHE = HERE / "cache"
|
||
LTF, HTF = "1m", "5m"
|
||
HTF_RATIO = 5
|
||
SYMS = ("BTC", "ETH", "SOL")
|
||
NUMERIC = ("open", "high", "low", "close", "volume")
|
||
|
||
|
||
def _exchange():
|
||
import ccxt
|
||
# 这台机器在新加坡,直连 Bitget 0.30s。绝不要照抄交接文档里 Mac 的代理配置,
|
||
# 代理会把延迟放大到秒级,测出来的滑点就是代理的账。
|
||
# rateLimit 默认 50ms,连拉上千页 history-candles 会被 429,放宽到 120ms
|
||
return ccxt.bitget({"options": {"defaultType": "swap"},
|
||
"enableRateLimit": True, "rateLimit": 120})
|
||
|
||
|
||
def fetch_bitget(sym: str, tf: str, days: int, refresh: bool = False) -> pd.DataFrame:
|
||
"""分页拉 Bitget 永续 K 线,落盘缓存,列结构对齐 lib/data.py。"""
|
||
CACHE.mkdir(parents=True, exist_ok=True)
|
||
path = CACHE / f"bitget_{sym}_{tf}_{days}d.feather"
|
||
if path.exists() and not refresh:
|
||
return pd.read_feather(path)
|
||
|
||
ex = _exchange()
|
||
pair = f"{sym}/USDT:USDT"
|
||
period_ms = ex.parse_timeframe(tf) * 1000
|
||
t0 = time.perf_counter()
|
||
calls = 0
|
||
|
||
# 远端 history-candles 每页硬上限 200 根,而 ccxt 会按 limit 推算 endTime,
|
||
# 只把窗口末尾的 200 根还给你。若照 limit=1000 步进,每页就白丢 800 根——
|
||
# 210 天曾因此只拿到应有量的 31%。故分页一律按 200 走。
|
||
PAGE = 200
|
||
|
||
def one(since: int) -> list:
|
||
"""单次取数并退避重试。history-candles 连拉上千次会触发 429。"""
|
||
for attempt in range(6):
|
||
try:
|
||
return ex.fetch_ohlcv(pair, tf, since=since, limit=PAGE)
|
||
except Exception as e:
|
||
if attempt == 5:
|
||
raise
|
||
wait = 2 ** attempt
|
||
print(f" {sym} {tf}: {type(e).__name__},{wait}s 后重试",
|
||
flush=True)
|
||
time.sleep(wait)
|
||
return []
|
||
|
||
def page(start: int, stop: int) -> list:
|
||
"""向前分页。Bitget 把 since 当开区间,故每次从上一批最后一根重取,
|
||
边界少的那一根靠去重消化。"""
|
||
nonlocal calls
|
||
got, since = [], start
|
||
while since < stop:
|
||
batch = one(since)
|
||
calls += 1
|
||
if len(batch) < 2:
|
||
break
|
||
got.extend(batch)
|
||
if batch[-1][0] <= since:
|
||
break
|
||
since = batch[-1][0]
|
||
if calls % 200 == 0:
|
||
print(f" {sym} {tf}: {len(got)} 根 / {calls} 次请求", flush=True)
|
||
return got
|
||
|
||
now = ex.milliseconds()
|
||
rows = page(now - days * 86_400_000, now)
|
||
|
||
# 补缺口:远端接口一次只给 200 根,个别区段仍可能漏,逐个补到补不动为止
|
||
for _ in range(5):
|
||
ts = np.unique(np.array([r[0] for r in rows], dtype="int64"))
|
||
if len(ts) < 2:
|
||
break
|
||
holes = np.where(np.diff(ts) > period_ms)[0]
|
||
if not len(holes):
|
||
break
|
||
before = len(ts)
|
||
for i in holes:
|
||
rows.extend(page(int(ts[i]), int(ts[i + 1])))
|
||
if len(np.unique([r[0] for r in rows])) <= before:
|
||
break
|
||
|
||
df = pd.DataFrame(rows, columns=["timestamp", *NUMERIC])
|
||
df["timestamp"] = df["timestamp"].astype("int64")
|
||
for c in NUMERIC:
|
||
df[c] = pd.to_numeric(df[c], errors="coerce")
|
||
df = (df.dropna(subset=list(NUMERIC))
|
||
.drop_duplicates(subset=["timestamp"])
|
||
.sort_values("timestamp")
|
||
.reset_index(drop=True))
|
||
df["date"] = (pd.to_datetime(df["timestamp"], unit="ms", utc=True)
|
||
.dt.tz_convert("Asia/Shanghai"))
|
||
df = df[["timestamp", "date", *NUMERIC]]
|
||
gap = int(((np.diff(df["timestamp"].to_numpy()) // period_ms) - 1).clip(0).sum())
|
||
print(f" {sym} {tf}: {len(df)} 根,{calls} 次请求,"
|
||
f"{time.perf_counter() - t0:.1f}s,残余缺口 {gap} 根", flush=True)
|
||
df.to_feather(path)
|
||
return df
|
||
|
||
|
||
def pipeline(df_l: pd.DataFrame, df_h: pd.DataFrame) -> pd.DataFrame:
|
||
"""全量口径跑一遍:原始信号 + 5m 同向过滤,返回带时间戳的信号表。"""
|
||
from chanlun import TF_DF
|
||
from lib.fast_bsp3 import find_fast_bsp3
|
||
from lib.fx_signal import extract_fx_signals, signals_to_frame
|
||
from lib.nested_bsp import attach_htf_context, htf_fx_timeline
|
||
from lib.nested_level import build_htf_zones
|
||
|
||
chan_l = TF_DF(df_l, 1, LTF)
|
||
cdf = chan_l.dataframe
|
||
zones = build_htf_zones(cdf, LTF, chan=chan_l)
|
||
if zones.empty:
|
||
return pd.DataFrame()
|
||
sig = find_fast_bsp3(cdf, zones.reset_index(drop=True))
|
||
if sig.empty:
|
||
return pd.DataFrame()
|
||
|
||
chan_h = TF_DF(df_h, 1, HTF)
|
||
hdf = chan_h.dataframe
|
||
tl = htf_fx_timeline(signals_to_frame(extract_fx_signals(chan_h, hdf)), hdf)
|
||
full = attach_htf_context(sig, cdf, tl, "h1")
|
||
full["entry_ts"] = cdf["timestamp"].to_numpy()[full["entry_idx"].astype(int)]
|
||
return full
|
||
|
||
|
||
def overlap(a: pd.DataFrame, b: pd.DataFrame, period_ms: int,
|
||
tol_bars: int = 1) -> dict:
|
||
"""按时间戳比对两个信号集。方向也必须一致才算命中。"""
|
||
if a.empty or b.empty:
|
||
return {"a": len(a), "b": len(b), "同根": np.nan, f"±{tol_bars}根": np.nan}
|
||
bt = b["entry_ts"].to_numpy()
|
||
bd = b["direction"].to_numpy()
|
||
exact = near = 0
|
||
for ts, d in zip(a["entry_ts"].to_numpy(), a["direction"].to_numpy()):
|
||
hit = np.where((bt == ts) & (bd == d))[0]
|
||
if len(hit):
|
||
exact += 1
|
||
near += 1
|
||
continue
|
||
if np.any((np.abs(bt - ts) <= tol_bars * period_ms) & (bd == d)):
|
||
near += 1
|
||
return {"a": len(a), "b": len(b),
|
||
"同根": exact / len(a), f"±{tol_bars}根": near / len(a)}
|
||
|
||
|
||
def main() -> None:
|
||
ap = argparse.ArgumentParser()
|
||
ap.add_argument("--days", type=int, default=30)
|
||
ap.add_argument("--symbols", default="BTC,ETH,SOL")
|
||
ap.add_argument("--warmup", type=int, default=2000,
|
||
help="丢弃前若干根的信号,避开中枢左边界效应")
|
||
ap.add_argument("--refresh", action="store_true")
|
||
args = ap.parse_args()
|
||
|
||
from lib.data import load_local
|
||
|
||
syms = [s.strip() for s in args.symbols.split(",")]
|
||
period_ms = 60_000
|
||
print(f"[venue 对齐] {syms} · 近 {args.days} 天 1m · "
|
||
f"预热丢弃 {args.warmup} 根\n", flush=True)
|
||
|
||
bar_rows, sig_rows = [], []
|
||
for sym in syms:
|
||
print(f"── {sym}", flush=True)
|
||
bg_l = fetch_bitget(sym, LTF, args.days, args.refresh)
|
||
bg_h = fetch_bitget(sym, HTF, args.days, args.refresh)
|
||
|
||
pair = f"{sym}/USDT:USDT"
|
||
bn_l_all = load_local(pair, LTF)
|
||
bn_h_all = load_local(pair, HTF)
|
||
if bn_l_all is None or bn_h_all is None:
|
||
print(f" 跳过:本地无 Binance 数据")
|
||
continue
|
||
|
||
# 只比两家都有的那段时间
|
||
lo = max(bg_l["timestamp"].min(), bn_l_all["timestamp"].min())
|
||
hi = min(bg_l["timestamp"].max(), bn_l_all["timestamp"].max())
|
||
bg_l = bg_l[(bg_l.timestamp >= lo) & (bg_l.timestamp <= hi)].reset_index(drop=True)
|
||
bn_l = bn_l_all[(bn_l_all.timestamp >= lo) & (bn_l_all.timestamp <= hi)].reset_index(drop=True)
|
||
bg_h = bg_h[bg_h.timestamp <= hi].reset_index(drop=True)
|
||
bn_h = bn_h_all[(bn_h_all.timestamp >= bg_h["timestamp"].min())
|
||
& (bn_h_all.timestamp <= hi)].reset_index(drop=True)
|
||
|
||
span_d = (hi - lo) / 86_400_000
|
||
expect = int((hi - lo) / period_ms) + 1
|
||
# K 线层比对
|
||
m = bg_l.merge(bn_l, on="timestamp", suffixes=("_bg", "_bn"))
|
||
dc = (m["close_bg"] - m["close_bn"]) / m["close_bn"] * 1e4
|
||
dh = (m["high_bg"] - m["high_bn"]) / m["high_bn"] * 1e4
|
||
dl = (m["low_bg"] - m["low_bn"]) / m["low_bn"] * 1e4
|
||
bar_rows.append({
|
||
"品种": sym, "重叠天数": round(span_d, 1),
|
||
"Bitget根数": len(bg_l), "Binance根数": len(bn_l),
|
||
"应有根数": expect,
|
||
"Bitget缺口": expect - len(bg_l), "Binance缺口": expect - len(bn_l),
|
||
"共有根数": len(m),
|
||
"close中位差": f"{dc.median():+.2f}bp",
|
||
"close绝对差P95": f"{dc.abs().quantile(.95):.2f}bp",
|
||
"high绝对差P95": f"{dh.abs().quantile(.95):.2f}bp",
|
||
"low绝对差P95": f"{dl.abs().quantile(.95):.2f}bp",
|
||
})
|
||
print(f" K线:重叠 {span_d:.1f} 天,共有 {len(m)} 根,"
|
||
f"close 中位差 {dc.median():+.2f}bp,P95 {dc.abs().quantile(.95):.2f}bp",
|
||
flush=True)
|
||
|
||
# 信号层比对
|
||
t0 = time.perf_counter()
|
||
s_bg = pipeline(bg_l, bg_h)
|
||
s_bn = pipeline(bn_l, bn_h)
|
||
print(f" 管线跑完 {time.perf_counter() - t0:.1f}s", flush=True)
|
||
if s_bg.empty or s_bn.empty:
|
||
print(" 信号为空,跳过信号层")
|
||
continue
|
||
|
||
cut_bg = bg_l["timestamp"].to_numpy()[min(args.warmup, len(bg_l) - 1)]
|
||
cut_bn = bn_l["timestamp"].to_numpy()[min(args.warmup, len(bn_l) - 1)]
|
||
cut = max(cut_bg, cut_bn)
|
||
s_bg = s_bg[s_bg.entry_ts >= cut]
|
||
s_bn = s_bn[s_bn.entry_ts >= cut]
|
||
f_bg = s_bg[s_bg["h1_agree"] == 1]
|
||
f_bn = s_bn[s_bn["h1_agree"] == 1]
|
||
|
||
for tag, x, y in (("原始", s_bg, s_bn), ("5m同向后", f_bg, f_bn)):
|
||
o1 = overlap(x, y, period_ms) # Bitget 的信号有多少在 Binance 也有
|
||
o2 = overlap(y, x, period_ms) # 反向
|
||
sig_rows.append({
|
||
"品种": sym, "口径": tag,
|
||
"Bitget信号": o1["a"], "Binance信号": o1["b"],
|
||
"BG→BN同根": f"{o1['同根'] * 100:.1f}%",
|
||
"BG→BN±1根": f"{o1['±1根'] * 100:.1f}%",
|
||
"BN→BG同根": f"{o2['同根'] * 100:.1f}%",
|
||
"BN→BG±1根": f"{o2['±1根'] * 100:.1f}%",
|
||
})
|
||
print(f" {tag}:Bitget {o1['a']} 笔 / Binance {o1['b']} 笔,"
|
||
f"同根重合 {o1['同根'] * 100:.1f}%,±1根 {o1['±1根'] * 100:.1f}%",
|
||
flush=True)
|
||
|
||
if not bar_rows:
|
||
print("无结果")
|
||
return
|
||
|
||
print("\n" + "=" * 120)
|
||
print("########## 1. K 线层 ##########")
|
||
tb_bar = pd.DataFrame(bar_rows)
|
||
print(tb_bar.to_string(index=False))
|
||
print(" 缺口是「应有根数 − 实际根数」,永续在极端行情或维护时会漏推。")
|
||
|
||
print("\n########## 2. 信号层 ##########")
|
||
tb_sig = pd.DataFrame(sig_rows)
|
||
print(tb_sig.to_string(index=False))
|
||
|
||
print("\n########## 结论 ##########")
|
||
fin = tb_sig[tb_sig["口径"] == "5m同向后"]
|
||
if not fin.empty:
|
||
v = fin["BG→BN同根"].str.rstrip("%").astype(float)
|
||
print(f" 过滤后口径的同根重合率:{v.min():.1f}% ~ {v.max():.1f}%,"
|
||
f"均值 {v.mean():.1f}%")
|
||
print(" 高 → 滑点可直接对照 Binance 回测的 3.9bp 预算;")
|
||
print(" 低 → 必须先补 Bitget 自己的 1m 回测基线,否则实验不可归因。")
|
||
|
||
out_dir = RESEARCH / "out"
|
||
tb_bar.to_csv(out_dir / "venue_parity_bars.csv", index=False)
|
||
tb_sig.to_csv(out_dir / "venue_parity_signals.csv", index=False)
|
||
print(f"\n产物写入 {out_dir}/venue_parity_bars.csv 与 venue_parity_signals.csv")
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|