"""影子交易器:在 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 os import socket 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 # 站点标识。跨地对比时两台机器的 CSV 要能合起来读,没有这一列就分不清哪行 # 来自哪台。默认取主机名,部署脚本会显式传 SHADOW_SITE(如 sg-hetzner) SITE = os.environ.get("SHADOW_SITE") or socket.gethostname() # 判「加 worker 有没有用」必须知道核数:CPU 密集的活,worker 超过核数不增吞吐 CORES = os.cpu_count() or 1 # 少于这么多根就不送去算。缠论要先有分型再有笔再有中枢,几十根出不来中枢, # 送过去只会白占一个计算槽 MIN_BARS = 200 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) # 回查点 # 主口径 10 万名义额(2026-08-28 定:不会有更大资金)。上下各留两档是为了 # 读出局部斜率——单点看不出「再大一倍会怎样」。 # **真正的答案在 shadow_books.jsonl.gz 里**:完整盘口已落盘,任意资金量级的 # 冲击都能离线重算,换规模不必重测,这里的档位只为让 CSV 直接可读 NOTIONALS = (25_000.0, 50_000.0, 100_000.0, 200_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" RUN_ID = time.strftime("%Y%m%dT%H%M%S", time.gmtime()) def run_path(path: Path) -> Path: """把 `x.jsonl.gz` 变成本轮专属的 `x..jsonl.gz`。 gzip 追加流在进程被杀后不可靠(见 BookLog 的说明),分文件是唯一能保证 历史数据不被后续运行连坐的办法。读侧 shadow_depth.read_jsonl_gz 会把 同前缀的所有文件一并读入,所以分文件对分析是透明的。 """ return path.with_name(path.name.replace(".jsonl.gz", f".{RUN_ID}.jsonl.gz")) 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 = _SiteWriter(csv.DictWriter(f, fieldnames=cols)) if fresh: w.writeheader() return f, w class _SiteWriter: """DictWriter 的薄包装,自动补上 site 列。 逐个 writerow 手加 site 有四处,漏一处就是静默的空值,而跨地对比正是靠 这一列区分数据来源。在这里注入,漏不掉。 """ def __init__(self, w: csv.DictWriter) -> None: self._w = w def writeheader(self) -> None: self._w.writeheader() def writerow(self, row: dict) -> None: row.setdefault("site", SITE) self._w.writerow(row) class BookLog: """把完整盘口快照落成 gzip JSONL。 只记「某几个仓位档的成交价」的话,这批数据的寿命就等于那几个档位的寿命: 换一次资金规模就得重跑一周。存完整深度后,任意仓位的冲击都能离线重算, 一次采集回答所有资金量级的问题——包括容量上限那个必须现在就算、 不该等实盘暴露的数。 **每轮运行单独一个文件**,不追加到同一个。追加看着更省事,实际很危险: 进程被 SIGKILL 时当前 gzip 成员停在 deflate 块中间,下一轮追加的新成员 接在这段垃圾字节之后,顺序解压会在损坏点抛 `invalid block type`,该点 之后的所有数据——包括后续每一轮写进去的——全都读不出来。已经因此丢过 一次。分文件后损坏最多只影响被杀那一轮的尾部。 """ def __init__(self, path: Path) -> None: self.path = run_path(path) self.fh = gzip.open(self.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 = {"site": SITE, "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 挂卖出止盈,成交靠的是主动**买盘** 打上来;把双边成交量合在一起会把成交率高估约一倍。 聚合到「根 × 价位」而不是逐笔:判据是「本根内有多少量在 ≥ 限价处成交」, 逐笔的时序对这个判据没有增量信息,而聚合能把体量压下两个数量级。 每轮运行单独一个文件,理由同 BookLog。 """ def __init__(self, path: Path) -> None: self.path = run_path(path) self.fh = gzip.open(self.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 = {"site": SITE, "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 # 排队 / 纯计算的滚动窗口,用来判断加核有没有用 self.q_hist: deque = deque(maxlen=90) self.i_hist: deque = deque(maxlen=90) # 每个收盘时刻「清空所有币」耗时。这才是决定信号何时可下单的量: # 币同一秒收盘,币数超 worker 数时后面的币串行等待,而这笔代价不 # 出现在任何单根的 queue_ms 或 inner_ms 里 self.clear_hist: deque = deque(maxlen=60) self._clear_cur: dict[int, float] = {} # 成交监听:已挂上的币,以及必须持有的 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", [ "site", "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", [ "site", "sym", "kline_ts", "t_close_ms", "t_data_ms", "t_signal_ms", "lag_data_ms", "lag_signal_ms", # compute_ms 含排队;queue_ms/inner_ms 把它拆开,用来判断加核有没有用 "compute_ms", "queue_ms", "inner_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", [ "site", "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] # WS 重连的瞬间 feed 的 deque 可能是空的。放行的话 worker 会抛 # 「DataFrame for 1m is empty」,白占一个计算槽(币数超核数时这笔 # 代价会推迟后面所有币),而报错文本还会让人以为是缺历史数据 if len(df_l) < MIN_BARS or len(df_h) < MIN_BARS: print(f" [{sym}] 窗口过短(1m {len(df_l)} / 5m {len(df_h)} 根)," f"跳过本根。feed 大概在重连", flush=True) return 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, time.time()) 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) if res.get("queue_ms") is not None: self.q_hist.append(res["queue_ms"]) if res.get("inner_ms") is not None: self.i_hist.append(res["inner_ms"]) # 同一 kline_ts 上取各币最大值即该时刻的清空耗时;只保留最近几个 # 时刻,否则这个 dict 会随运行时长无界增长 cur = self._clear_cur cur[kline_ts] = max(cur.get(kline_ts, 0.0), float(compute_ms)) if len(cur) > 3: done = min(cur) self.clear_hist.append(cur.pop(done)) 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, "queue_ms": res.get("queue_ms"), "inner_ms": res.get("inner_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.q_hist and self.i_hist: q, i = float(np.median(self.q_hist)), float(np.median(self.i_hist)) # 建议要看绝对量级:lean + 新引擎后纯计算约 128ms,此时再提 clear = float(np.median(self.clear_hist)) if self.clear_hist \ else float("nan") print(f" [计算] 每币排队 {q:.0f}ms · 纯计算 {i:.0f}ms · " f"清空全部 {len(SYMS)} 币 {clear:.0f}ms" f"(worker {self.workers} / 核 {CORES})", flush=True) print(f" → {self._compute_verdict(q, i, clear)}", 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 def _compute_verdict(self, q: float, i: float, clear: float) -> str: """给出唯一可行的出路,而不是「哪一项数字更大」。 旧版比逐根的 q 与 i,结构上错了两处: 1. 判据错。真正要紧的是**一个收盘时刻清空所有币要多久**(clear), 不是单币的 q 或 i。所有币同一秒收盘,币数超过 worker 数时后面的 币必然串行等待,而这笔代价不出现在任何单根的 q 或 i 里。 2. 出路错。「排队为主 → 加核」只在还有空闲核时成立。worker 已等于 核数时,加 worker 不会增加吞吐——CPU 密集的活变不出来,只会把 等待从 queue_ms 挪到 inner_ms。十币实测正是如此:inner 被争抢从 144ms 抬到 192ms,反而超过 queue 135ms,于是判定落到「量级已低、 无需优化」,而此时最后一个币已经落在 1376ms。 所以币数超过核数时,唯一的出路是压单币耗时(增量计算),不是加 worker。 """ if not np.isfinite(clear): return "样本不足,暂不判定" if clear < 300: return f"清空 {clear:.0f}ms,宽裕,无需优化" if self.workers < CORES and q > i: return (f"排队为主且还有 {CORES - self.workers} 个空闲核 → " f"--workers 加到 {CORES}") if len(SYMS) > CORES: per = clear / max(len(SYMS) / max(self.workers, 1), 1) return (f"币数 {len(SYMS)} > 核数 {CORES},加 worker 无用(CPU 密集)" f"。唯一出路是压单币耗时 {per:.0f}ms → 走增量 " f"init_stream/append_bar(HANDOFF §5.5,实测约 3.7x)") return f"清空 {clear:.0f}ms 偏高,但币数未超核数,先查是否有别的争抢" 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: global SYMS ap = argparse.ArgumentParser() ap.add_argument("--hours", type=float, default=24.0) ap.add_argument("--workers", type=int, default=2) # 币数直接决定排队:所有币在同一秒收盘,worker 少于币数就必然排队, # 最后一个币的信号要等 ceil(n/worker) 轮计算。TRX 不在默认池里—— # 实盘口径 208 天只有 5 笔,ATR 门控几乎全刷掉(HANDOFF §step48) ap.add_argument("--syms", default=",".join(SYMS), help="逗号分隔。十币池:BTC,ETH,SOL,BNB,XRP,DOGE,ADA," "AVAX,LINK,LTC") a = ap.parse_args() SYMS = tuple(s.strip().upper() for s in a.syms.split(",") if s.strip()) # 进程池必须在事件循环和任何 WS 连接之前建好:fork 一个已带活跃 socket # 的进程会把连接状态一起复制过去,后果不可预测 with ProcessPoolExecutor(max_workers=a.workers) as pool: asyncio.run(main_async(a.workers, a.hours, pool)) if __name__ == "__main__": main()