From 181bca303ff97a2eca0a2bd95fa5fed4e1851a4e Mon Sep 17 00:00:00 2001 From: jack Date: Fri, 28 Aug 2026 01:57:59 +0800 Subject: [PATCH] =?UTF-8?q?=E5=BD=B1=E5=AD=90=E6=B5=8B=E9=87=8F=E6=94=B9?= =?UTF-8?q?=E7=94=A8=E6=A1=86=E6=9E=B6=E5=90=83=E5=8D=95=E5=8E=9F=E8=AF=AD?= =?UTF-8?q?=EF=BC=8C=E5=B9=B6=E8=A1=A5=E9=BD=90=E5=AE=B9=E9=87=8F=E4=B8=8E?= =?UTF-8?q?=20maker=20=E6=88=90=E4=BA=A4=E7=8E=87=E4=B8=A4=E9=A1=B9?= =?UTF-8?q?=E6=B5=8B=E7=AE=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 吃单查询换成 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 --- .gitignore | 3 + research/live/latency_compare.py | 15 +- research/live/patched_candles.py | 61 ++++ research/live/penetration.py | 152 ++++++++ research/live/probe_depth.py | 68 ++++ research/live/shadow_depth.py | 287 +++++++++++++++ research/live/shadow_hb.py | 543 +++++++++++++++++++++++----- research/live/shadow_report.py | 232 +++++++++--- research/live/shadow_signal.py | 111 ++++-- research/live/verify_signal_path.py | 205 +++++++++++ 10 files changed, 1503 insertions(+), 174 deletions(-) create mode 100644 research/live/penetration.py create mode 100644 research/live/probe_depth.py create mode 100644 research/live/shadow_depth.py create mode 100644 research/live/verify_signal_path.py diff --git a/.gitignore b/.gitignore index 2861ef6..d839a37 100644 --- a/.gitignore +++ b/.gitignore @@ -42,3 +42,6 @@ venv/ # Local tooling .gstack/ +research/out/*.jsonl.gz +research/out/penetration.csv +research/out/shadow_*.csv diff --git a/research/live/latency_compare.py b/research/live/latency_compare.py index a0e1fe1..cfaf060 100644 --- a/research/live/latency_compare.py +++ b/research/live/latency_compare.py @@ -5,9 +5,11 @@ 的标准差。入场方向上还有系统性追价(信号触发往往伴随同向动量),所以随机 漂移只是下限,真实成本更高——这也是为什么最终仍要用真实盘口测滑点。 -余量(bitget_baseline.py 得出,Bitget 原生基线减去手续费后剩下的空间): - BTC -0.13bp ETH +4.02bp SOL +2.92bp -BTC 本就为负,留着只作延迟测量的参照物,不作交易标的。 +预算从 `lib/shadow_budget` import,不在这里写死。曾经写死的 +`{BTC: -0.13, ETH: 4.02, SOL: 2.92}` 是错的——那是 `bitget_baseline.py` +按「毛均 − 6bp 双边 taker」算的,六处口径叠加(费率档位记高、余量没除 +taker 名义额、用了 5m~30m 的出场参数、只有同向没有阶梯与 ATR 门控)。 +正确值是 8.58 / 20.64 / 16.83,ETH 差了五倍。 .venv/bin/python research/live/latency_compare.py """ @@ -23,8 +25,9 @@ HERE = Path(__file__).resolve().parent RESEARCH = HERE.parent sys.path.insert(0, str(RESEARCH)) +from lib.shadow_budget import budget_of # noqa: E402 + OUT = RESEARCH / "out" -BUDGET_BP = {"BTC": -0.13, "ETH": 4.02, "SOL": 2.92} SYMS = ("BTC", "ETH", "SOL") @@ -67,8 +70,8 @@ def describe(df: pd.DataFrame, name: str, vols: dict) -> None: med, p90 = float(np.median(v)), float(np.percentile(v, 90)) vol = vols.get(s) d = drift_bp(med, vol) if vol else float("nan") - b = BUDGET_BP[s] - share = f"{d / b * 100:.0f}%" if b > 0 else "—(负)" + b = budget_of(s) + share = f"{d / b * 100:.0f}%" if np.isfinite(b) and b > 0 else "—" print(f"{s:<5}{len(v):>5}{med:>9.0f}{p90:>9.0f}{v.max():>9.0f}" f"{d:>12.2f}{b:>9.2f}{share:>9}") diff --git a/research/live/patched_candles.py b/research/live/patched_candles.py index 61fe510..586aeb8 100644 --- a/research/live/patched_candles.py +++ b/research/live/patched_candles.py @@ -14,9 +14,21 @@ 再把最后一根交给基类走正常的 append 流程。 已向上游反馈前,本地用子类覆盖,不改动镜像。 + +## 为什么必须有启动断言 + +子类覆盖的失效方式是**静默**的:上游若把 `_parse_websocket_message` 改名、 +或改走别的钩子,我们的覆盖就成了死代码,行情悄悄退回慢 1.06 秒,不崩、 +不报错、不留日志,只会让收益慢慢变差,几周后才从统计里看出来。 + +`assert_patch_effective()` 不做名字检查——名字对不上未必失效,名字对得上 +也未必生效。它喂一条合成的两元素消息,直接验证行为:基类返回首元素(bug +仍在、覆盖仍有必要),子类返回末元素(覆盖确实生效)。再加一条源码检查 +确认基类的收包循环还在调这个钩子。任一不满足就在启动时抛错。 """ from __future__ import annotations +import inspect from typing import Any, Dict, Optional import numpy as np @@ -24,6 +36,7 @@ import numpy as np from hummingbot.data_feed.candles_feed.bitget_perpetual_candles import ( BitgetPerpetualCandles, ) +from hummingbot.data_feed.candles_feed.candles_base import CandlesBase def _row_to_dict(row: list, ensure_s) -> Dict[str, Any]: @@ -63,3 +76,51 @@ class PatchedBitgetPerpetualCandles(BitgetPerpetualCandles): d["volume"], d["quote_asset_volume"], d["n_trades"], d["taker_buy_base_volume"], d["taker_buy_quote_volume"]] ).astype(float) + + +# 换根时 Bitget 推的就是这个形状:[上一根, 新一根] +_PROBE = { + "action": "update", + "arg": {"instType": "USDT-FUTURES", "channel": "candle1m", + "instId": "BTCUSDT"}, + "data": [ + ["1700000040000", "1", "1", "1", "1", "1", "1", "1"], + ["1700000100000", "2", "2", "2", "2", "2", "2", "2"], + ], +} + + +def assert_patch_effective() -> None: + """启动即验证覆盖真的生效,否则抛错。让静默失效变成启动失败。""" + src = inspect.getsource(CandlesBase._process_websocket_messages_task) + if "_parse_websocket_message" not in src: + raise RuntimeError( + "上游收包循环已不再调用 _parse_websocket_message," + "patched_candles 的覆盖失效。需重新定位钩子后再启动。") + + stock = BitgetPerpetualCandles("BTC-USDT", "1m", 20) + ours = PatchedBitgetPerpetualCandles("BTC-USDT", "1m", 20) + got_stock = stock._parse_websocket_message(_PROBE) + got_ours = ours._parse_websocket_message(_PROBE) + + head_ts = stock.ensure_timestamp_in_seconds(int(_PROBE["data"][0][0])) + tail_ts = stock.ensure_timestamp_in_seconds(int(_PROBE["data"][-1][0])) + + if not got_ours or int(got_ours["timestamp"]) != int(tail_ts): + raise RuntimeError( + f"覆盖未生效:子类返回 {got_ours and got_ours.get('timestamp')}," + f"应为末元素 {tail_ts}。") + if got_stock and int(got_stock["timestamp"]) == int(tail_ts): + # 上游自己修好了。此时覆盖无害但已多余,明确说出来,免得以后 + # 有人以为那 1.06 秒还是靠我们拿回来的 + print(" [补丁] 上游已自行修正换根解析,本地覆盖现为冗余,可移除", + flush=True) + elif not got_stock or int(got_stock["timestamp"]) != int(head_ts): + raise RuntimeError( + f"基类行为与预期不符:返回 " + f"{got_stock and got_stock.get('timestamp')}," + f"既非首元素 {head_ts} 也非末元素 {tail_ts}。" + f"上游改了解析逻辑,补丁的前提需重新确认。") + else: + print(f" [补丁] 覆盖生效:基类取首元素 {int(head_ts)}、" + f"本地取末元素 {int(tail_ts)}", flush=True) diff --git a/research/live/penetration.py b/research/live/penetration.py new file mode 100644 index 0000000..7dd9fb3 --- /dev/null +++ b/research/live/penetration.py @@ -0,0 +1,152 @@ +"""止盈位被首次触及那一根,价格穿透了多深。 + +为什么要这个数:影子成交流显示,限价单若正好落在某根的最高价,该价位之上 +的主动买成交额只有几十到几千美元——对十万量级的仓位等于不成交。但那是 +最坏情形。真实成交率取决于**止盈位被穿透了多深**:若价格一路冲过目标, +成交没问题;若只是上影线点一下就回落,就成交不了。 + +这个分布不需要再采数据,历史 K 线里就有:给定入场价与 ATR,目标位是 +`entry + T×ATR`,找到首次 `high ≥ target` 的那根,穿透深度就是 +`high − target`。把它折成「占该根价格区间的比例」,就能直接对上成交流 +那条「≥ 限价的成交额 vs 限价在区间中的位置」曲线。 + +口径与 lib/exit_model.walk_exits 对齐:入场取信号次根开盘价,ATR 取信号 +根的 Wilder ATR-14,上限 48 根。 + +**取样方式的局限**:这里用全体 K 线做候选入场点,而非真实的三滤网信号。 +真实信号是按结构条件挑出来的,入场时刻可能与波动率状态相关。以 ATR 归一 +后的穿透深度对波动率状态应当不敏感,但要精确到信号级别,得重跑一次 +step42 的缠论链路(单币约 370s、峰值 24.5GB)。 +""" +from __future__ import annotations + +import argparse +from pathlib import Path + +import numpy as np +import pandas as pd + +SCALE_AT = 3.0 # 分批减仓位 +RUNNER = 8.0 # 剩余半仓目标 +MAX_BARS = 48 + + +def wilder_atr(high, low, close, period: int = 14) -> np.ndarray: + """与 chanlun.indicators.ta.ATR 逐位一致的 Wilder ATR。""" + n = high.size + out = np.full(n, np.nan) + if n <= period: + return out + prev_close = close[:-1] + tr = np.maximum.reduce([high[1:] - low[1:], + np.abs(high[1:] - prev_close), + np.abs(low[1:] - prev_close)]) + sm = np.empty(tr.size) + sm[period - 1] = tr[:period].mean() + a = 1.0 / period + for k in range(period, tr.size): + sm[k] = sm[k - 1] + a * (tr[k] - sm[k - 1]) + out[period:] = sm[period - 1:] + return out + + +def penetration(df: pd.DataFrame, target_atr: float, + stride: int = 1) -> pd.DataFrame: + """对每个候选入场点,求首次触及 `target_atr` 时的穿透深度。 + + 只统计**触及了**的那些(未触及的属止损或超时出场,不涉及 maker 腿)。 + """ + high = df["high"].to_numpy(float) + low = df["low"].to_numpy(float) + close = df["close"].to_numpy(float) + open_ = df["open"].to_numpy(float) + atr = wilder_atr(high, low, close) + n = len(df) + rows = [] + for i in range(20, n - MAX_BARS - 2, stride): + a = atr[i] + if not np.isfinite(a) or a <= 0: + continue + e = i + 1 + entry = open_[e] + # 多头:目标在上方。空头对称,穿透深度分布按对称性等价,故只算一边 + target = entry + target_atr * a + end = e + MAX_BARS + seg_hi = high[e:end + 1] + hit = np.flatnonzero(seg_hi >= target) + if not hit.size: + continue + j = e + int(hit[0]) + rng = high[j] - low[j] + if rng <= 0: + continue + pen = high[j] - target + rows.append({ + "bar": j, + # 限价在该根价格区间中的位置:0 = 正好在最高价(最坏), + # 1 = 在最低价(该根全部成交都在限价之上) + "k": min(1.0, pen / rng), + "pen_bp": pen / target * 1e4, + "pen_atr": pen / a, + "range_bp": rng / target * 1e4, + }) + return pd.DataFrame(rows) + + +def report(sym: str, df: pd.DataFrame) -> dict: + print(f"\n{'=' * 68}\n{sym} 共 {len(df):,} 根 1m") + out = {} + for tgt, name in ((SCALE_AT, f"减仓位 {SCALE_AT:g}ATR"), + (RUNNER, f"目标位 {RUNNER:g}ATR")): + p = penetration(df, tgt) + if p.empty: + print(f" {name}: 无触及样本") + continue + k = p["k"].to_numpy() + print(f"\n {name} · 触及 {len(p):,} 次") + print(f" 穿透深度 中位 {p['pen_bp'].median():.2f}bp " + f"({p['pen_atr'].median():.2f} ATR) · " + f"P25 {p['pen_bp'].quantile(.25):.2f}bp · " + f"P75 {p['pen_bp'].quantile(.75):.2f}bp") + print(f" 限价在区间中的位置 k(0=正好在最高价,越大越靠下越易成交)") + print(f" 中位 {np.median(k):.3f} · P10 {np.percentile(k, 10):.3f}" + f" · P25 {np.percentile(k, 25):.3f}" + f" · P75 {np.percentile(k, 75):.3f}") + for thr in (0.05, 0.10, 0.25, 0.50): + print(f" k ≤ {thr:.2f}(限价挤在该根顶部 {thr * 100:.0f}% 内):" + f"{float((k <= thr).mean()) * 100:5.1f}% 的触及") + out[tgt] = p + return out + + +def main() -> None: + ap = argparse.ArgumentParser() + ap.add_argument("--syms", default="BTC,ETH,SOL") + ap.add_argument("--cache", default="research/live/cache") + ap.add_argument("--save", default="research/out/penetration.csv") + a = ap.parse_args() + + root = Path(a.cache) + allp = [] + for sym in a.syms.split(","): + cands = sorted(root.glob(f"bitget_{sym}_1m_*.feather"), + key=lambda p: p.stat().st_size, reverse=True) + if not cands: + print(f"{sym}: 找不到 1m 缓存,跳过") + continue + df = pd.read_feather(cands[0]) + got = report(sym, df) + for tgt, p in got.items(): + p = p.copy() + p["sym"], p["target_atr"] = sym, tgt + allp.append(p) + if allp: + out = pd.concat(allp, ignore_index=True) + Path(a.save).parent.mkdir(parents=True, exist_ok=True) + out.to_csv(a.save, index=False) + print(f"\n已存 {a.save}({len(out):,} 行)," + f"供 shadow_depth.py 合并成交量曲线") + + +if __name__ == "__main__": + main() diff --git a/research/live/probe_depth.py b/research/live/probe_depth.py new file mode 100644 index 0000000..4a08b63 --- /dev/null +++ b/research/live/probe_depth.py @@ -0,0 +1,68 @@ +"""真实 Bitget 盘口在 BOOK_DEPTH 档内能不能吃下各个名义额档位。 + +shadow_hb 把吃单换成了框架的 get_vwap_for_volume,深度不足时它返回 nan、 +整行标 depth_ok=0。所以「档数够不够」直接决定某个仓位档会不会整段丢失, +不是个可以事后补救的参数——先量出来再定 BOOK_DEPTH。 +""" +from __future__ import annotations + +import asyncio +import sys + +import numpy as np + +sys.path.insert(0, "/repo/research/live") + +from shadow_hb import BOOK_DEPTH, NOTIONALS, SYMS, book_from + + +async def run() -> None: + from hummingbot.connector.derivative.bitget_perpetual.bitget_perpetual_derivative import ( + BitgetPerpetualDerivative, + ) + conn = 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 conn.start_network() + print(f"连接器已启动,等盘口(档数上限 {BOOK_DEPTH})") + for _ in range(60): + await asyncio.sleep(1) + try: + if all(conn.get_order_book(f"{s}-USDT") is not None for s in SYMS): + break + except Exception: + continue + + for s in SYMS: + ob = conn.get_order_book(f"{s}-USDT") + 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)))]) + mid = (bids[0][0] + asks[0][0]) / 2.0 + snap = book_from(bids, asks) + ask_notional = float((asks[:, 0] * asks[:, 1]).sum()) + print(f"\n{s} 中价 {mid:.2f} · 取到 {len(asks)} 档 · " + f"卖盘 {len(asks)} 档合计 {ask_notional:,.0f} USDT") + print(f" 最深一档距中价 " + f"{(asks[-1][0] / mid - 1) * 1e4:.1f}bp") + for notional in NOTIONALS: + base = notional / mid + r = snap.get_vwap_for_volume(True, base) + px = float(r.result_price) + ok = float(r.result_volume) >= base * 0.999 + if ok: + print(f" 名义 {notional:>7,.0f} → {base:.6f} 币 · " + f"冲击 {(px / mid - 1) * 1e4:6.2f}bp · 吃得下") + else: + print(f" 名义 {notional:>7,.0f} → {base:.6f} 币 · " + f"深度不足,仅 {r.result_volume:.6f} 币 · " + f"这一档会整段丢失") + await conn.stop_network() + + +if __name__ == "__main__": + asyncio.run(run()) diff --git a/research/live/shadow_depth.py b/research/live/shadow_depth.py new file mode 100644 index 0000000..0c92bb6 --- /dev/null +++ b/research/live/shadow_depth.py @@ -0,0 +1,287 @@ +"""从落盘的完整盘口与成交流,算资金容量与 maker 腿成交率。 + +这两个数都不该等实盘暴露: + +**容量**。预算 20bp 意味着存在一个资金上限,超过它策略就不工作。既然完整 +深度已落盘,任意仓位的冲击都能重算——一次采集回答所有资金量级,换个规模 +不必重测一周。绑定约束是**薄盘时段**而非中位盘口,所以按分位数报。 + +**maker 成交率**。回测假设 3ATR / 8ATR 的限价单全额成交。深度回答不了这个 +问题:深度说的是「现在挂着多少」,成交率问的是「之后打过来多少」。只有 +成交流能回答,而且买卖必须分开——多头在 3ATR 挂卖出,靠主动买盘成交。 + +读 gzip 时必须容忍末尾成员不完整:采集进程还在写,最后一个 gzip 成员没有 +结尾标记,直接遍历会在文件尾抛 EOFError 而丢掉**全部**已读记录。 +""" +from __future__ import annotations + +import argparse +import gzip +import json +from pathlib import Path + +import numpy as np +import pandas as pd + + +def out_dir() -> Path: + p = Path("/out") + return p if p.is_dir() else Path(__file__).resolve().parents[1] / "out" + + +def read_jsonl_gz(path: Path): + """逐行读 gzip JSONL,末尾截断则静默停止。 + + 采集仍在进行时,最后一个 gzip 成员缺结尾标记;不接这个异常的话, + 整个分析会因为文件尾而失败,前面几万条完好记录一起丢掉。 + """ + if not path.exists(): + return + n_ok = 0 + try: + with gzip.open(path, "rt", encoding="utf-8") as fh: + for line in fh: + try: + rec = json.loads(line) + except json.JSONDecodeError: + break # 半行,说明写到这里被打断 + n_ok += 1 + yield rec + except (EOFError, OSError, gzip.BadGzipFile): + # 采集进程正在写,尾部不完整属正常 + pass + + +def impact_bp(levels: list, notional: float, mid: float) -> float | None: + """吃掉 notional 计价币后的加权均价相对中间价,bp。深度不足返回 None。""" + need = notional + cost = 0.0 + qty = 0.0 + for px, amt in levels: + avail = px * amt + take = min(avail, need) + q = take / px + cost += q * px + qty += q + need -= take + if need <= 1e-9: + break + if need > 1e-9 or qty <= 0: + return None + return (cost / qty / mid - 1.0) * 1e4 + + +def capacity(books_path: Path, budgets: dict[str, float], + pctl: float = 10.0) -> None: + """报各币的深度曲线与「冲击吃掉预算多少」的资金上限。""" + grid = [1e4, 5e4, 1e5, 2e5, 3.2e5, 5.3e5, 1e6, 2e6, 5e6] + per: dict[str, dict[float, list[float]]] = {} + n = 0 + for r in read_jsonl_gz(books_path): + asks, bids = r["asks"], r["bids"] + if not asks or not bids: + continue + mid = (asks[0][0] + bids[0][0]) / 2.0 + d = per.setdefault(r["sym"], {g: [] for g in grid}) + for g in grid: + v = impact_bp(asks, g, mid) + d[g].append(np.nan if v is None else v) + n += 1 + + if not n: + print("没有盘口快照,先跑采集") + return + + print(f"\n########## 资金容量 ##########") + print(f" 基于 {n:,} 份完整盘口快照(单边买入方向)\n") + for sym, d in per.items(): + b = budgets.get(sym) + print(f" {sym} 预算 {b:.2f}bp" if b else f" {sym}") + print(f" {'名义额':>12} {'冲击中位':>10} {'冲击P90':>10} " + f"{'吃满深度率':>10} {'占预算':>8}") + for g in grid: + a = np.array(d[g], dtype=float) + fill = float(np.isfinite(a).mean()) + if fill == 0: + print(f" {g:>12,.0f} {'—— 50 档吃不下 ——':>30}") + continue + med = float(np.nanmedian(a)) + p90 = float(np.nanpercentile(a, 90)) + share = f"{med / b * 100:6.1f}%" if b else " na" + print(f" {g:>12,.0f} {med:>10.2f} {p90:>10.2f} " + f"{fill * 100:>9.1f}% {share:>8}") + if b: + # 上限:冲击的 P90(薄盘时段)吃掉预算三成为止。三成是留给 + # 漂移与价差的余地——它们才是主项,冲击不该独占预算 + cap = None + for g in grid: + a = np.array(d[g], dtype=float) + if not np.isfinite(a).any(): + break + if float(np.nanpercentile(a, 90)) > b * 0.30: + break + cap = g + if cap is None: + print(f" → 连最小档 {grid[0]:,.0f} 的薄盘冲击都超预算三成") + else: + print(f" → 资金上限约 {cap:,.0f} USDT" + f"(薄盘 P90 冲击 ≤ 预算 30%)") + print() + + +def maker_fill(tape_path: Path, mults=(3.0, 8.0), + notionals=(1e5, 3.2e5, 5.3e5)) -> None: + """限价单挂在离场目标位,本根内有多少主动量打到那里。 + + 这里只回答「量够不够」。真实成交还要看排队位置——我们的单排在该价位 + 已有挂单之后,所以这是**上界**:量不够则必然不能全成交,量够也未必成交。 + """ + rows = list(read_jsonl_gz(tape_path)) + if not rows: + print("没有成交流数据,先跑采集") + return + print(f"\n########## maker 腿成交量上界 ##########") + print(f" 基于 {len(rows):,} 根的逐价位成交聚合") + print(f" 多头在目标位挂卖出,成交靠主动**买**盘,故只计买方向\n") + + # 限价单只能被**价格 ≥ 限价**的主动买成交打到。而止盈位被触及的那一根, + # 限价往往就落在该根价格区间的顶部——最高价刚好碰到目标位是最典型的 + # 情形。所以按「限价距最高价多近」分层:depth=0 表示限价正好在最高价 + # (只有打在最高价那一档的量算数),depth=0.25 表示限价在区间顶部 25% 处 + depths = (0.0, 0.10, 0.25, 1.0) + per: dict[str, dict[float, list[float]]] = {} + for r in rows: + buys = {float(p): v for p, v in r["buys"].items()} + d = per.setdefault(r["sym"], {k: [] for k in depths}) + if not buys: + for k in depths: + d[k].append(0.0) + continue + hi, lo = max(buys), min(buys) + rng = hi - lo + for k in depths: + floor_px = hi - k * rng + d[k].append(sum(p * v for p, v in buys.items() if p >= floor_px)) + + for sym, d in per.items(): + print(f" {sym} ≥ 限价的主动买成交额(USDT),按限价所处位置分层") + print(f" {'限价位置':>16} {'中位':>12} {'P25':>12} " + + " ".join(f"{n:>9,.0f}全仓" for n in notionals)) + for k in depths: + a = np.array(d[k], dtype=float) + where = ("正好在最高价" if k == 0 else + "整根全部成交" if k == 1.0 else + f"区间顶部 {k * 100:.0f}%") + cells = " ".join(f"{float((a >= n).mean()) * 100:8.1f}%" + for n in notionals) + print(f" {where:>16} {np.median(a):>12,.0f} " + f"{np.percentile(a, 25):>12,.0f} {cells}") + print() + print(" 「正好在最高价」那一行才是止盈被刚好触及时的真实处境;") + print(" 「整根全部成交」是最宽松的上界。两行差多少,就是回测那个") + print(" 「限价单全额成交」假设虚了多少。而且这仍未计排队——我们的单") + print(" 排在该价位既有挂单之后,所以真实成交率比表里更低") + + +def tape_shape(tape_path: Path, kgrid: np.ndarray) -> dict[str, np.ndarray]: + """成交流给「形状」:一根的主动买成交额里,有多少比例落在区间顶部 k 之内。 + + 形状与规模分开是为了绕开成交流样本小的限制——形状是微观结构性质, + 几十根就相当稳定;规模(每根成交多少钱)则由 210 天历史成交量提供。 + """ + acc: dict[str, list[np.ndarray]] = {} + for r in read_jsonl_gz(tape_path): + buys = {float(p): v for p, v in r["buys"].items()} + if len(buys) < 2: + continue + hi, lo = max(buys), min(buys) + rng = hi - lo + if rng <= 0: + continue + tot = sum(p * v for p, v in buys.items()) + if tot <= 0: + continue + frac = np.array([sum(p * v for p, v in buys.items() + if p >= hi - k * rng) / tot for k in kgrid]) + acc.setdefault(r["sym"], []).append(frac) + return {s: np.mean(np.vstack(v), axis=0) for s, v in acc.items() if v} + + +def composite_fill(tape_path: Path, pen_path: Path, cache: Path, + notionals=(1e5, 3.2e5, 5.3e5)) -> None: + """把穿透深度分布与成交量曲线合并,出真实 maker 成交率。""" + if not pen_path.exists(): + print("\n没有 penetration.csv,先跑 penetration.py") + return + kgrid = np.linspace(0.0, 1.0, 51) + shape = tape_shape(tape_path, kgrid) + if not shape: + print("\n成交流样本不足,无法定形状") + return + pen = pd.read_csv(pen_path) + + print(f"\n\n########## maker 腿真实成交率 ##########") + print(f" 穿透深度分布(历史 63 万次触及)× 每根成交额(210 天)") + print(f" × 区间内成交分布形状(影子成交流)\n") + + for sym in sorted(shape): + cands = sorted(cache.glob(f"bitget_{sym}_1m_*.feather"), + key=lambda p: p.stat().st_size, reverse=True) + if not cands: + continue + bars = pd.read_feather(cands[0]) + # 每根的主动买成交额。取总成交额的一半——买卖大致均衡,且这与 + # 成交流实测的买卖比一致 + bar_notional = (bars["volume"].to_numpy(float) + * bars["close"].to_numpy(float)) * 0.5 + bar_notional = bar_notional[np.isfinite(bar_notional) + & (bar_notional > 0)] + f = shape[sym] + for tgt in sorted(pen["target_atr"].unique()): + k = pen[(pen["sym"] == sym) + & (pen["target_atr"] == tgt)]["k"].to_numpy(float) + if not k.size: + continue + # 独立配对:穿透位置与该根成交额各自抽样。真实触及根多为放量根, + # 故此处偏**保守**(低估可成交量) + rng = np.random.default_rng(0) + m = 200_000 + ks = rng.choice(k, m) + ns = rng.choice(bar_notional, m) + avail = np.interp(ks, kgrid, f) * ns + print(f" {sym} · 目标 {tgt:g}ATR · 每根主动买额中位 " + f"{np.median(bar_notional):,.0f} USDT") + for nt in notionals: + full = float((avail >= nt).mean()) + half = float((avail >= nt / 2).mean()) + print(f" 仓位 {nt:>9,.0f}:全额成交 {full * 100:5.1f}%" + f" · 至少半额 {half * 100:5.1f}%" + f" · 可成交额中位 {np.median(avail):>10,.0f}") + # 成交率反推的资金上限。这才是绑定约束——它比冲击反推的上限 + # 低一到两个数量级,而后者才是通常被当作「容量」的那个数 + for want in (0.80, 0.90): + cap = float(np.quantile(avail, 1.0 - want)) + print(f" → 要 {want * 100:.0f}% 的止盈全额成交," + f"仓位须 ≤ {cap:,.0f} USDT") + print() + print(" 未计排队(我们的单排在该价位既有挂单之后),故仍是上界。") + print(" 回测把这些止盈按「全额成交在目标价」计,差多少就是收益虚多少") + + +def main() -> None: + ap = argparse.ArgumentParser() + ap.add_argument("--books", default=None) + ap.add_argument("--tape", default=None) + a = ap.parse_args() + d = out_dir() + from lib.shadow_budget import BUDGET_BP + tape = Path(a.tape) if a.tape else d / "shadow_tape.jsonl.gz" + capacity(Path(a.books) if a.books else d / "shadow_books.jsonl.gz", + BUDGET_BP) + maker_fill(tape) + composite_fill(tape, d / "penetration.csv", + Path(__file__).resolve().parent / "cache") + + +if __name__ == "__main__": + main() diff --git a/research/live/shadow_hb.py b/research/live/shadow_hb.py index 7cf1b91..7abe2c4 100644 --- a/research/live/shadow_hb.py +++ b/research/live/shadow_hb.py @@ -10,8 +10,25 @@ 为什么不能用 paper trade 的成交:那是 Hummingbot 自己的撮合模型模拟的, 测出来是模型行为不是市场行为。 -口径对齐 aggregate_robustness.py:回测假设成交在**信号次根的开盘价** +口径对齐 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 字段上。 ### 盘口滚动缓冲把延迟变成自变量 @@ -38,22 +55,32 @@ 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_DEPTH = 25 +BOOK_TOL_MS = 250 # 回查容差:10Hz 正常 ≤100ms,留些余量 +BOOK_DEPTH = 50 # 双边各 50 档,实测能撑 78 万~261 万美元 DELAYS_S = (0.5, 1.0, 2.0, 5.0) # 回查点 -NOTIONALS = (1_000.0, 5_000.0, 20_000.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"] @@ -62,6 +89,119 @@ def out_dir() -> Path: 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 的列结构。 @@ -77,38 +217,32 @@ def hb_to_research(cdf: pd.DataFrame) -> pd.DataFrame: .sort_values("timestamp").reset_index(drop=True) -def walk_book(levels: list[tuple[float, float]], notional: float - ) -> tuple[float, float]: - """吃单到 notional(计价币)为止,返回 (加权成交价, 实际吃到的额度)。 +def book_from(bids: np.ndarray, asks: np.ndarray): + """用缓冲里的快照临时搭一个 OrderBook,以便调用框架自带的吃单查询。 - 深度不足时返回吃到的部分,由调用方按 filled < notional 判断是否可信。 + 自己手写吃单曾经踩过两个坑,框架版都没有:`get_vwap_for_volume` 返回的 + 是真加权均价(市价单的实际成交价),而 `get_price_for_quote_volume` 返回 + 的是**边际价**,用后者会高估冲击;深度不足时框架返回 nan 而不是一个 + 「看起来很正常」的部分成交均价,靠 query_volume/result_volume 判断。 """ - if not levels: - return float("nan"), 0.0 - got = 0.0 - cost = 0.0 - qty = 0.0 - for px, sz in levels: - avail = px * sz - take = min(avail, notional - got) - if take <= 0: - break - q = take / px - cost += q * px - qty += q - got += take - if got >= notional - 1e-9: - break - if qty <= 0: - return float("nan"), 0.0 - return cost / qty, got + 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] @@ -118,12 +252,13 @@ class BookBuffer: d.popleft() def at(self, sym: str, t_ms: int) -> tuple | None: - best = None for snap in self.buf[sym]: if snap[0] >= t_ms: - best = snap - break - return best + if snap[0] - t_ms > BOOK_TOL_MS: + self.n_stale += 1 + return None + return snap + return None class Shadow: @@ -136,30 +271,84 @@ class Shadow: 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() # 追加模式:长跑期间若重启,已收集的样本不该被清掉 - p_sig = d / "shadow_signals.csv" - new_sig = not p_sig.exists() or p_sig.stat().st_size == 0 - self.f_sig = p_sig.open("a", newline="") - self.w_sig = csv.DictWriter(self.f_sig, fieldnames=[ - "sym", "kline_ts", "direction", "h1_agree", + 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", "notional", - "baseline_px", "mid", "best_px", "fill_px", "filled", + "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"]) - if new_sig: - self.w_sig.writeheader() - p_lat = d / "shadow_latency.csv" - new_lat = not p_lat.exists() or p_lat.stat().st_size == 0 - self.f_lat = p_lat.open("a", newline="") - self.w_lat = csv.DictWriter(self.f_lat, fieldnames=[ + 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"]) - if new_lat: - self.w_lat.writeheader() + "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) # ---------- 启动 ---------- @@ -167,7 +356,11 @@ class Shadow: from hummingbot.connector.derivative.bitget_perpetual.bitget_perpetual_derivative import ( BitgetPerpetualDerivative, ) - from patched_candles import PatchedBitgetPerpetualCandles + from patched_candles import (PatchedBitgetPerpetualCandles, + assert_patch_effective) + + # 覆盖失效是静默的(悄悄退回慢 1.06 秒,不报错),所以在启动就验一次 + assert_patch_effective() for s in SYMS: self.feeds_l[s] = PatchedBitgetPerpetualCandles( @@ -187,6 +380,7 @@ class Shadow: trading_required=False) await self.connector.start_network() print(" 连接器已启动,等盘口与历史回填", flush=True) + self._hook_trades() t0 = time.time() while time.time() - t0 < 600: @@ -200,6 +394,39 @@ class Shadow: 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") @@ -207,25 +434,41 @@ class Shadow: return None if ob is None: return None - bids = [(float(r.price), float(r.amount)) - for r, _ in zip(ob.bid_entries(), range(BOOK_DEPTH))] - asks = [(float(r.price), float(r.amount)) - for r, _ in zip(ob.ask_entries(), range(BOOK_DEPTH))] - if not bids or not asks: + # 存成 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]) - await asyncio.sleep(period) + 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]) @@ -239,7 +482,11 @@ class Shadow: if last[s] is not None and newest > last[s]: t_data = int(time.time() * 1000) kts = newest * 1000 if newest < 1e12 else newest - asyncio.create_task(self.on_bar(s, kts, t_data)) + # 刚收盘那根的成交聚合先落盘,再算信号 + 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) @@ -252,15 +499,28 @@ class Shadow: 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()) + df_h[NUM_COLS].values.tolist(), baseline) loop = asyncio.get_running_loop() from shadow_signal import compute_packed - res = await loop.run_in_executor(self.pool, compute_packed, payload) + 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, @@ -268,13 +528,18 @@ class Shadow: "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(res.get("hits", []))}) + "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 - hits = res.get("hits", []) if not hits: return if baseline is None or not np.isfinite(baseline): @@ -283,20 +548,73 @@ class Shadow: for h in hits: self.n_signal += 1 - print(f" ★ [{sym}] {kline_ts} 方向 {h['direction']:+d} " - f"h1_agree={h['h1_agree']} · 数据 {t_data - kline_ts}ms " + 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,此刻尚未发生;等它过去再一次性落盘 - asyncio.create_task( - self._record_later(sym, kline_ts, h, t_data, t_signal, baseline)) + self._spawn( + self._record_later(sym, kline_ts, h, t_data, t_signal, + baseline, atr_pct, lag_ok), + f"record {sym}") - async def _record_later(self, sym: str, kline_ts: int, hit: dict, - t_data: int, t_signal: int, baseline: float) -> None: - target = kline_ts + int(max(DELAYS_S) * 1000) + 500 - wait = target / 1000.0 - time.time() + 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) - self._record(sym, kline_ts, hit, t_data, t_signal, baseline) + + 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: """次根开盘价 = 回测假设的成交价。""" @@ -309,25 +627,56 @@ class Shadow: 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) -> None: - points = [("actual", t_signal - kline_ts)] - points += [(f"{d}s", int(d * 1000)) for d in DELAYS_S] + 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 points: - snap = self.books.at(sym, kline_ts + delay_ms) + 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 - _, bids, asks = snap - best_bid, best_ask = bids[0][0], asks[0][0] + book_ts, bids, asks = snap + best_bid, best_ask = float(bids[0][0]), float(asks[0][0]) mid = (best_bid + best_ask) / 2.0 - # 多头吃卖盘,空头吃买盘 - side = asks if d_sign > 0 else bids 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: - fill, filled = walk_book(side, notional) + # 名义额按基准价折成基础币再下单——真实委托是基础币计价的, + # 框架的 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 @@ -336,14 +685,21 @@ class Shadow: 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, - "notional": notional, "baseline_px": baseline, + "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(filled, 2), + "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)}) @@ -353,8 +709,22 @@ class Shadow: while not self.stop.is_set(): await asyncio.sleep(300) depth = {s: len(self.books.buf[s]) for s in SYMS} - print(f" [心跳] 已处理 {self.n_bars} 根 · 命中 {self.n_signal} 个 " - f"· 盘口缓冲 {depth}", flush=True) + 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() @@ -370,9 +740,12 @@ class Shadow: for f in list(self.feeds_l.values()) + list(self.feeds_h.values()): f.stop() await self.connector.stop_network() - self.f_sig.close() - self.f_lat.close() - print(f"\n收工:{self.n_bars} 根 · {self.n_signal} 个信号", flush=True) + 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: diff --git a/research/live/shadow_report.py b/research/live/shadow_report.py index f89466c..97fd5da 100644 --- a/research/live/shadow_report.py +++ b/research/live/shadow_report.py @@ -1,27 +1,44 @@ -"""影子交易器的报表:首日延迟门槛 + 滑点对延迟曲线。 +"""影子交易器的报表:延迟门槛 + 滑点对延迟曲线。 -两份产物对应计划里的两件事。 +### 判据常数一律从研究侧 import,不在这里写死 -### 延迟门槛(提前止损用) +`BUDGET_BP` 等常数留在 `research/lib/shadow_budget.py`。理由是这些数会变—— +2026-08-27 一天之内预算就动了四次(3.91 → 11.06 → 14.25 → 15.19bp), +费率也改了一次。本文件曾经写死过 BTC −0.13 / ETH 4.02 / SOL 2.92,那三个数 +由六处差异叠加而来(只有同向没有阶梯、费率按 6bp 双边 taker、TP=3.0、 +余量没除 taker 名义额、无 ATR 门控,且 BTC/ETH/SOL 恰是 ATR 最低的三个币)。 +正确值是 8.58 / 20.64 / 16.83——**ETH 差了五倍**。 -跑满 24 小时先看这个。若**总延迟已令预期漂移超过余量**,说明方案在这台机器 -上就不成立,不必等两周样本再停。余量取 bitget_baseline.py 的实测值: -BTC −0.13bp(本就为负,只作参照)、ETH +4.02bp、SOL +2.92bp。 +这件事要紧是因为下面的判读是自动停机开关:用 4.02 当 ETH 的预算,真实滑点 +只要到 2.4bp 就会报「需要压延迟或放弃」,会误杀一个可行的策略。 -漂移按随机游走折算:σ_1m · √(t/60)。这是下限——入场条件是「收盘突破转强」, -那一刻价格正朝我们方向跑,延迟造成的是系统性追价,不会正负抵消。所以实测 -滑点理应比这个折算值更差,两者对照本身就是个校验。 +### 什么时候能判什么 + +| | 一天的样本量 | 够不够 | +|---|---|---| +| 延迟 | 1440 根/币 | 够,统计上很厚 | +| 滑点 | 门控后 4~6 笔/天 | **不够**,判据要 30 笔以上,即一周起步 | + +所以首日只能判延迟和管道通不通。滑点那一节在样本不足时会明说。 + +### 延迟门槛(提前止损用) + +若**总延迟已令预期漂移超过预算**,说明方案在这台机器上就不成立,不必等 +两周样本再停。漂移按随机游走折算 σ_1m · √(t/60)。这是下限——入场条件是 +「收盘突破转强」,那一刻价格正朝我们方向跑,延迟造成的是系统性追价, +不会正负抵消。所以实测滑点理应比折算值更差,两者对照本身就是个校验。 ### 滑点对延迟曲线 把延迟当自变量:0.5s / 1s / 2s / 5s 各一个滑点值,外加「actual」= 本机实际 -算完的时刻。这样即便本机算得慢,也能读出「若延迟压到 X 秒,滑点是多少」, +算完的时刻。同信号内的受控对比,能直接读出「若延迟压到 X 秒,滑点是多少」, 决策不被当前实现拖累。 .venv/bin/python research/live/shadow_report.py """ from __future__ import annotations +import sys from pathlib import Path import numpy as np @@ -29,8 +46,16 @@ import pandas as pd HERE = Path(__file__).resolve().parent OUT = HERE.parent / "out" +sys.path.insert(0, str(HERE.parent)) + +from lib.shadow_budget import ( # noqa: E402 + ATR_GATE_BP, BUDGET_PORTFOLIO_2026, LAG_ALARM_MS, budget_of, lag_healthy, + verdict, +) + SYMS = ("BTC", "ETH", "SOL") -BUDGET_BP = {"BTC": -0.13, "ETH": 4.02, "SOL": 2.92} +ORDER = ["0.5s", "1.0s", "2.0s", "5.0s", "actual"] +MIN_N = 30 # 滑点判据的最低笔数,低于此只报数不下结论 def vol_bp() -> dict[str, float]: @@ -46,13 +71,38 @@ def vol_bp() -> dict[str, float]: return v +def lag_health(lat: pd.DataFrame) -> None: + """运行时 lag 探针的回看。补丁后实测 506~642ms,理论下限约 500ms。""" + print("\n\n########## 二、lag 探针(>%.0fms 该停开仓)##########" + % LAG_ALARM_MS) + print(f"{'币':<5}{'根数':>6}{'中位ms':>9}{'P90ms':>8}{'最差30根中位':>14}" + f"{'超阈根数':>10}{'判定':>8}") + for s in SYMS: + g = lat[lat["sym"] == s].sort_values("kline_ts") + if g.empty: + continue + x = g["lag_data_ms"].to_numpy(float) + roll = pd.Series(x).rolling(30).median() + worst = float(np.nanmax(roll)) if roll.notna().any() else float("nan") + ok = lag_healthy(x) + print(f"{s:<5}{len(g):>6}{np.median(x):>9.0f}" + f"{np.percentile(x, 90):>8.0f}{worst:>14.0f}" + f"{int((x > LAG_ALARM_MS).sum()):>10}" + f"{'健康' if ok else '退化':>8}") + if "lag_ok" in lat.columns: + bad = int((lat["lag_ok"] == 0).sum()) + if bad: + print(f"\n 采集期间有 {bad} 根被判不健康,那些根上的信号" + f"(lag_ok=0)在真实运行下不会开仓,统计时应排除") + + def latency_gate(lat: pd.DataFrame, vols: dict) -> None: print("########## 一、延迟门槛 ##########") span_h = (lat["t_close_ms"].max() - lat["t_close_ms"].min()) / 3.6e6 print(f"样本跨度 {span_h:.1f} 小时 · 共 {len(lat)} 根\n") print(f"{'币':<5}{'根数':>6}{'数据ms':>9}{'计算ms':>9}{'总延迟ms':>10}" - f"{'P90ms':>8}{'折算漂移bp':>12}{'余量bp':>9}{'占余量':>9}") - verdicts = {} + f"{'P90ms':>8}{'折算漂移bp':>12}{'预算bp':>9}{'占预算':>9}") + rows = {} for s in SYMS: g = lat[lat["sym"] == s] if g.empty: @@ -63,48 +113,110 @@ def latency_gate(lat: pd.DataFrame, vols: dict) -> None: p90 = float(np.percentile(g["lag_signal_ms"], 90)) vol = vols.get(s) drift = vol * np.sqrt(t / 60_000) if vol else float("nan") - b = BUDGET_BP[s] - share = drift / b if b > 0 else float("nan") - verdicts[s] = (drift, b) - txt = f"{share * 100:.0f}%" if b > 0 else "—(负)" + b = budget_of(s) + share = drift / b if np.isfinite(b) and b > 0 else float("nan") + rows[s] = (drift, b) + txt = f"{share * 100:.0f}%" if np.isfinite(share) else "—" print(f"{s:<5}{len(g):>6}{d:>9.0f}{c:>9.0f}{t:>10.0f}{p90:>8.0f}" f"{drift:>12.2f}{b:>9.2f}{txt:>9}") - print("\n判读:") - for s, (drift, b) in verdicts.items(): - if b <= 0: - print(f" {s}: 余量本就为负,不作交易标的,仅作延迟参照") + print("\n判读(预算来自 lib/shadow_budget,2026 年口径):") + for s, (drift, b) in rows.items(): + if not np.isfinite(b): + print(f" {s}: 当前环境预算不足,不作交易标的,仅作延迟参照") + elif not np.isfinite(drift): + # 不特判的话 nan > b 是 False,会一路落到「尚有空间」说反话 + print(f" {s}: 缺 1m 波动率缓存,折算不出漂移,无法判读") elif drift > b: - print(f" {s}: 折算漂移 {drift:.2f}bp 已超余量 {b:.2f}bp —— 停下改方案") + print(f" {s}: 折算漂移 {drift:.2f}bp 已超预算 {b:.2f}bp —— 停下改方案") elif drift > b * 0.6: - print(f" {s}: 折算漂移 {drift:.2f}bp 吃掉余量 {b:.2f}bp 的六成以上," + print(f" {s}: 折算漂移 {drift:.2f}bp 吃掉预算 {b:.2f}bp 的六成以上," f"需要压延迟或放弃") else: - print(f" {s}: 折算漂移 {drift:.2f}bp 对余量 {b:.2f}bp 尚有空间,继续收集") + print(f" {s}: 折算漂移 {drift:.2f}bp 对预算 {b:.2f}bp 尚有空间,继续收集") + + +def book_quality(sig: pd.DataFrame, drf: pd.DataFrame) -> None: + """回查到的盘口比目标时刻晚多少。晚太多的已在采集侧丢弃,这里做复核。""" + frames = [d for d in (sig, drf) if not d.empty and "book_lag_ms" in d] + if not frames: + return + x = pd.concat([d["book_lag_ms"] for d in frames]).astype(float) + print("\n\n########## 二·五、盘口回查质量 ##########") + print(f" 回查 {len(x)} 次 · 中位 {x.median():.0f}ms · " + f"P90 {np.percentile(x, 90):.0f}ms · 最大 {x.max():.0f}ms") + print(f" 10Hz 采样下这个值应在 0~100ms。它是所有延迟点的同向偏置," + f"不改变曲线形状,但要确认没有异常长尾") + + +def drift_split(sig: pd.DataFrame, drf: pd.DataFrame) -> None: + """条件漂移 vs 无条件漂移。两者的差就是「系统性追价」的大小。""" + print("\n\n########## 三、条件漂移 vs 无条件漂移 ##########") + if drf.empty: + print(" 尚无逐根漂移数据(shadow_drift.csv 由本轮起才开始记)") + return + print(" 无条件 = 每根 K 线,方向未知故取 |漂移|;" + "条件 = 信号根按下单方向定号") + print(f"\n{'延迟':<8}{'无条件n':>9}{'无条件|漂移|':>14}" + f"{'条件n':>7}{'条件漂移':>10}{'追价差':>9}") + cond = sig[(sig["notional"] == sig["notional"].min())] if not sig.empty \ + else sig + for lb in ORDER: + u = drf[drf["delay_label"] == lb]["drift_bp_long"].abs() + c = cond[cond["delay_label"] == lb]["drift_bp"] if not cond.empty \ + else pd.Series(dtype=float) + if u.empty and c.empty: + continue + um = u.mean() if not u.empty else float("nan") + cm = c.mean() if not c.empty else float("nan") + print(f"{lb:<8}{len(u):>9}{um:>14.2f}{len(c):>7}{cm:>10.2f}" + f"{cm - um:>9.2f}") + if len(cond) and len(cond[cond["delay_label"] == "1.0s"]) < MIN_N: + print(f"\n 条件侧样本不足 {MIN_N},差值还读不出方向") def slippage_curve(sig: pd.DataFrame) -> None: - print("\n\n########## 二、滑点对延迟曲线 ##########") + print("\n\n########## 四、滑点对延迟曲线 ##########") if sig.empty: print(" 尚无信号样本") return - n_sig = sig.groupby(["sym", "kline_ts", "direction"]).ngroups - n_agree = sig[sig["h1_agree"] == 1].groupby( - ["sym", "kline_ts", "direction"]).ngroups - print(f"信号总数 {n_sig}(其中 h1_agree=1 的 {n_agree} 个)\n") - order = ["0.5s", "1.0s", "2.0s", "5.0s", "actual"] - for scope, sub in (("全部信号", sig), - ("仅 h1_agree=1(主口径)", sig[sig["h1_agree"] == 1])): + def n_of(df): + return df.groupby(["sym", "kline_ts", "direction"]).ngroups + + has_flags = "pass_all" in sig.columns + if not has_flags: + print(" ⚠ 数据来自旧版采集(只有 h1_agree,无阶梯与 ATR 门控)。" + "这批不是我们要交易的那批信号,只能作管道验证,不能对预算判读。\n") + main_scope, main_name = sig[sig["h1_agree"] == 1], "仅 h1_agree=1(旧口径)" + else: + # depth_ok=0 是 25 档吃不满该仓位,均价按部分成交算会**低估**冲击 + ok = (sig["pass_all"] == 1) & (sig["lag_ok"] == 1) + if "depth_ok" in sig.columns: + thin = int((sig["depth_ok"] == 0).sum()) + ok &= sig["depth_ok"] == 1 + if thin: + print(f" ({thin} 行深度吃不满,已排除;这些行会低估冲击)") + print(f"信号总数 {n_of(sig)} · 同向 {n_of(sig[sig['h1_agree'] == 1])}" + f" · 同向+阶梯 " + f"{n_of(sig[(sig['h1_agree'] == 1) & (sig['ladder_ok'] == 1)])}" + f" · 三项全过 {n_of(sig[sig['pass_all'] == 1])}" + f" · 再要求 lag 健康 {n_of(sig[ok])}") + print(f"(门控阈值 ATR ≥ {ATR_GATE_BP:.0f}bp,是费率的函数不是市场常数)\n") + main_scope = sig[ok] + main_name = "三项滤网全过 + lag 健康 + 深度吃满(主口径)" + + scopes = [("全部信号(含不会下单的,仅作提前读数)", sig), + (main_name, main_scope)] + for scope, sub in scopes: if sub.empty: + print(f"--- {scope} ---\n 尚无样本\n") continue print(f"--- {scope} ---") print(f"{'延迟':<8}{'仓位':>9}{'n':>5}{'滑点均值bp':>12}" f"{'中位bp':>9}{'漂移bp':>9}{'价差bp':>9}{'冲击bp':>9}") - for lb in order: + for lb in ORDER: g0 = sub[sub["delay_label"] == lb] - if g0.empty: - continue for nt in sorted(sub["notional"].unique()): g = g0[g0["notional"] == nt] if g.empty: @@ -117,19 +229,36 @@ def slippage_curve(sig: pd.DataFrame) -> None: f"{g['impact_bp'].mean():>9.2f}") print() - print("--- 分币种(仅 h1_agree=1,仓位 5000)---") - m = sig[(sig["h1_agree"] == 1) & (sig["notional"] == 5000.0)] + print("--- 分币种判读(主口径,仓位 5000)---") + m = main_scope[main_scope["notional"] == 5000.0] if not main_scope.empty \ + else main_scope if m.empty: print(" 尚无样本") return - print(f"{'币':<5}{'延迟':<8}{'n':>5}{'滑点均值bp':>12}{'余量bp':>9}") + print(f"{'币':<5}{'延迟':<8}{'n':>5}{'滑点中位bp':>12}{'预算bp':>9} 判读") for s in SYMS: - for lb in order: + for lb in ORDER: g = m[(m["sym"] == s) & (m["delay_label"] == lb)] if g.empty: continue - print(f"{s:<5}{lb:<8}{len(g):>5}{g['slip_bp'].mean():>12.2f}" - f"{BUDGET_BP[s]:>9.2f}") + med = float(g["slip_bp"].median()) + b = budget_of(s) + note = verdict(med, s) if len(g) >= MIN_N \ + else f"n={len(g)} < {MIN_N},不下结论" + print(f"{s:<5}{lb:<8}{len(g):>5}{med:>12.2f}{b:>9.2f} {note}") + + n_main = len(m[m["delay_label"] == "actual"]) + if n_main < MIN_N: + print(f"\n ⚠ 主口径仅 {n_main} 笔。门控后约 4~6 笔/天/全部币种," + f"滑点判据要 {MIN_N} 笔以上——**一周起步**。首日只能判延迟和管道。") + print(f" 单币样本薄时可先看组合口径:2026 预算 {BUDGET_PORTFOLIO_2026}bp") + + +def _load(name: str) -> pd.DataFrame: + f = OUT / name + if not f.exists() or f.stat().st_size == 0: + return pd.DataFrame() + return pd.read_csv(f) def main() -> None: @@ -137,20 +266,17 @@ def main() -> None: print("1m 收益标准差(bp/分钟,Bitget 实测):" + " ".join(f"{s} {v:.2f}" for s, v in vols.items()) + "\n") - f_lat = OUT / "shadow_latency.csv" - if f_lat.exists() and f_lat.stat().st_size > 0: - lat = pd.read_csv(f_lat) - if not lat.empty: - latency_gate(lat, vols) - else: + lat = _load("shadow_latency.csv") + if lat.empty: print("尚无延迟数据") - - f_sig = OUT / "shadow_signals.csv" - if f_sig.exists() and f_sig.stat().st_size > 0: - sig = pd.read_csv(f_sig) - slippage_curve(sig) else: - print("\n尚无信号数据") + latency_gate(lat, vols) + lag_health(lat) + + sig, drf = _load("shadow_signals.csv"), _load("shadow_drift.csv") + book_quality(sig, drf) + drift_split(sig, drf) + slippage_curve(sig) if __name__ == "__main__": diff --git a/research/live/shadow_signal.py b/research/live/shadow_signal.py index c65ed79..e1c72fa 100644 --- a/research/live/shadow_signal.py +++ b/research/live/shadow_signal.py @@ -3,15 +3,30 @@ 单次调用约 0.26s 的纯 CPU,且 chanlun 是纯 Python 受 GIL 限制,放进 Hummingbot 的 asyncio 循环里会把行情处理一起卡住,所以必须隔离到独立进程。 -口径与 [aggregate_robustness.run_once] 逐行对齐:同样的 build_htf_zones → -find_fast_bsp3 → attach_htf_context(h1) 链路,同样的 h1_agree 过滤。两边 -必须一致,否则影子测出来的滑点没法和 3.9bp 预算对照。 +## 口径必须与预算同源,缺一项数就不可比 -与回测的唯一差别是这里只关心**最后一根已收盘 K 线**上有没有信号—— -实盘只能在当下下单,历史信号无意义。 +预算(`lib/shadow_budget.BUDGET_BP`)算在 step42 的这套滤网上, +本文件逐行对齐 `step42_exit_tp_1m.run_one`: -未过滤信号也一并返回:过滤后样本太稀,先用未过滤的当提前读数, -两者都记,靠 h1_agree 字段区分。 + 同向 h1_agree == 1 + 中枢阶梯 多头要求当前中枢整体高于前一个(zd > 前 zg),空头反之 + ATR 门控 atr_pct ≥ ATR_GATE_BP(当前 8bp) + +早先这里只有 h1_agree。缺阶梯与门控测的就不是我们要交易的那批信号, +而这一项改常数解决不了——必须改信号路径本身。 + +门控阈值是**费率的函数**不是市场常数(低 ATR 信号的毛质量反而最好, +断崖只在扣费后出现),换 VIP 档或换交易所要重扫,不要抄 8bp。 + +## 与回测的两点差别 + +其一,这里只关心**最后一根已收盘 K 线**上有没有信号——实盘只能在当下下单。 +其二,`atr_pct` 的分母取次根开盘价,与 `exit_model.walk_exits` 一致, +所以调用方必须把次根开盘价传进来。 + +未通过滤网的信号也一并返回并打上标志:过滤后样本很稀(门控后 8 个币 +合计约 38 笔/周),未过滤的可作提前读数。但**统计主口径只能用 pass_all**, +在我们根本不会下单的根上测滑点会把判据算宽。 """ from __future__ import annotations @@ -23,15 +38,20 @@ for _v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS"): os.environ.setdefault(_v, "1") -def compute(df_l, df_h) -> dict: +def compute(df_l, df_h, entry_px: float | None = None) -> dict: """在 df_l 的最后一根上找信号。df_l/df_h 都只含已收盘 K 线。 + entry_px 是次根开盘价(回测 entry_delay=1 的成交价),用作 atr_pct 的 + 分母。取不到时退回用信号根收盘价,并在返回里标 atr_ref="close"。 + 返回 dict: last_idx 最后一根在 chanlun 处理后 dataframe 里的下标 n_bars 实际参与计算的根数 - hits 命中列表,每项 {direction, h1_agree} + atr_pct 信号根 ATR / 次根开盘价 + hits 命中列表,每项含方向与三个滤网标志、pass_all error 出错时的说明,正常为 None """ + import numpy as np import pandas as pd try: @@ -40,45 +60,76 @@ def compute(df_l, df_h) -> dict: 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 + from lib.shadow_budget import ATR_GATE_BP chan_l = TF_DF(df_l, 1, "1m") cdf = chan_l.dataframe last = len(cdf) - 1 base = {"last_idx": last, "n_bars": int(len(df_l)), "hits": [], - "error": None} + "atr_pct": None, "atr_ref": None, "error": None} - zones = build_htf_zones(cdf, "1m", chan=chan_l) + # ATR 门控。分母与 exit_model.walk_exits 一致,取次根开盘价。 + # 放在任何早退之前——无信号的根也要记,才能在线看到门控的真实刷除率 + atr = float(cdf["atr"].to_numpy(dtype=float)[last]) \ + if "atr" in cdf.columns else float("nan") + ref = entry_px if (entry_px and np.isfinite(entry_px)) else \ + float(cdf["close"].to_numpy(dtype=float)[last]) + atr_pct = atr / ref if (np.isfinite(atr) and ref) else float("nan") + gate_ok = bool(np.isfinite(atr_pct) and atr_pct * 1e4 >= ATR_GATE_BP) + base["atr_pct"] = None if not np.isfinite(atr_pct) else float(atr_pct) + base["atr_ref"] = "next_open" if (entry_px and np.isfinite(entry_px)) \ + else "close" + + zones = build_htf_zones(cdf, "1m", chan=chan_l).reset_index(drop=True) if zones.empty: return base - sig = find_fast_bsp3(cdf, zones.reset_index(drop=True)) + + # 中枢阶梯:当前中枢是否整体脱离前一个。与 step42 同一算法 + z = zones.copy() + prev_zg, prev_zd = z["zg"].shift(), z["zd"].shift() + z["z_above"], z["z_below"] = z["zd"] > prev_zg, z["zg"] < prev_zd + z["zone_i"] = np.arange(len(z)) + + sig = find_fast_bsp3(cdf, zones) if sig is None or sig.empty: return base + sig = sig.merge(z[["zone_i", "z_above", "z_below"]], + on="zone_i", how="left") - # 只留落在最后一根上的信号,其余是历史,实盘下不了 - cur = sig[sig["entry_idx"].astype(int) == last] - if cur.empty: - return base - - # 5m 同向过滤:算得出就标 h1_agree,算不出就当未过滤照记 - agree_map: dict[int, int] = {} + # 5m 同向。算不出时 h1_agree 记 0,该信号自然不会通过 pass_all if df_h is not None and len(df_h) > 0: chan_h = TF_DF(df_h, 1, "5m") 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") - f_cur = full[full["entry_idx"].astype(int) == last] - for _, r in f_cur.iterrows(): - agree_map[int(r["direction"])] = int(r.get("h1_agree", 0)) + sig = attach_htf_context(sig, cdf, tl, "h1") + else: + sig["h1_agree"] = 0 - base["hits"] = [{"direction": int(r["direction"]), - "h1_agree": agree_map.get(int(r["direction"]), 0)} - for _, r in cur.iterrows()] + cur = sig[sig["entry_idx"].astype(int) == last] + if cur.empty: + return base + + hits = [] + for _, r in cur.iterrows(): + d = int(r["direction"]) + push = r["z_above"] if d == 1 else r["z_below"] + ladder_ok = bool(pd.notna(push) and bool(push)) + # attach_htf_context 在入场时刻之前没有大级别分型时写 NaN。 + # 不能写成 `int(x or 0)`——NaN 是真值,会走到 int(nan) 抛异常, + # 整根的信号就此丢掉,只留一行报错 + raw = r.get("h1_agree", 0) + agree = int(raw) if pd.notna(raw) else 0 + hits.append({"direction": d, "h1_agree": agree, + "ladder_ok": int(ladder_ok), "gate_ok": int(gate_ok), + "pass_all": int(agree == 1 and ladder_ok and gate_ok)}) + base["hits"] = hits return base except Exception as e: # 子进程里异常必须带回主进程,否则只见超时不见原因 import traceback return {"last_idx": -1, "n_bars": int(len(df_l)) if df_l is not None else 0, - "hits": [], "error": f"{type(e).__name__}: {e}", + "hits": [], "atr_pct": None, "atr_ref": None, + "error": f"{type(e).__name__}: {e}", "traceback": traceback.format_exc()} @@ -102,8 +153,8 @@ def _rebuild(rows) -> "object": def compute_packed(payload: tuple) -> dict: - """ProcessPoolExecutor 的入口:收 (l_rows, h_rows) 两组数值行。""" - l_rows, h_rows = payload + """ProcessPoolExecutor 的入口:收 (l_rows, h_rows, entry_px)。""" + l_rows, h_rows, entry_px = payload df_l = _rebuild(l_rows) df_h = _rebuild(h_rows) if h_rows else None - return compute(df_l, df_h) + return compute(df_l, df_h, entry_px) diff --git a/research/live/verify_signal_path.py b/research/live/verify_signal_path.py new file mode 100644 index 0000000..a1f8093 --- /dev/null +++ b/research/live/verify_signal_path.py @@ -0,0 +1,205 @@ +"""验证 live 信号路径与 step42 的批量过滤等价。 + +live 侧每根只看最后一根、且只喂 2000 根窗口;研究侧一次性跑全量。两者 +用同一套滤网(同向 + 中枢阶梯 + ATR 门控),但**不保证逐笔一致**—— +缠论结构依赖历史,2000 根窗口是 step39 定的命中率饱和点,不是无损截断。 + +每个信号根上比三种口径: + + 批量 全量历史 + 完整 5m。这是预算的来源,是基准 + 剔partial 窗口 2000 根 1m + 800 根**已收盘** 5m + 含partial 窗口,5m 末尾保留那根**尚未收盘**的。这是 live 现行做法 + +## 结论:partial 根要保留,不能剔 + +live 的 `df_h[df_h["timestamp"] < kline_ts]` 里 kline_ts 是 1m 的收盘时刻, +而正在走的那根 5m 开盘更早,于是被保留下来——五根里有四根如此。乍看像是 +「把未收盘的根当完整根用」的口径错误,实测**反过来**: + + BTC 25/25、ETH 24/25、SOL 23/25 与批量一致(合计 96%) + 剔掉则只有 21/25、21/25、18/25(合计 80%) + +原因是批量口径里那根 5m 是存在的。缠论的包含处理与分型检测吃整条序列, +凭空少一根会把结构整体挪位;保留一根「开盘价正确、高低点尚不完整」的 +近似根,比直接删掉更接近批量。 + +这也暴露了研究侧的一处残留:批量的 HTF 结构用到了那根 5m 的**最终**高低点, +而实盘在该时刻不可能知道。`htf_fx_timeline` 的 `confirm_ts += period` 只挡住了 +分型**选取**上的未来函数,挡不住结构构建。live 用 partial 根逼近,落在 96%, +差的那 4% 是这条残留的下界,不是可以修掉的 bug。 + + .venv/bin/python research/live/verify_signal_path.py --symbol BTC +""" +from __future__ import annotations + +import argparse +import sys +import warnings +from pathlib import Path + +warnings.filterwarnings("ignore") +HERE = Path(__file__).resolve().parent +sys.path.insert(0, str(HERE)) +sys.path.insert(0, str(HERE.parent)) +sys.path.insert(0, str(HERE.parents[1])) + +import numpy as np # noqa: E402 +import pandas as pd # noqa: E402 + +WINDOW_L, WINDOW_H = 2000, 800 + + +HTF_MS = 300_000 + + +def partial_htf_bar(df_l: pd.DataFrame, kline_ts: int) -> dict | None: + """用 1m 合成「此刻正在走的那根 5m」,复现 live 曾经喂进去的 partial 根。""" + bucket = kline_ts // HTF_MS * HTF_MS + if bucket >= kline_ts: # 正好落在 5m 边界,没有未收盘的根 + return None + part = df_l[(df_l["timestamp"] >= bucket) & (df_l["timestamp"] < kline_ts)] + if part.empty: + return None + date = pd.to_datetime(bucket, unit="ms", utc=True) \ + .tz_convert("Asia/Shanghai") + return {"timestamp": bucket, "date": date, + "open": float(part["open"].iloc[0]), + "high": float(part["high"].max()), "low": float(part["low"].min()), + "close": float(part["close"].iloc[-1]), + "volume": float(part["volume"].sum())} + + +def load(sym: str, tf: str, days: int) -> pd.DataFrame: + f = HERE / "cache" / f"bitget_{sym}_{tf}_{days}d.feather" + if not f.exists(): + raise SystemExit(f"缺数据 {f}") + return pd.read_feather(f) + + +def batch_flags(df_l: pd.DataFrame, df_h: pd.DataFrame) -> pd.DataFrame: + """step42_exit_tp_1m.run_one 的滤网,逐行照搬。""" + 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 + from lib.shadow_budget import ATR_GATE_BP + + chan_l = TF_DF(df_l, 1, "1m") + cdf = chan_l.dataframe + zones = build_htf_zones(cdf, "1m", chan=chan_l).reset_index(drop=True) + z = zones.copy() + pg, pdn = z["zg"].shift(), z["zd"].shift() + z["z_above"], z["z_below"] = z["zd"] > pg, z["zg"] < pdn + z["zone_i"] = np.arange(len(z)) + + chan_h = TF_DF(df_h, 1, "5m") + hdf = chan_h.dataframe + tl = htf_fx_timeline( + signals_to_frame(extract_fx_signals(chan_h, hdf)), hdf) + + sig = find_fast_bsp3(cdf, zones) + sig = sig.merge(z[["zone_i", "z_above", "z_below"]], on="zone_i", how="left") + sig = attach_htf_context(sig, cdf, tl, "h1") + + d = sig["direction"].astype(int) + push = np.where(d == 1, sig["z_above"], sig["z_below"]) + idx = sig["entry_idx"].astype(int).to_numpy() + entry = cdf["open"].to_numpy(float)[np.minimum(idx + 1, len(cdf) - 1)] + atr_pct = cdf["atr"].to_numpy(float)[idx] / entry + + out = pd.DataFrame({ + "entry_idx": idx, + "ts": cdf["timestamp"].to_numpy()[idx], + "direction": d.to_numpy(), + "h1_agree": sig["h1_agree"].fillna(0).astype(int).to_numpy(), + "ladder_ok": pd.Series(push).fillna(False).astype(int).to_numpy(), + "atr_bp": atr_pct * 1e4, + }) + out["gate_ok"] = (out["atr_bp"] >= ATR_GATE_BP).astype(int) + out["pass_all"] = ((out["h1_agree"] == 1) & (out["ladder_ok"] == 1) + & (out["gate_ok"] == 1)).astype(int) + return out, cdf + + +def main() -> None: + ap = argparse.ArgumentParser() + ap.add_argument("--symbol", default="BTC") + ap.add_argument("--days", type=int, default=30) + ap.add_argument("--checks", type=int, default=12) + a = ap.parse_args() + + df_l, df_h = load(a.symbol, "1m", a.days), load(a.symbol, "5m", a.days) + print(f"[{a.symbol}] 1m {len(df_l)} 根 / 5m {len(df_h)} 根,跑批量滤网…", + flush=True) + batch, cdf = batch_flags(df_l, df_h) + n = len(batch) + print(f" 原始 B4/S4 {n} 个 · 同向 {int((batch.h1_agree == 1).sum())}" + f" · 同向+阶梯 " + f"{int(((batch.h1_agree == 1) & (batch.ladder_ok == 1)).sum())}" + f" · 三项全过 {int(batch.pass_all.sum())}") + print(f" ATR 中位 {batch.atr_bp.median():.2f}bp · " + f"门控刷掉 {(1 - batch.gate_ok.mean()) * 100:.1f}%\n") + + # 挑最近的若干个信号根做窗口复现 + from shadow_signal import compute + cand = batch[batch["entry_idx"] >= WINDOW_L].tail(a.checks) + if cand.empty: + raise SystemExit("窗口内没有可核对的信号") + + def run_window(r, with_partial: bool): + i = int(r["entry_idx"]) + kline_ts = int(r["ts"]) + 60_000 # 信号根的收盘时刻 + wl = df_l[df_l["timestamp"] < kline_ts].tail(WINDOW_L) + wh = df_h[df_h["timestamp"] + HTF_MS <= kline_ts].tail(WINDOW_H) + if with_partial: + p = partial_htf_bar(df_l, kline_ts) + if p is not None: + wh = pd.concat([wh, pd.DataFrame([p])], ignore_index=True) + entry_px = float(cdf["open"].to_numpy(float)[min(i + 1, len(cdf) - 1)]) + res = compute(wl.reset_index(drop=True), wh.reset_index(drop=True), + entry_px) + if res.get("error"): + raise SystemExit(f"compute 报错,测试本身有问题:{res['error']}\n" + f"{res.get('traceback', '')}") + return next((h for h in res["hits"] + if h["direction"] == int(r["direction"])), None) + + def fmt(h): + if h is None: + return f"{'未复现':>18}" + return (f"{h['h1_agree']:>6}{h['ladder_ok']:>5}{h['gate_ok']:>5}") + + print(f"{'K线时刻':<15}{'方向':>4}{' 批量':>18}{' 窗口':>18}" + f"{' 含未收盘':>18}{' 截断':>7}{'partial':>9}") + print(f"{'':<15}{'':>4}{'同向 阶梯 门控':>20}{'同向 阶梯 门控':>20}" + f"{'同向 阶梯 门控':>20}") + n_trunc = n_part = n_ok = n_bad = 0 + for _, r in cand.iterrows(): + h_ok = run_window(r, False) + h_bad = run_window(r, True) + t = pd.to_datetime(r["ts"], unit="ms", utc=True) \ + .tz_convert("Asia/Shanghai").strftime("%m-%d %H:%M") + ref = (int(r.h1_agree), int(r.ladder_ok), int(r.gate_ok)) + got = None if h_ok is None else (h_ok["h1_agree"], h_ok["ladder_ok"], + h_ok["gate_ok"]) + bad = None if h_bad is None else (h_bad["h1_agree"], h_bad["ladder_ok"], + h_bad["gate_ok"]) + n_trunc += got != ref + n_part += bad != got + n_ok += got == ref + n_bad += bad == ref + print(f"{t:<15}{int(r.direction):>+4}" + f"{ref[0]:>6}{ref[1]:>5}{ref[2]:>5}" + f"{fmt(h_ok):>18}{fmt(h_bad):>18}" + f"{'' if got == ref else '差':>7}" + f"{'' if bad == got else '差':>9}") + n = len(cand) + print(f"\n 与批量(预算口径)一致:") + print(f" 剔掉未收盘 5m 根 {n_ok}/{n}") + print(f" 保留未收盘 5m 根 {n_bad}/{n} ← live 现行做法") + print(f" 两种窗口口径互不相同 {n_part}/{n}") + + +if __name__ == "__main__": + main()