Files
Chan/research/live/venue_parity.py
T
UbuntuandCursor 7e339d2a54 research: 影子交易器落在 Hummingbot 上,并修掉 Bitget 连接器的换根延迟
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>
2026-08-27 23:51:57 +08:00

307 lines
12 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""前置测量二: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}bpP95 {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()