From f5f456c523c830d6870f502788ec9a7e65af25d5 Mon Sep 17 00:00:00 2001 From: jack Date: Fri, 28 Aug 2026 05:19:14 +0800 Subject: [PATCH] =?UTF-8?q?=E8=AE=A1=E7=AE=97=E5=88=A4=E5=AE=9A=E6=94=B9?= =?UTF-8?q?=E7=9C=8B=E3=80=8C=E6=B8=85=E7=A9=BA=E5=85=A8=E9=83=A8=E5=B8=81?= =?UTF-8?q?=E3=80=8D=EF=BC=8C=E5=B9=B6=E6=8C=A1=E6=8E=89=E7=A9=BA=E7=AA=97?= =?UTF-8?q?=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 判定原先比逐根的 queue_ms 与 inner_ms,结构上错了两处,十币下直接指反: 判据错。要紧的是一个收盘时刻清空所有币要多久,不是单币的 q 或 i。币同一秒 收盘,币数超 worker 数时后面的币串行等待,这笔代价不出现在任何单根的 q 或 i 里。现在按 kline_ts 聚合取各币最大 compute_ms,落到 clear_hist。 出路错。「排队为主 → 加核」只在还有空闲核时成立。worker 已等于核数时加 worker 不增吞吐,只把等待从 queue 挪到 inner。十币实测正是如此:inner 被 争抢从 144 抬到 192ms 反超 queue 135ms,于是判定落到「量级已低、无需优化」 ——而此时最后一个币已在 1376ms。现在币数超核数就直接指向增量路径。 顺带修一个瞬时故障:WS 重连瞬间 feed 的 deque 可能为空,空窗口放行会让 worker 抛「DataFrame for 1m is empty」,白占一个计算槽(币数超核数时会推迟 后面所有币),而报错文本还会让人以为是缺历史数据。加 MIN_BARS 守卫。 Co-authored-by: Cursor --- research/live/deploy/start.sh | 9 ++++- research/live/shadow_hb.py | 70 ++++++++++++++++++++++++++++++----- 2 files changed, 67 insertions(+), 12 deletions(-) diff --git a/research/live/deploy/start.sh b/research/live/deploy/start.sh index 1f193c6..7a9f351 100755 --- a/research/live/deploy/start.sh +++ b/research/live/deploy/start.sh @@ -18,8 +18,13 @@ NAME="${NAME:-shadow}" # 设 0 可退回 full,用来复量两模式的耗时差。 SHADOW_LEAN="${SHADOW_LEAN:-1}" # 币池。默认三个流动性最好的做滑点测量;十币池是实际要交易的那批(TRX 剔除, -# ATR 门控几乎全刷掉)。币数直接决定排队:所有币同一秒收盘,2 核上 10 个币 -# 需要约 640ms 墙钟才算完,最后一个币的信号会落在 800ms 哨兵线之外。 +# ATR 门控几乎全刷掉)。 +# +# 币数超过核数时排队会成为主项:所有币同一秒收盘,2 核上十币实测清空要 +# 约 640ms,最后一个币的信号落在 1376ms。此时**加 worker 没用**——CPU 密集 +# 的活,worker 超过核数不增吞吐,只会把等待从 queue_ms 挪到 inner_ms。 +# 唯一出路是压单币耗时,走 init_stream/append_bar 增量路径(实测约 3.7x, +# 换算后十币 / 2 核清空约 265ms)。 SYMS="${SYMS:-BTC,ETH,SOL}" IMAGE="${SHADOW_IMAGE:-hummingbot/hummingbot:latest}" HOURS="${HOURS:-168}" diff --git a/research/live/shadow_hb.py b/research/live/shadow_hb.py index f6f771e..b7fae66 100644 --- a/research/live/shadow_hb.py +++ b/research/live/shadow_hb.py @@ -74,6 +74,11 @@ from lib.shadow_budget import LAG_ALARM_MS, LAG_WINDOW, lag_healthy # 站点标识。跨地对比时两台机器的 CSV 要能合起来读,没有这一列就分不清哪行 # 来自哪台。默认取主机名,部署脚本会显式传 SHADOW_SITE(如 sg-hetzner) SITE = os.environ.get("SHADOW_SITE") or socket.gethostname() +# 判「加 worker 有没有用」必须知道核数:CPU 密集的活,worker 超过核数不增吞吐 +CORES = os.cpu_count() or 1 +# 少于这么多根就不送去算。缠论要先有分型再有笔再有中枢,几十根出不来中枢, +# 送过去只会白占一个计算槽 +MIN_BARS = 200 SYMS = ("BTC", "ETH", "SOL") # 多存一根:deque 尾部是尚未收盘的当前根,剔除后正好剩 step39 定下的窗口 @@ -324,6 +329,11 @@ class Shadow: # 排队 / 纯计算的滚动窗口,用来判断加核有没有用 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() @@ -548,6 +558,13 @@ class Shadow: # 末行是刚开始的那根,未收盘,必须剔除,否则等于用未来数据 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) @@ -570,6 +587,13 @@ class Shadow: 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") @@ -779,16 +803,12 @@ class Shadow: if self.q_hist and self.i_hist: q, i = float(np.median(self.q_hist)), float(np.median(self.i_hist)) # 建议要看绝对量级:lean + 新引擎后纯计算约 128ms,此时再提 - # 「改增量计算」是误导——尾部已由数据腿主导,压计算换不到东西 - if q > i: - verdict = "排队为主 → 加 worker/加核直接见效" - elif i > 400: - verdict = "纯计算为主且偏高 → 加核帮不上,需改增量计算" - else: - verdict = "纯计算为主但量级已低 → 无需再优化" - print(f" [计算] 排队中位 {q:.0f}ms · 纯计算中位 {i:.0f}ms" - f" · {verdict}(worker {self.workers} 个 / 币 {len(SYMS)} 个)", - flush=True) + 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 线" @@ -796,6 +816,36 @@ class Shadow: flush=True) self._hb_last_bars = self.n_bars + def _compute_verdict(self, q: float, i: float, clear: float) -> str: + """给出唯一可行的出路,而不是「哪一项数字更大」。 + + 旧版比逐根的 q 与 i,结构上错了两处: + + 1. 判据错。真正要紧的是**一个收盘时刻清空所有币要多久**(clear), + 不是单币的 q 或 i。所有币同一秒收盘,币数超过 worker 数时后面的 + 币必然串行等待,而这笔代价不出现在任何单根的 q 或 i 里。 + 2. 出路错。「排队为主 → 加核」只在还有空闲核时成立。worker 已等于 + 核数时,加 worker 不会增加吞吐——CPU 密集的活变不出来,只会把 + 等待从 queue_ms 挪到 inner_ms。十币实测正是如此:inner 被争抢从 + 144ms 抬到 192ms,反而超过 queue 135ms,于是判定落到「量级已低、 + 无需优化」,而此时最后一个币已经落在 1376ms。 + + 所以币数超过核数时,唯一的出路是压单币耗时(增量计算),不是加 worker。 + """ + if not np.isfinite(clear): + return "样本不足,暂不判定" + if clear < 300: + return f"清空 {clear:.0f}ms,宽裕,无需优化" + if self.workers < CORES and q > i: + return (f"排队为主且还有 {CORES - self.workers} 个空闲核 → " + f"--workers 加到 {CORES}") + if len(SYMS) > CORES: + per = clear / max(len(SYMS) / max(self.workers, 1), 1) + return (f"币数 {len(SYMS)} > 核数 {CORES},加 worker 无用(CPU 密集)" + f"。唯一出路是压单币耗时 {per:.0f}ms → 走增量 " + f"init_stream/append_bar(HANDOFF §5.5,实测约 3.7x)") + return f"清空 {clear:.0f}ms 偏高,但币数未超核数,先查是否有别的争抢" + async def run(self) -> None: await self.start() tasks = [asyncio.create_task(self.sample_books()),