服务端跑 probe_inner.py 与本机对表:分档占比一致(TF_DF 两腿 68~74%, 信号链 16~21ms),此前报的「chan 构建 22ms / 信号链 86ms」是无效减法 ——孤立环境的 1m append 减在场十币的 inner_ms,还漏了 5m 腿。 真实构成:inner_ms 100ms = 约 30ms 计算 × 3.3 倍争抢。逐个排除 _rebuild(1.6ms)、币种差异、周期重建(21.5%/1.5x)、批内次序(相关-0.083) 后,用外生的到达密集度确认(67→117ms),并以注入合成负载复现 (2 个满载进程 2.68x,3 个 2.86x,在场实测 2.24~3.34x)。 心跳判定不再推荐「继续压算法」,改为加核或减币。 第一版用 t_signal-inner_ms 反推并发区间得相关 0.533,是循环构造, 已废弃并在 §5.72 记下这个坑。 Co-authored-by: Cursor <cursoragent@cursor.com>
934 lines
45 KiB
Python
934 lines
45 KiB
Python
"""影子交易器:在 Hummingbot 运行时上测 1m 腿的真实入场滑点。
|
||
|
||
不下单。用 Hummingbot 的 Bitget 连接器取真实盘口,按信号方向和仓位吃单深度
|
||
算出「若此刻市价单进场会成交在哪」,再与回测假设的成交价相减。
|
||
|
||
为什么必须跑在 Hummingbot 上而不是自写脚本:要测的是**生产路径**的滑点。
|
||
决定成交价的是实际执行链路的延迟,换个运行时测出来的数就不作数了。
|
||
(连接器的换根解析 bug 见 patched_candles.py,已修,拿回约 1.06 秒。)
|
||
|
||
为什么不能用 paper trade 的成交:那是 Hummingbot 自己的撮合模型模拟的,
|
||
测出来是模型行为不是市场行为。
|
||
|
||
口径对齐 step42_exit_tp_1m.py:回测假设成交在**信号次根的开盘价**
|
||
(entry_delay=1),所以基准价就是换根后新一根的 open。滑点为正表示比回测差。
|
||
滤网(同向 + 中枢阶梯 + ATR 门控)在 shadow_signal.py 里,必须与预算同源。
|
||
|
||
### 统计口径三条硬要求
|
||
|
||
1. **主口径只用 pass_all 的信号根**。未过滤的照记但只作提前读数——
|
||
在我们根本不会下单的根上测滑点会把判据算宽
|
||
2. **条件漂移与无条件漂移分开报**。所以每根 K 线都记一份漂移
|
||
(shadow_drift.csv),不只信号根。两者的差就是「系统性追价」的大小
|
||
3. **出场腿按 maker/taker 分开**。止盈挂限价不吃滑点,把那 60% 混进
|
||
平均值会低估真实成本。出场腿属持仓管理,尚未实现
|
||
|
||
### lag 探针:超阈值要停开仓,不能只打日志
|
||
|
||
补丁只防得住「上游代码变了」,防不住 Bitget 再改一次消息格式。每根记
|
||
本地接收 − K 线收盘,近 30 根取中位数,超 800ms 即判该币不健康。
|
||
要停开仓是因为这种退化是**经济性且静默**的:不崩不报错,只让收益慢慢
|
||
变差,几周后才从统计里看得出来。影子期不下单,故落到 lag_ok 字段上。
|
||
|
||
### 盘口滚动缓冲把延迟变成自变量
|
||
|
||
每 100ms 存一份盘口。信号触发后,不只记「我们实际算完时」的滑点,而是回查
|
||
t_close+0.5s / 1s / 2s / 5s 各一个。这样即使本机算得慢,也能读出「若延迟为
|
||
X 秒,滑点是多少」,决策不被自身实现拖累。
|
||
|
||
### 滑点分解
|
||
|
||
延迟漂移 中间价相对次根开盘价的偏移——主项,且入场方向上系统性追价
|
||
盘口价差 最优价相对中间价
|
||
深度冲击 吃单加权价相对最优价
|
||
|
||
Bitget 永续实测价差仅约 0.01bp、100 档深度,故预期延迟漂移占绝大部分。
|
||
|
||
docker run -d --name shadow -w /home/hummingbot \\
|
||
-e PYTHONPATH=/home/hummingbot:/repo/research:/repo/research/live:/repo \\
|
||
-v $PWD:/repo:ro -v $PWD/research/out:/out \\
|
||
--entrypoint /opt/conda/envs/hummingbot/bin/python \\
|
||
hummingbot/hummingbot:latest /repo/research/live/shadow_hb.py --hours 24
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import asyncio
|
||
import csv
|
||
import gzip
|
||
import json
|
||
import math
|
||
import os
|
||
import socket
|
||
import time
|
||
from collections import deque
|
||
from concurrent.futures import ProcessPoolExecutor
|
||
from concurrent.futures.process import BrokenProcessPool
|
||
from pathlib import Path
|
||
|
||
import numpy as np
|
||
import pandas as pd
|
||
|
||
from lib.shadow_budget import LAG_ALARM_MS, LAG_WINDOW, lag_healthy
|
||
|
||
import signal_bus
|
||
import tg_notify
|
||
|
||
# 站点标识。跨地对比时两台机器的 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)")
|
||
# 增量已生效时,单币耗时的大头不再是算法而是争抢:离线同一窗口 30ms,
|
||
# 在场 100ms,注入合成负载能复现这个倍数(HANDOFF §5.72)。所以这里不能
|
||
# 再推荐「继续改算法」——那是按 3.3 倍虚高的数去优化 30ms 的真实计算
|
||
return (f"币数 {len(SYMS)} > 核数 {CORES},增量已生效,加 worker 无用"
|
||
f"(worker 已等于核数)。单币 {i:.0f}ms 里真实计算约 30ms,"
|
||
f"其余是 {CORES} 核上争抢的放大 —— 杠杆是加核或减币,"
|
||
f"不是继续压算法")
|
||
|
||
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()
|