吃单查询换成 Hummingbot 的 OrderBook.get_vwap_for_volume:手写的 walk_book 返回的是按计价币吃单的加权均价,但框架的 get_price_for_quote_volume 返回 边际价、get_vwap_for_volume 收基础币量,两者语义不同。改为按基础币下单 (真实委托与 PositionExecutor.amount 均是基础币计价),深度不足由 query_volume/result_volume 判定,框架此时返回 nan 而非一个看似正常的 部分成交均价。 落盘完整盘口(双边 50 档)。此前只记三个固定名义额的成交价,这批数据的 寿命就等于那几个档位的寿命;存完整深度后任意资金量级的冲击都能离线重算。 仓位档同时从 1k/5k/20k 提到十万量级,此前低估真实仓位约两个数量级。 订阅成交流,按根按价位聚合。买卖分开存——多头在目标位挂卖出靠主动买盘 成交,混在一起会把成交率高估约一倍。BTC 每根总成交额中位与 210 天历史 的 volume×close 差 0.3%,可确认采集完整。 新增两项测算: - 冲击不是绑定约束。32 万仓位单边冲击 0.19~2.39bp,对 8.58~20.64bp 的 预算只占 1.6~14.2%,冲击反推的资金上限 100~500 万。 - maker 成交率才是。止盈位被首次触及时,限价在该根价格区间中的位置 中位 k=0.28(63.9 万次触及,三币一致);合并每根成交额后,32 万仓位 的全额成交率仅 30.1%/15.6%/1.5%。要 80% 全额成交,仓位须 ≤ 4.7 万 /1.4 万/0.26 万——比冲击反推的上限低 40~370 倍。 回测把这些止盈按「全额成交在目标价」计,故预算所依据的收益流本身需重估。 Co-authored-by: Cursor <cursoragent@cursor.com>
770 lines
35 KiB
Python
770 lines
35 KiB
Python
"""影子交易器:在 Hummingbot 运行时上测 1m 腿的真实入场滑点。
|
||
|
||
不下单。用 Hummingbot 的 Bitget 连接器取真实盘口,按信号方向和仓位吃单深度
|
||
算出「若此刻市价单进场会成交在哪」,再与回测假设的成交价相减。
|
||
|
||
为什么必须跑在 Hummingbot 上而不是自写脚本:要测的是**生产路径**的滑点。
|
||
决定成交价的是实际执行链路的延迟,换个运行时测出来的数就不作数了。
|
||
(连接器的换根解析 bug 见 patched_candles.py,已修,拿回约 1.06 秒。)
|
||
|
||
为什么不能用 paper trade 的成交:那是 Hummingbot 自己的撮合模型模拟的,
|
||
测出来是模型行为不是市场行为。
|
||
|
||
口径对齐 step42_exit_tp_1m.py:回测假设成交在**信号次根的开盘价**
|
||
(entry_delay=1),所以基准价就是换根后新一根的 open。滑点为正表示比回测差。
|
||
滤网(同向 + 中枢阶梯 + ATR 门控)在 shadow_signal.py 里,必须与预算同源。
|
||
|
||
### 统计口径三条硬要求
|
||
|
||
1. **主口径只用 pass_all 的信号根**。未过滤的照记但只作提前读数——
|
||
在我们根本不会下单的根上测滑点会把判据算宽
|
||
2. **条件漂移与无条件漂移分开报**。所以每根 K 线都记一份漂移
|
||
(shadow_drift.csv),不只信号根。两者的差就是「系统性追价」的大小
|
||
3. **出场腿按 maker/taker 分开**。止盈挂限价不吃滑点,把那 60% 混进
|
||
平均值会低估真实成本。出场腿属持仓管理,尚未实现
|
||
|
||
### lag 探针:超阈值要停开仓,不能只打日志
|
||
|
||
补丁只防得住「上游代码变了」,防不住 Bitget 再改一次消息格式。每根记
|
||
本地接收 − K 线收盘,近 30 根取中位数,超 800ms 即判该币不健康。
|
||
要停开仓是因为这种退化是**经济性且静默**的:不崩不报错,只让收益慢慢
|
||
变差,几周后才从统计里看得出来。影子期不下单,故落到 lag_ok 字段上。
|
||
|
||
### 盘口滚动缓冲把延迟变成自变量
|
||
|
||
每 100ms 存一份盘口。信号触发后,不只记「我们实际算完时」的滑点,而是回查
|
||
t_close+0.5s / 1s / 2s / 5s 各一个。这样即使本机算得慢,也能读出「若延迟为
|
||
X 秒,滑点是多少」,决策不被自身实现拖累。
|
||
|
||
### 滑点分解
|
||
|
||
延迟漂移 中间价相对次根开盘价的偏移——主项,且入场方向上系统性追价
|
||
盘口价差 最优价相对中间价
|
||
深度冲击 吃单加权价相对最优价
|
||
|
||
Bitget 永续实测价差仅约 0.01bp、100 档深度,故预期延迟漂移占绝大部分。
|
||
|
||
docker run -d --name shadow -w /home/hummingbot \\
|
||
-e PYTHONPATH=/home/hummingbot:/repo/research:/repo/research/live:/repo \\
|
||
-v $PWD:/repo:ro -v $PWD/research/out:/out \\
|
||
--entrypoint /opt/conda/envs/hummingbot/bin/python \\
|
||
hummingbot/hummingbot:latest /repo/research/live/shadow_hb.py --hours 24
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import csv
|
||
import gzip
|
||
import json
|
||
import math
|
||
import time
|
||
from collections import deque
|
||
from concurrent.futures import ProcessPoolExecutor
|
||
from concurrent.futures.process import BrokenProcessPool
|
||
from pathlib import Path
|
||
|
||
import numpy as np
|
||
import pandas as pd
|
||
|
||
from lib.shadow_budget import LAG_ALARM_MS, LAG_WINDOW, lag_healthy
|
||
|
||
SYMS = ("BTC", "ETH", "SOL")
|
||
# 多存一根:deque 尾部是尚未收盘的当前根,剔除后正好剩 step39 定下的窗口
|
||
LTF_BARS, HTF_BARS = 2001, 801 # 有效窗口 2000 / 800,命中率在此饱和
|
||
BOOK_HZ = 10 # 盘口采样 10Hz
|
||
BOOK_KEEP_S = 30 # 缓冲保留 30 秒,够回查到 +5s
|
||
BOOK_TOL_MS = 250 # 回查容差:10Hz 正常 ≤100ms,留些余量
|
||
BOOK_DEPTH = 50 # 双边各 50 档,实测能撑 78 万~261 万美元
|
||
DELAYS_S = (0.5, 1.0, 2.0, 5.0) # 回查点
|
||
# 便利视图用的仓位档。**真正的答案在 shadow_books.jsonl.gz 里**——完整盘口
|
||
# 落了盘,任意资金量级的冲击都能离线算,换个规模不必重测。这里的档位只是
|
||
# 为了让 CSV 直接可读,覆盖到按 2ATR 止损反推的十万量级真实仓位
|
||
NOTIONALS = (50_000.0, 100_000.0, 320_000.0, 530_000.0)
|
||
NUM_COLS = ["timestamp", "open", "high", "low", "close", "volume"]
|
||
|
||
|
||
def out_dir() -> Path:
|
||
p = Path("/out")
|
||
return p if p.is_dir() else Path(__file__).resolve().parents[1] / "out"
|
||
|
||
|
||
def _writer(path: Path, cols: list[str]):
|
||
"""追加模式打开;表头对不上就先把旧文件归档。
|
||
|
||
不校验的话,列一改,DictWriter 会按新顺序把行写到旧表头下面——
|
||
读出来整片错位,而且没有任何报错。长跑靠追加续命,这个校验是必需的。
|
||
"""
|
||
if path.exists() and path.stat().st_size > 0:
|
||
with path.open(newline="") as fh:
|
||
old = next(csv.reader(fh), [])
|
||
if old != cols:
|
||
arch = path.parent / "archive"
|
||
arch.mkdir(exist_ok=True)
|
||
dst = arch / f"{path.stem}_{time.strftime('%Y%m%d_%H%M%S')}.csv"
|
||
path.rename(dst)
|
||
print(f" [CSV] {path.name} 表头已变,旧数据归档为 {dst.name}",
|
||
flush=True)
|
||
fresh = not path.exists() or path.stat().st_size == 0
|
||
f = path.open("a", newline="")
|
||
w = csv.DictWriter(f, fieldnames=cols)
|
||
if fresh:
|
||
w.writeheader()
|
||
return f, w
|
||
|
||
|
||
class BookLog:
|
||
"""把完整盘口快照落成 gzip JSONL。
|
||
|
||
只记「某几个仓位档的成交价」的话,这批数据的寿命就等于那几个档位的寿命:
|
||
换一次资金规模就得重跑一周。存完整深度后,任意仓位的冲击都能离线重算,
|
||
一次采集回答所有资金量级的问题——包括容量上限那个必须现在就算、
|
||
不该等实盘暴露的数。
|
||
|
||
用 gzip 追加(多个 gzip 成员首尾相接仍可正常解压),进程被杀也只丢最后
|
||
一个缓冲块,不会毁掉整个文件。
|
||
"""
|
||
|
||
def __init__(self, path: Path) -> None:
|
||
self.path = path
|
||
self.fh = gzip.open(path, "at", encoding="utf-8")
|
||
self.n = 0
|
||
|
||
def write(self, sym: str, kline_ts: int, label: str, delay_ms: int,
|
||
target: int, book_ts: int, bids: np.ndarray,
|
||
asks: np.ndarray) -> None:
|
||
# 只留价与量两列,update_id 对离线分析没用。round 到 10 位避免
|
||
# float repr 把文件撑大一倍
|
||
rec = {"sym": sym, "kline_ts": kline_ts, "label": label,
|
||
"delay_ms": delay_ms, "target": target, "book_ts": book_ts,
|
||
"bids": [[round(float(p), 10), round(float(a), 10)]
|
||
for p, a, *_ in bids],
|
||
"asks": [[round(float(p), 10), round(float(a), 10)]
|
||
for p, a, *_ in asks]}
|
||
self.fh.write(json.dumps(rec, separators=(",", ":")) + "\n")
|
||
self.n += 1
|
||
|
||
def flush(self) -> None:
|
||
self.fh.flush()
|
||
|
||
def close(self) -> None:
|
||
try:
|
||
self.fh.close()
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
class TapeLog:
|
||
"""按 K 线、按价位聚合成交量,用来判 maker 腿能不能全额成交。
|
||
|
||
回测假设 3ATR 和 8ATR 的限价单全额成交。十万量级挂在那里,全成交还是
|
||
部分成交是完全不同的事——部分成交会把分批出场的收益结构改掉,而这个
|
||
问题盘口深度回答不了:深度说的是「现在有多少人挂着」,成交率问的是
|
||
「之后有多少人打过来」。只有成交流能回答。
|
||
|
||
**买卖必须分开存。** 多头在 3ATR 挂卖出止盈,成交靠的是主动**买盘**
|
||
打上来;把双边成交量合在一起会把成交率高估约一倍。
|
||
|
||
聚合到「根 × 价位」而不是逐笔:判据是「本根内有多少量在 ≥ 限价处成交」,
|
||
逐笔的时序对这个判据没有增量信息,而聚合能把体量压下两个数量级。
|
||
"""
|
||
|
||
def __init__(self, path: Path) -> None:
|
||
self.fh = gzip.open(path, "at", encoding="utf-8")
|
||
# sym -> side('b'/'s') -> price -> 累计基础币量
|
||
self.acc: dict[str, dict[str, dict[float, float]]] = {}
|
||
self.n_trades = 0
|
||
|
||
def add(self, sym: str, is_buy: bool, price: float, amount: float) -> None:
|
||
d = self.acc.setdefault(sym, {"b": {}, "s": {}})
|
||
side = d["b"] if is_buy else d["s"]
|
||
side[price] = side.get(price, 0.0) + amount
|
||
self.n_trades += 1
|
||
|
||
def flush_bar(self, sym: str, bar_ts: int) -> None:
|
||
"""一根走完就把这根的聚合结果落盘并清空。"""
|
||
d = self.acc.get(sym)
|
||
if not d or (not d["b"] and not d["s"]):
|
||
return
|
||
rec = {"sym": sym, "bar_ts": bar_ts,
|
||
"buys": {f"{p:.10g}": round(v, 10)
|
||
for p, v in sorted(d["b"].items())},
|
||
"sells": {f"{p:.10g}": round(v, 10)
|
||
for p, v in sorted(d["s"].items())}}
|
||
self.fh.write(json.dumps(rec, separators=(",", ":")) + "\n")
|
||
self.fh.flush()
|
||
self.acc[sym] = {"b": {}, "s": {}}
|
||
|
||
def close(self) -> None:
|
||
try:
|
||
self.fh.close()
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def hb_to_research(cdf: pd.DataFrame) -> pd.DataFrame:
|
||
"""Hummingbot 的 candles_df 转成 research/lib/data.py 的列结构。
|
||
|
||
HB 的 timestamp 是秒且无 date 列;chanlun 的 kline builder 需要真 datetime,
|
||
时区跟 lib/data.py 取 Asia/Shanghai,保证与回测同一口径。
|
||
"""
|
||
ts_ms = (cdf["timestamp"].astype("int64") * 1000)
|
||
date = pd.to_datetime(ts_ms, unit="ms", utc=True).dt.tz_convert("Asia/Shanghai")
|
||
out = pd.DataFrame({"timestamp": ts_ms.astype("int64"), "date": date})
|
||
for c in ("open", "high", "low", "close", "volume"):
|
||
out[c] = pd.to_numeric(cdf[c], errors="coerce")
|
||
return out.dropna().drop_duplicates(subset=["timestamp"]) \
|
||
.sort_values("timestamp").reset_index(drop=True)
|
||
|
||
|
||
def book_from(bids: np.ndarray, asks: np.ndarray):
|
||
"""用缓冲里的快照临时搭一个 OrderBook,以便调用框架自带的吃单查询。
|
||
|
||
自己手写吃单曾经踩过两个坑,框架版都没有:`get_vwap_for_volume` 返回的
|
||
是真加权均价(市价单的实际成交价),而 `get_price_for_quote_volume` 返回
|
||
的是**边际价**,用后者会高估冲击;深度不足时框架返回 nan 而不是一个
|
||
「看起来很正常」的部分成交均价,靠 query_volume/result_volume 判断。
|
||
"""
|
||
from hummingbot.core.data_type.order_book import OrderBook
|
||
ob = OrderBook()
|
||
ob.apply_numpy_snapshot(bids, asks)
|
||
return ob
|
||
|
||
|
||
class BookBuffer:
|
||
"""每币一份滚动盘口。按时间戳回查,取第一个不早于目标时刻的快照。
|
||
|
||
回查必须有容差上界。10Hz 下正常落在目标后 100ms 内(所有延迟点同向
|
||
偏约 +50ms,不影响曲线形状),但采样一旦卡顿,标着「0.5s」的那行可能
|
||
用的是 +3s 的盘口——数据看不出异常,判读却已经错了。超容差宁可丢弃,
|
||
并且把快照实际时刻写进 CSV,让这件事事后可查。
|
||
"""
|
||
|
||
def __init__(self) -> None:
|
||
self.buf: dict[str, deque] = {s: deque() for s in SYMS}
|
||
self.n_stale = 0 # 因超容差被丢弃的回查次数
|
||
|
||
def push(self, sym: str, t_ms: int, bids: list, asks: list) -> None:
|
||
d = self.buf[sym]
|
||
d.append((t_ms, bids, asks))
|
||
cutoff = t_ms - BOOK_KEEP_S * 1000
|
||
while d and d[0][0] < cutoff:
|
||
d.popleft()
|
||
|
||
def at(self, sym: str, t_ms: int) -> tuple | None:
|
||
for snap in self.buf[sym]:
|
||
if snap[0] >= t_ms:
|
||
if snap[0] - t_ms > BOOK_TOL_MS:
|
||
self.n_stale += 1
|
||
return None
|
||
return snap
|
||
return None
|
||
|
||
|
||
class Shadow:
|
||
def __init__(self, workers: int, hours: float) -> None:
|
||
self.workers = workers
|
||
self.deadline = time.time() + hours * 3600
|
||
self.books = BookBuffer()
|
||
self.pool: ProcessPoolExecutor | None = None
|
||
self.feeds_l: dict = {}
|
||
self.feeds_h: dict = {}
|
||
self.connector = None
|
||
self.stop = asyncio.Event()
|
||
# 持有 fire-and-forget 任务的强引用。只 create_task 不留引用的话,
|
||
# 任务可能在完成前被 GC 掉,asyncio 官方文档明确警告过这一点
|
||
self._tasks: set = set()
|
||
self.n_broken = 0
|
||
self._hb_last_bars = 0
|
||
# 成交监听:已挂上的币,以及必须持有的 forwarder 强引用
|
||
# (PubSub 只存弱引用,不持有的话监听会被 GC 静默摘掉)
|
||
self._hooked: set[str] = set()
|
||
self._trade_fwd: dict = {}
|
||
self.n_signal = 0
|
||
self.n_pass = 0
|
||
self.n_bars = 0
|
||
# lag 探针的滚动窗口,逐币独立:一个币的行情退化不该连累其他币
|
||
self.lag_hist: dict[str, deque] = {
|
||
s: deque(maxlen=LAG_WINDOW) for s in SYMS}
|
||
self.lag_ok: dict[str, bool] = {s: True for s in SYMS}
|
||
d = out_dir()
|
||
# 追加模式:长跑期间若重启,已收集的样本不该被清掉
|
||
self.f_sig, self.w_sig = _writer(d / "shadow_signals.csv", [
|
||
"sym", "kline_ts", "direction",
|
||
"h1_agree", "ladder_ok", "gate_ok", "pass_all", "lag_ok",
|
||
"atr_pct", "atr_bp",
|
||
"t_close_ms", "t_data_ms", "t_signal_ms",
|
||
"lag_data_ms", "lag_signal_ms",
|
||
"delay_label", "delay_ms", "book_ts", "book_lag_ms",
|
||
"notional", "base_amt",
|
||
"baseline_px", "mid", "best_px", "fill_px", "filled", "depth_ok",
|
||
"slip_bp", "drift_bp", "spread_bp", "impact_bp"])
|
||
self.f_lat, self.w_lat = _writer(d / "shadow_latency.csv", [
|
||
"sym", "kline_ts", "t_close_ms", "t_data_ms", "t_signal_ms",
|
||
"lag_data_ms", "lag_signal_ms", "compute_ms", "n_bars", "n_hits",
|
||
"n_pass", "atr_bp", "lag_med_ms", "lag_ok"])
|
||
# 无条件漂移:每根都记,用来和信号根上的条件漂移对照
|
||
self.f_drf, self.w_drf = _writer(d / "shadow_drift.csv", [
|
||
"sym", "kline_ts", "delay_label", "delay_ms",
|
||
"book_ts", "book_lag_ms", "baseline_px", "mid", "drift_bp_long"])
|
||
# 完整深度。挂在无条件漂移那条路径上,所以每根 K 线的四个固定延迟点
|
||
# 都有一份,信号根上再补一份 actual 点
|
||
self.blog = BookLog(d / "shadow_books.jsonl.gz")
|
||
self.tape = TapeLog(d / "shadow_tape.jsonl.gz")
|
||
|
||
def _spawn(self, coro, what: str) -> None:
|
||
"""起一个后台任务,但异常要吼出来。
|
||
|
||
裸 create_task 的异常只在对象被 GC 时才由 asyncio 打一句
|
||
「Task exception was never retrieved」,很容易整晚没人发现。
|
||
这套东西最怕的就是不崩不报错的静默退化。
|
||
"""
|
||
async def guard():
|
||
try:
|
||
await coro
|
||
except asyncio.CancelledError:
|
||
raise
|
||
except Exception as e:
|
||
import traceback
|
||
print(f" [异常] {what}: {type(e).__name__}: {e}", flush=True)
|
||
traceback.print_exc()
|
||
|
||
t = asyncio.create_task(guard())
|
||
self._tasks.add(t)
|
||
t.add_done_callback(self._tasks.discard)
|
||
|
||
def _restart_pool(self) -> None:
|
||
"""进程池坏了之后重建。
|
||
|
||
用 spawn 而非 fork:此刻进程里已经有活跃的 WS 连接,fork 会把连接
|
||
状态一起复制进子进程。spawn 启动慢几秒,但只在故障时走这条路。
|
||
"""
|
||
import multiprocessing
|
||
self.n_broken += 1
|
||
try:
|
||
self.pool.shutdown(wait=False, cancel_futures=True)
|
||
except Exception:
|
||
pass
|
||
self.pool = ProcessPoolExecutor(
|
||
max_workers=self.workers,
|
||
mp_context=multiprocessing.get_context("spawn"))
|
||
print(f" [进程池] 已重建(第 {self.n_broken} 次)", flush=True)
|
||
|
||
# ---------- 启动 ----------
|
||
|
||
async def start(self) -> None:
|
||
from hummingbot.connector.derivative.bitget_perpetual.bitget_perpetual_derivative import (
|
||
BitgetPerpetualDerivative,
|
||
)
|
||
from patched_candles import (PatchedBitgetPerpetualCandles,
|
||
assert_patch_effective)
|
||
|
||
# 覆盖失效是静默的(悄悄退回慢 1.06 秒,不报错),所以在启动就验一次
|
||
assert_patch_effective()
|
||
|
||
for s in SYMS:
|
||
self.feeds_l[s] = PatchedBitgetPerpetualCandles(
|
||
f"{s}-USDT", "1m", LTF_BARS)
|
||
self.feeds_h[s] = PatchedBitgetPerpetualCandles(
|
||
f"{s}-USDT", "5m", HTF_BARS)
|
||
self.feeds_l[s].start()
|
||
self.feeds_h[s].start()
|
||
print(f"[影子] {SYMS} · 1m×{LTF_BARS} + 5m×{HTF_BARS} · "
|
||
f"{self.workers} 个计算进程", flush=True)
|
||
|
||
# 只取公开数据:无密钥 + trading_required=False
|
||
self.connector = BitgetPerpetualDerivative(
|
||
bitget_perpetual_api_key="", bitget_perpetual_secret_key="",
|
||
bitget_perpetual_passphrase="",
|
||
trading_pairs=[f"{s}-USDT" for s in SYMS],
|
||
trading_required=False)
|
||
await self.connector.start_network()
|
||
print(" 连接器已启动,等盘口与历史回填", flush=True)
|
||
self._hook_trades()
|
||
|
||
t0 = time.time()
|
||
while time.time() - t0 < 600:
|
||
ready = all(f.ready for f in
|
||
list(self.feeds_l.values()) + list(self.feeds_h.values()))
|
||
books = all(self._snapshot(s) is not None for s in SYMS)
|
||
if ready and books:
|
||
break
|
||
await asyncio.sleep(1)
|
||
print(f" 就绪 {time.time() - t0:.1f}s · "
|
||
f"1m {[len(self.feeds_l[s]._candles) for s in SYMS]} 根 · "
|
||
f"5m {[len(self.feeds_h[s]._candles) for s in SYMS]} 根", flush=True)
|
||
|
||
def _hook_trades(self) -> None:
|
||
"""给每个盘口挂成交监听。
|
||
|
||
盘口对象可能还没建好(订阅是异步的),所以挂不上的先记下来,由
|
||
watch_bars 那圈重试;一直挂不上会在心跳里显示成交笔数为 0。
|
||
"""
|
||
from hummingbot.core.event.event_forwarder import EventForwarder
|
||
from hummingbot.core.event.events import OrderBookEvent
|
||
from hummingbot.core.data_type.common import TradeType
|
||
|
||
def make(sym: str):
|
||
def cb(ev) -> None:
|
||
self.tape.add(sym, ev.type == TradeType.BUY,
|
||
float(ev.price), float(ev.amount))
|
||
return EventForwarder(cb)
|
||
|
||
for s in SYMS:
|
||
if s in self._hooked: # 重复挂会让同一笔成交被记两次
|
||
continue
|
||
try:
|
||
ob = self.connector.get_order_book(f"{s}-USDT")
|
||
except Exception:
|
||
ob = None
|
||
if ob is None:
|
||
continue
|
||
fwd = make(s)
|
||
ob.add_listener(OrderBookEvent.TradeEvent, fwd)
|
||
self._trade_fwd[s] = fwd
|
||
self._hooked.add(s)
|
||
miss = [s for s in SYMS if s not in self._hooked]
|
||
print(f" 成交流已挂 {sorted(self._hooked)}"
|
||
+ (f",待重试 {miss}" if miss else ""), flush=True)
|
||
|
||
def _snapshot(self, sym: str):
|
||
try:
|
||
ob = self.connector.get_order_book(f"{sym}-USDT")
|
||
except Exception:
|
||
return None
|
||
if ob is None:
|
||
return None
|
||
# 存成 apply_numpy_snapshot 要的 [价, 量, update_id] 三列,
|
||
# 回查时才能直接搭 OrderBook 调框架的吃单查询
|
||
bids = np.array([(float(r.price), float(r.amount), i)
|
||
for i, (r, _) in enumerate(
|
||
zip(ob.bid_entries(), range(BOOK_DEPTH)))])
|
||
asks = np.array([(float(r.price), float(r.amount), i)
|
||
for i, (r, _) in enumerate(
|
||
zip(ob.ask_entries(), range(BOOK_DEPTH)))])
|
||
if not len(bids) or not len(asks):
|
||
return None
|
||
return bids, asks
|
||
|
||
# ---------- 三个循环 ----------
|
||
|
||
async def sample_books(self) -> None:
|
||
"""按截止时刻补睡,且对齐到墙钟 100ms 网格。
|
||
|
||
补睡是因为「干完活再睡固定时长」的实际周期是 100ms 加采样耗时,
|
||
名义 10Hz 到不了 10Hz。
|
||
|
||
对齐是因为回查目标都是 `kline_ts + n×500ms`,而 kline_ts 是整分钟,
|
||
所以目标必然落在墙钟 100ms 的整数倍上。采样相位若随启动时刻漂移,
|
||
每个回查点就会固定晚半个采样周期(实测 52ms)——四个固定延迟点
|
||
同向偏置,虽不改曲线形状,但白白多算了 50ms 的漂移。
|
||
"""
|
||
period = 1.0 / BOOK_HZ
|
||
nxt = math.ceil(time.time() / period) * period
|
||
while not self.stop.is_set():
|
||
t = int(time.time() * 1000)
|
||
for s in SYMS:
|
||
snap = self._snapshot(s)
|
||
if snap:
|
||
self.books.push(s, t, snap[0], snap[1])
|
||
nxt += period
|
||
await asyncio.sleep(max(0.0, nxt - time.time()))
|
||
|
||
async def watch_bars(self) -> None:
|
||
last = {s: (int(self.feeds_l[s]._candles[-1][0])
|
||
if len(self.feeds_l[s]._candles) else None) for s in SYMS}
|
||
while not self.stop.is_set():
|
||
for s in SYMS:
|
||
c = self.feeds_l[s]._candles
|
||
if not len(c):
|
||
continue
|
||
newest = int(c[-1][0])
|
||
if last[s] is not None and newest > last[s]:
|
||
t_data = int(time.time() * 1000)
|
||
kts = newest * 1000 if newest < 1e12 else newest
|
||
# 刚收盘那根的成交聚合先落盘,再算信号
|
||
self.tape.flush_bar(s, kts)
|
||
if len(self._hooked) < len(SYMS):
|
||
self._hook_trades() # 换根时才重试,避免重复挂
|
||
self._spawn(self.on_bar(s, kts, t_data), f"on_bar {s}")
|
||
last[s] = newest
|
||
await asyncio.sleep(0.01)
|
||
|
||
async def on_bar(self, sym: str, kline_ts: int, t_data: int) -> None:
|
||
"""kline_ts 是新一根的开盘时刻,也就是上一根的收盘时刻 t_close。"""
|
||
df_l = hb_to_research(self.feeds_l[sym].candles_df)
|
||
df_h = hb_to_research(self.feeds_h[sym].candles_df)
|
||
# 末行是刚开始的那根,未收盘,必须剔除,否则等于用未来数据
|
||
df_l = df_l[df_l["timestamp"] < kline_ts]
|
||
df_h = df_h[df_h["timestamp"] < kline_ts]
|
||
baseline = self._new_bar_open(sym, kline_ts)
|
||
|
||
lag_med, lag_ok = self._probe_lag(sym, t_data - kline_ts)
|
||
|
||
t0 = time.perf_counter()
|
||
payload = (df_l[NUM_COLS].values.tolist(),
|
||
df_h[NUM_COLS].values.tolist(), baseline)
|
||
loop = asyncio.get_running_loop()
|
||
from shadow_signal import compute_packed
|
||
try:
|
||
res = await loop.run_in_executor(self.pool, compute_packed, payload)
|
||
except BrokenProcessPool as e:
|
||
# 不重建的话,之后每一根都会走到这里,采集静默停摆到跑完为止
|
||
print(f" [{sym}] 进程池损坏 {e},重建后跳过本根", flush=True)
|
||
self._restart_pool()
|
||
return
|
||
compute_ms = int((time.perf_counter() - t0) * 1000)
|
||
t_signal = int(time.time() * 1000)
|
||
|
||
hits = res.get("hits", [])
|
||
atr_pct = res.get("atr_pct")
|
||
atr_bp = round(atr_pct * 1e4, 3) if atr_pct else ""
|
||
n_pass = sum(h["pass_all"] for h in hits)
|
||
|
||
self.n_bars += 1
|
||
self.w_lat.writerow({
|
||
"sym": sym, "kline_ts": kline_ts, "t_close_ms": kline_ts,
|
||
"t_data_ms": t_data, "t_signal_ms": t_signal,
|
||
"lag_data_ms": t_data - kline_ts,
|
||
"lag_signal_ms": t_signal - kline_ts,
|
||
"compute_ms": compute_ms, "n_bars": res.get("n_bars", 0),
|
||
"n_hits": len(hits), "n_pass": n_pass, "atr_bp": atr_bp,
|
||
"lag_med_ms": lag_med, "lag_ok": int(lag_ok)})
|
||
self.f_lat.flush()
|
||
|
||
if baseline is not None and np.isfinite(baseline):
|
||
# 无条件漂移:每根都记,不管有没有信号
|
||
self._spawn(self._drift_later(sym, kline_ts, baseline),
|
||
f"drift {sym}")
|
||
|
||
if res.get("error"):
|
||
print(f" [{sym}] 信号计算出错 {res['error']}", flush=True)
|
||
return
|
||
if not hits:
|
||
return
|
||
if baseline is None or not np.isfinite(baseline):
|
||
print(f" [{sym}] 有信号但拿不到次根开盘价,跳过", flush=True)
|
||
return
|
||
|
||
for h in hits:
|
||
self.n_signal += 1
|
||
self.n_pass += h["pass_all"]
|
||
mark = "★" if h["pass_all"] else "·"
|
||
print(f" {mark} [{sym}] {kline_ts} 方向 {h['direction']:+d} "
|
||
f"同向{h['h1_agree']} 阶梯{h['ladder_ok']} 门控{h['gate_ok']} "
|
||
f"(ATR {atr_bp or 'na'}bp) · 数据 {t_data - kline_ts}ms "
|
||
f"信号 {t_signal - kline_ts}ms", flush=True)
|
||
# 最远的回查点在 t_close+5s,此刻尚未发生;等它过去再一次性落盘
|
||
self._spawn(
|
||
self._record_later(sym, kline_ts, h, t_data, t_signal,
|
||
baseline, atr_pct, lag_ok),
|
||
f"record {sym}")
|
||
|
||
def _probe_lag(self, sym: str, lag_ms: int) -> tuple[float, bool]:
|
||
"""记一根的到达延迟,返回 (滚动中位数, 该币是否健康)。
|
||
|
||
不健康时应停止开新仓;影子期不下单,故只落到 lag_ok 字段并告警。
|
||
"""
|
||
self.lag_hist[sym].append(lag_ms)
|
||
ok = lag_healthy(self.lag_hist[sym])
|
||
med = float(np.median(self.lag_hist[sym]))
|
||
if ok != self.lag_ok[sym]:
|
||
state = "恢复" if ok else f"退化,超 {LAG_ALARM_MS:.0f}ms 阈值,停开新仓"
|
||
print(f" [lag] {sym} {state}:近 {len(self.lag_hist[sym])} 根"
|
||
f"中位 {med:.0f}ms", flush=True)
|
||
self.lag_ok[sym] = ok
|
||
return round(med, 1), ok
|
||
|
||
async def _wait_for_delays(self, kline_ts: int) -> None:
|
||
"""最远回查点是 t_close+5s,等它过去(多留 0.5s 给采样)。"""
|
||
wait = (kline_ts + int(max(DELAYS_S) * 1000) + 500) / 1000.0 - time.time()
|
||
if wait > 0:
|
||
await asyncio.sleep(wait)
|
||
|
||
async def _record_later(self, sym: str, kline_ts: int, hit: dict,
|
||
t_data: int, t_signal: int, baseline: float,
|
||
atr_pct: float | None, lag_ok: bool) -> None:
|
||
await self._wait_for_delays(kline_ts)
|
||
self._record(sym, kline_ts, hit, t_data, t_signal, baseline,
|
||
atr_pct, lag_ok)
|
||
|
||
async def _drift_later(self, sym: str, kline_ts: int,
|
||
baseline: float) -> None:
|
||
await self._wait_for_delays(kline_ts)
|
||
for label, delay_ms in self._points(None):
|
||
target = kline_ts + delay_ms
|
||
snap = self.books.at(sym, target)
|
||
if snap is None:
|
||
continue
|
||
book_ts, bids, asks = snap
|
||
mid = (float(bids[0][0]) + float(asks[0][0])) / 2.0
|
||
self.blog.write(sym, kline_ts, label, delay_ms, target,
|
||
book_ts, bids, asks)
|
||
self.w_drf.writerow({
|
||
"sym": sym, "kline_ts": kline_ts, "delay_label": label,
|
||
"delay_ms": delay_ms, "book_ts": book_ts,
|
||
"book_lag_ms": book_ts - target, "baseline_px": baseline,
|
||
"mid": mid,
|
||
"drift_bp_long": round((mid - baseline) / baseline * 1e4, 4)})
|
||
self.f_drf.flush()
|
||
self.blog.flush() # 每根冲刷一次,进程被杀最多丢一根
|
||
|
||
@staticmethod
|
||
def _points(t_signal_delay: int | None) -> list[tuple[str, int]]:
|
||
pts = [(f"{d}s", int(d * 1000)) for d in DELAYS_S]
|
||
if t_signal_delay is not None:
|
||
pts.insert(0, ("actual", t_signal_delay))
|
||
return pts
|
||
|
||
def _new_bar_open(self, sym: str, kline_ts: int) -> float | None:
|
||
"""次根开盘价 = 回测假设的成交价。"""
|
||
c = self.feeds_l[sym]._candles
|
||
if not len(c):
|
||
return None
|
||
row = c[-1]
|
||
ts = int(row[0])
|
||
ts = ts * 1000 if ts < 1e12 else ts
|
||
return float(row[1]) if ts == kline_ts else None
|
||
|
||
def _record(self, sym: str, kline_ts: int, hit: dict,
|
||
t_data: int, t_signal: int, baseline: float,
|
||
atr_pct: float | None, lag_ok: bool) -> None:
|
||
d_sign = hit["direction"]
|
||
|
||
for label, delay_ms in self._points(t_signal - kline_ts):
|
||
target = kline_ts + delay_ms
|
||
snap = self.books.at(sym, target)
|
||
if snap is None:
|
||
continue
|
||
book_ts, bids, asks = snap
|
||
best_bid, best_ask = float(bids[0][0]), float(asks[0][0])
|
||
mid = (best_bid + best_ask) / 2.0
|
||
best_px = best_ask if d_sign > 0 else best_bid
|
||
ob = book_from(bids, asks)
|
||
if label == "actual":
|
||
# 四个固定点已由无条件漂移那条路径落过,只补这一个
|
||
self.blog.write(sym, kline_ts, label, delay_ms, target,
|
||
book_ts, bids, asks)
|
||
|
||
for notional in NOTIONALS:
|
||
# 名义额按基准价折成基础币再下单——真实委托是基础币计价的,
|
||
# 框架的 get_vwap_for_volume 也收基础币量。名义额那一栏留着
|
||
# 是为了跨币可比(1 BTC 和 1 SOL 没法横向比)
|
||
base_amt = notional / baseline
|
||
r = ob.get_vwap_for_volume(d_sign > 0, base_amt)
|
||
fill = float(r.result_price)
|
||
depth_ok = int(float(r.result_volume) >= base_amt * 0.999)
|
||
if not np.isfinite(fill):
|
||
# 25 档吃不下这个量,框架直接给 nan。记一行标明深度不足,
|
||
# 免得「某个仓位档在薄盘时段整段消失」看不出来
|
||
self.w_sig.writerow({
|
||
"sym": sym, "kline_ts": kline_ts, "direction": d_sign,
|
||
"h1_agree": hit["h1_agree"],
|
||
"ladder_ok": hit["ladder_ok"],
|
||
"gate_ok": hit["gate_ok"], "pass_all": hit["pass_all"],
|
||
"lag_ok": int(lag_ok),
|
||
"atr_pct": atr_pct if atr_pct else "",
|
||
"atr_bp": round(atr_pct * 1e4, 3) if atr_pct else "",
|
||
"t_close_ms": kline_ts, "t_data_ms": t_data,
|
||
"t_signal_ms": t_signal,
|
||
"lag_data_ms": t_data - kline_ts,
|
||
"lag_signal_ms": t_signal - kline_ts,
|
||
"delay_label": label, "delay_ms": delay_ms,
|
||
"book_ts": book_ts, "book_lag_ms": book_ts - target,
|
||
"notional": notional, "base_amt": round(base_amt, 8),
|
||
"baseline_px": baseline, "mid": mid,
|
||
"best_px": best_px, "fill_px": "",
|
||
"filled": round(float(r.result_volume), 8),
|
||
"depth_ok": 0, "slip_bp": "", "drift_bp": "",
|
||
"spread_bp": "", "impact_bp": ""})
|
||
continue
|
||
slip = d_sign * (fill - baseline) / baseline * 1e4
|
||
drift = d_sign * (mid - baseline) / baseline * 1e4
|
||
spread = d_sign * (best_px - mid) / mid * 1e4
|
||
impact = d_sign * (fill - best_px) / best_px * 1e4
|
||
self.w_sig.writerow({
|
||
"sym": sym, "kline_ts": kline_ts,
|
||
"direction": d_sign, "h1_agree": hit["h1_agree"],
|
||
"ladder_ok": hit["ladder_ok"], "gate_ok": hit["gate_ok"],
|
||
"pass_all": hit["pass_all"], "lag_ok": int(lag_ok),
|
||
"atr_pct": atr_pct if atr_pct else "",
|
||
"atr_bp": round(atr_pct * 1e4, 3) if atr_pct else "",
|
||
"t_close_ms": kline_ts, "t_data_ms": t_data,
|
||
"t_signal_ms": t_signal,
|
||
"lag_data_ms": t_data - kline_ts,
|
||
"lag_signal_ms": t_signal - kline_ts,
|
||
"delay_label": label, "delay_ms": delay_ms,
|
||
"book_ts": book_ts, "book_lag_ms": book_ts - target,
|
||
"notional": notional, "base_amt": round(base_amt, 8),
|
||
"baseline_px": baseline,
|
||
"mid": mid, "best_px": best_px, "fill_px": fill,
|
||
"filled": round(float(r.result_volume), 8),
|
||
"depth_ok": depth_ok,
|
||
"slip_bp": round(slip, 4), "drift_bp": round(drift, 4),
|
||
"spread_bp": round(spread, 4),
|
||
"impact_bp": round(impact, 4)})
|
||
self.f_sig.flush()
|
||
|
||
async def heartbeat(self) -> None:
|
||
while not self.stop.is_set():
|
||
await asyncio.sleep(300)
|
||
depth = {s: len(self.books.buf[s]) for s in SYMS}
|
||
lag = {s: (f"{np.median(h):.0f}ms" if h else "na")
|
||
+ ("" if self.lag_ok[s] else "!")
|
||
for s, h in self.lag_hist.items()}
|
||
print(f" [心跳] 已处理 {self.n_bars} 根 · 命中 {self.n_signal} 个"
|
||
f"(过全部滤网 {self.n_pass}) · lag {lag} · 盘口缓冲 {depth}"
|
||
f" · 回查超容差 {self.books.n_stale} 次"
|
||
f" · 在途任务 {len(self._tasks)}"
|
||
f" · 盘口落盘 {self.blog.n} 份"
|
||
f" · 成交 {self.tape.n_trades} 笔{'' if self.tape.n_trades else ' ⚠监听未生效'}",
|
||
flush=True)
|
||
# 五分钟一根都没进来,说明管道断了。不喊一声就只能靠人翻日志
|
||
if self.n_bars == self._hb_last_bars:
|
||
print(f" ⚠ [停滞] 距上次心跳未处理任何 K 线"
|
||
f"(进程池重建 {self.n_broken} 次),管道可能已断",
|
||
flush=True)
|
||
self._hb_last_bars = self.n_bars
|
||
|
||
async def run(self) -> None:
|
||
await self.start()
|
||
tasks = [asyncio.create_task(self.sample_books()),
|
||
asyncio.create_task(self.watch_bars()),
|
||
asyncio.create_task(self.heartbeat())]
|
||
while time.time() < self.deadline:
|
||
await asyncio.sleep(5)
|
||
self.stop.set()
|
||
for t in tasks:
|
||
t.cancel()
|
||
await asyncio.gather(*tasks, return_exceptions=True)
|
||
for f in list(self.feeds_l.values()) + list(self.feeds_h.values()):
|
||
f.stop()
|
||
await self.connector.stop_network()
|
||
for f in (self.f_sig, self.f_lat, self.f_drf):
|
||
f.close()
|
||
self.blog.close()
|
||
self.tape.close()
|
||
print(f"\n收工:{self.n_bars} 根 · {self.n_signal} 个信号"
|
||
f"(过全部滤网 {self.n_pass})", flush=True)
|
||
|
||
|
||
async def main_async(workers: int, hours: float, pool) -> None:
|
||
sh = Shadow(workers, hours)
|
||
sh.pool = pool
|
||
await sh.run()
|
||
|
||
|
||
def main() -> None:
|
||
ap = argparse.ArgumentParser()
|
||
ap.add_argument("--hours", type=float, default=24.0)
|
||
ap.add_argument("--workers", type=int, default=2)
|
||
a = ap.parse_args()
|
||
# 进程池必须在事件循环和任何 WS 连接之前建好:fork 一个已带活跃 socket
|
||
# 的进程会把连接状态一起复制过去,后果不可预测
|
||
with ProcessPoolExecutor(max_workers=a.workers) as pool:
|
||
asyncio.run(main_async(a.workers, a.hours, pool))
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|