Files
jackandCursor 335d891478 生产与研究分家:实盘执行器独立成 live/ 子树,加 AWS 部署
实盘要跑在 AWS(API key 绑了 IP 白名单),而信号在新加坡那台算。借这次
把生产从研究侧摘出来,四条具体代价里第一条已经咬过:

1. live_state.json 原先落在 research/out/,而那里 shadow_hb 会自动 rename
   归档、研究脚本会写、人也手工清过。那文件装的是 MAX_DAY_LOSS 累计与已
   处理信号键,被清掉不报错,只是两道闸静默失效。改到 LIVE_HOME。
2. 采集器十币清空 300~560ms 直接叠在信号到达执行器的延迟上。
3. 研究侧探针 OOM 过一次(14.9GB),当时若有仓位在场会连坐执行器。
4. 为读两个常量 import 研究侧 step43,把 numpy/pandas/pyarrow 拖进实盘
   进程。抽出 stdlib-only 的 live/exit_params.py,install.sh 加断言挡回归。

新增 live/ship_signals.py:AWS 侧 ssh tail 拉总线,每次重连从文件头重放
+ 按幂等键去重,断线期间的信号自愈;旧信号由 staleness 闸挡掉不补做。
带时钟倒流检测——两机时钟不同步会让那道闸静默放宽。

部署件:systemd 两单元(搬运挂了执行器仍管在场仓位的超时平仓)、
install.sh、dryrun.sh(验密钥/白名单/时钟/ssh/取整)、status.sh、README。

验证:live_exec 重构后端到端空跑,SOL 多头与 ADA 空头的止损/两级止盈/
数量取整逐项核对正确,isolated + post_only + reduceOnly 都在;搬运的去重、
重启不重复追加、脏数据跳过、断线重连重放均已测。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 17:03:03 +08:00

938 lines
45 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""影子交易器:在 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 sys
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
# 总线模块住在生产子树 live/ 下。方向是刻意的:**生产不 import 研究侧**
# 研究侧反过来读生产持有的契约。见 live/live_exec.py 文件头
sys.path.insert(0, str(Path(__file__).resolve().parents[2] / "live"))
import signal_bus # noqa: E402
import tg_notify # noqa: E402
# 站点标识。跨地对比时两台机器的 CSV 要能合起来读,没有这一列就分不清哪行
# 来自哪台。默认取主机名,部署脚本会显式传 SHADOW_SITE(如 sg-hetzner
SITE = os.environ.get("SHADOW_SITE") or socket.gethostname()
# 判「加 worker 有没有用」必须知道核数:CPU 密集的活,worker 超过核数不增吞吐
CORES = os.cpu_count() or 1
# 判定要知道增量开没开,否则增量已生效时还会继续推荐「走增量」
INCR_ON = os.environ.get("SHADOW_INCR", "1") not in ("0", "", "false")
# 少于这么多根就不送去算。缠论要先有分型再有笔再有中枢,几十根出不来中枢,
# 送过去只会白占一个计算槽
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.<RUN_ID>.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",
"stream_bars"])
# 无条件漂移:每根都记,用来和信号根上的条件漂移对照
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(), sym)
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),
# 增量流当前窗口。恒等于 2001 说明缺口判定在每根都
# 回退重建,增量静默失效——只从耗时上看不出是哪一环
"stream_bars": res.get("stream_bars")})
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}")
# 手工执行的推送。只推过全部滤网的,且 lag 退化时不推——那与
# 「停开新仓」是同一条规则,不能只在自动化里执行
if h["pass_all"] and atr_pct:
if not lag_ok:
print(f" [TG] {sym} lag 退化,按停开新仓规则不推",
flush=True)
else:
# 总线先写、推送后发。写盘是同步的且已 fsync,实盘据此
# 下单;推送要走网络,不能让它的延迟挡在下单前面
signal_bus.emit(sym, kline_ts, h["direction"],
float(baseline), float(atr_pct),
t_data - kline_ts)
self._spawn(
tg_notify.push_signal(
sym, h["direction"], float(baseline),
float(atr_pct), kline_ts, t_data - kline_ts),
f"tg {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 不增吞吐。出路有两级:先上增量把真实计算
压下来;增量之后剩的是争抢放大(实测 3.3 倍,§5.72),那一级只能加核或
减币,继续改算法收益有限。
"""
if not np.isfinite(clear):
return "样本不足,暂不判定"
if clear < 400:
return f"清空 {clear:.0f}ms,宽裕,无需优化"
if self.workers < CORES and q > i:
return (f"排队为主且还有 {CORES - self.workers} 个空闲核 → "
f"--workers 加到 {CORES}")
if len(SYMS) <= CORES:
return f"清空 {clear:.0f}ms 偏高,但币数未超核数,先查别的争抢"
# 币数超核数:加 worker 不增吞吐,只能压单币耗时。但要看增量开没开,
# 否则会在增量已生效时继续推荐「走增量」——上线后实测踩到过
if not INCR_ON:
return (f"币数 {len(SYMS)} > 核数 {CORES},加 worker 无用(CPU 密集)"
f"。压单币耗时 → 开 SHADOW_INCR=1 走增量(实测 3.56x")
# 增量已生效时(§5.72 口径对齐后):单币 100ms ≈ 信号链 30ms + 追加
# 2.81 根 37ms + 争抢 33ms。争抢只有 1.49x,加核收益有限;而追加那 37ms
# 里约 22ms 纯属浪费——symbol 随机落 worker,各缓存都漏掉对方处理的根
return (f"币数 {len(SYMS)} > 核数 {CORES},增量已生效,加 worker 无用"
f"worker 已等于核数)。单币 {i:.0f}ms 里争抢只占约 1.5x"
f"最便宜的一刀是按币绑定 worker(每次只追 1 根,省约三分之一),"
f"其次是 add_indicators 增量化")
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()