"""延迟基准 A:ccxt.pro 侧的数据到手时刻。 要测的事件只有一个:**我们在什么时候得知第 N 根 1m 已经收盘**。 它等价于 t_data − t_close,是下单延迟里最先发生、也往往最大的一段。 与 latency_hummingbot.py 测的是同一个事件,两边跑同一段时间才可比, 所以两个脚本都按整分钟对齐输出,事后按 K 线时间戳 join。 判据:ETH 的 Bitget 原生滑点余量只有 4.02bp、SOL 2.92bp,而本机 1m 波动 ETH 8.6bp/分钟、SOL 9.6bp/分钟。按 √t 折算,1 秒延迟就是 1.1~1.2bp 的 随机漂移,且入场方向上还有系统性追价。所以这个数值本身就可能决定生死。 输出 out/latency_ccxt.csv,每根一行。 """ from __future__ import annotations import argparse import asyncio import csv import signal import sys import time from pathlib import Path HERE = Path(__file__).resolve().parent RESEARCH = HERE.parent sys.path.insert(0, str(RESEARCH)) sys.path.insert(0, str(RESEARCH.parent)) SYMS = ("BTC", "ETH", "SOL") OUT = RESEARCH / "out" / "latency_ccxt.csv" PERIOD_MS = 60_000 _stop = False def _on_signal(*_): global _stop _stop = True async def watch(ex, sym: str, writer, fh, stats: dict) -> None: """watch_ohlcv 每次推送都带整段最近 K 线,靠时间戳前进判断上一根已收盘。""" pair = f"{sym}/USDT:USDT" last_ts = None while not _stop: try: o = await ex.watch_ohlcv(pair, "1m") except Exception as e: print(f" {sym} watch 异常 {type(e).__name__}: {e}", flush=True) await asyncio.sleep(1) continue if not o: continue now_ms = int(time.time() * 1000) newest = int(o[-1][0]) if last_ts is None: last_ts = newest continue if newest > last_ts: # newest 是刚开始的那根,故 last_ts 那根在 newest 时刻收盘 closed_ts = newest lag_ms = now_ms - closed_ts writer.writerow({"kline_ts": closed_ts, "sym": sym, "t_data_ms": now_ms, "lag_ms": lag_ms}) fh.flush() stats.setdefault(sym, []).append(lag_ms) n = len(stats[sym]) if n % 5 == 1: med = sorted(stats[sym])[n // 2] print(f" {sym}: 第 {n} 根,本次 lag {lag_ms}ms,中位 {med}ms", flush=True) last_ts = newest async def main_async(minutes: int) -> None: import ccxt.pro as ccxtpro ex = ccxtpro.bitget({"options": {"defaultType": "swap"}, "enableRateLimit": True}) OUT.parent.mkdir(parents=True, exist_ok=True) fh = OUT.open("w", newline="") writer = csv.DictWriter(fh, fieldnames=["kline_ts", "sym", "t_data_ms", "lag_ms"]) writer.writeheader() stats: dict = {} print(f"[ccxt.pro 延迟] {SYMS} · 计划跑 {minutes} 分钟 · 输出 {OUT.name}", flush=True) tasks = [asyncio.create_task(watch(ex, s, writer, fh, stats)) for s in SYMS] deadline = time.time() + minutes * 60 while time.time() < deadline and not _stop: await asyncio.sleep(1) for t in tasks: t.cancel() await asyncio.gather(*tasks, return_exceptions=True) await ex.close() fh.close() print("\n########## t_data − t_close(ms)##########") for s in SYMS: v = sorted(stats.get(s, [])) if not v: print(f" {s}: 无样本") continue print(f" {s}: n={len(v)} 中位 {v[len(v) // 2]}ms " f"P90 {v[int(len(v) * .9)]}ms 最大 {v[-1]}ms") print(f"\n产物写入 {OUT}") def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--minutes", type=int, default=20) args = ap.parse_args() signal.signal(signal.SIGINT, _on_signal) signal.signal(signal.SIGTERM, _on_signal) asyncio.run(main_async(args.minutes)) if __name__ == "__main__": main()