"""前置测量二: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()