From 54792fe015cab6d2428b924c96b15ef0d47d0f41 Mon Sep 17 00:00:00 2001 From: jack Date: Fri, 28 Aug 2026 03:02:06 +0800 Subject: [PATCH] =?UTF-8?q?=E6=8A=8A=20compute=5Fms=20=E6=8B=86=E6=88=90?= =?UTF-8?q?=E6=8E=92=E9=98=9F=E4=B8=8E=E7=BA=AF=E8=AE=A1=E7=AE=97=EF=BC=8C?= =?UTF-8?q?=E6=8D=AE=E6=AD=A4=E5=90=A6=E6=8E=89=E6=8D=A2=E6=9C=BA=E5=99=A8?= =?UTF-8?q?=E8=BF=99=E4=B8=AA=E6=96=B9=E5=90=91?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit compute_ms 一直是「提交进程池到拿到结果」的墙钟时间,排队和纯计算混在一个 数里,所以「加核有没有用」只能靠猜——这也是原先打算在 AWS 开第二台比 CPU 的依据。 worker 内部自己计时,连同父进程传入的提交时刻一起回传,拆出 queue_ms 与 inner_ms。实测中位 3ms / 695ms:2 个 worker 跑 3 个币并不排队,因为三个币的 收盘消息错峰到达。瓶颈全在单线程,加核压不到。 顺带把 README 里的内存数据从臆测的 1.5GB 改成实测 410MiB,并注明跨站点比 CPU 收益有限。 Co-authored-by: Cursor --- research/live/deploy/README.md | 31 +++++++++++++++++++++++++++---- research/live/shadow_hb.py | 26 +++++++++++++++++++++++--- research/live/shadow_signal.py | 26 +++++++++++++++++++++++--- 3 files changed, 73 insertions(+), 10 deletions(-) diff --git a/research/live/deploy/README.md b/research/live/deploy/README.md index 18c8889..aa3b034 100644 --- a/research/live/deploy/README.md +++ b/research/live/deploy/README.md @@ -117,10 +117,33 @@ python research/live/compare_sites.py \ 同理,`lag_data_ms` 两站也应当几乎相同(网络那段只有 2ms 空间)。真正该出现 差异的是 `compute_ms`。如果 `lag_data_ms` 差很多,先查时钟——比查网络更可能。 +⚠ 但先读下面「资源占用」一节:`compute_ms` 的差异几乎全部来自单核性能,而 +现役机型之间单核差距很小。**跨站点比 CPU 这件事本身收益有限**,本节流程保留 +是为了比网络与时钟,不建议为了比 CPU 单独开机器。 + ## 资源占用 -本机实测:内存约 1.5GB(两个计算进程 + 盘口缓冲),CPU 单核不满。 -盘口与成交流落盘约 15MB/天(gzip)。一周 168 小时的量级在百 MB 内。 +本机实测(2 vCPU EPYC 9K65 / 3 个币 / 2 worker):内存 **410MiB**,CPU 均值 +1~3%。盘口与成交流落盘约 15MB/天(gzip),一周在百 MB 内。 -`--workers 2` 是因为信号计算走独立进程池、不能阻塞事件循环。核数少的机型 -可以给 1,但要看心跳里的 `compute_ms`:若接近 60 秒就会开始堆积。 +内存和平均 CPU 都不是约束。约束是**单根 K 线的计算延迟**,而它是纯单线程的: + +``` +compute_ms 中位 764ms 父进程测的墙钟,含排队 + queue_ms 中位 3ms 等空闲 worker + inner_ms 中位 695ms 进程内真正在算 +``` + +`queue_ms` 只有 3ms,说明 **2 个 worker 跑 3 个币并不排队**——三个币的收盘消息 +错峰到达(SOL 最晚,排 66ms),没有真正的并发争抢。 + +**这条结论直接否掉了「换更强机器」这个方向。** 加核只能压 queue_ms,而它已经 +是 3ms;695ms 全在单线程里,取决于单核性能。t3a.medium(Zen 1,2017)单核比 +本机 Zen 5 慢 1.8~2 倍,换过去 compute_ms 会涨到 1200ms 以上。c7a / c6a 这类 +现代机型单核与本机相当,也换不到东西。 + +要压这 695ms 只有算法一条路:现在每分钟把 2001 根从头算一遍,其中 2000 根的 +结构与上一分钟完全相同。 + +`--workers 2` 是因为信号计算走独立进程池、不能阻塞事件循环。币数超过 worker +数才会看到 queue_ms 上来;届时加 worker 有效,加到与币数相等即可。 diff --git a/research/live/shadow_hb.py b/research/live/shadow_hb.py index 708ba86..a02c558 100644 --- a/research/live/shadow_hb.py +++ b/research/live/shadow_hb.py @@ -301,6 +301,9 @@ class Shadow: 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) # 成交监听:已挂上的币,以及必须持有的 forwarder 强引用 # (PubSub 只存弱引用,不持有的话监听会被 GC 静默摘掉) self._hooked: set[str] = set() @@ -326,7 +329,10 @@ class Shadow: "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", "n_bars", "n_hits", + "lag_data_ms", "lag_signal_ms", + # compute_ms 含排队;queue_ms/inner_ms 把它拆开,用来判断加核有没有用 + "compute_ms", "queue_ms", "inner_ms", + "n_bars", "n_hits", "n_pass", "atr_bp", "lag_med_ms", "lag_ok"]) # 无条件漂移:每根都记,用来和信号根上的条件漂移对照 self.f_drf, self.w_drf = _writer(d / "shadow_drift.csv", [ @@ -528,7 +534,7 @@ class Shadow: t0 = time.perf_counter() payload = (df_l[NUM_COLS].values.tolist(), - df_h[NUM_COLS].values.tolist(), baseline) + df_h[NUM_COLS].values.tolist(), baseline, time.time()) loop = asyncio.get_running_loop() from shadow_signal import compute_packed try: @@ -540,6 +546,10 @@ class Shadow: 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"]) hits = res.get("hits", []) atr_pct = res.get("atr_pct") @@ -552,7 +562,9 @@ class Shadow: "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, "n_bars": res.get("n_bars", 0), + "compute_ms": compute_ms, + "queue_ms": res.get("queue_ms"), "inner_ms": res.get("inner_ms"), + "n_bars": res.get("n_bars", 0), "n_hits": len(hits), "n_pass": n_pass, "atr_bp": atr_bp, "lag_med_ms": lag_med, "lag_ok": int(lag_ok)}) self.f_lat.flush() @@ -744,6 +756,14 @@ class Shadow: 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)) + verdict = ("排队为主 → 加 worker/加核直接见效" + if q > i else + "纯计算为主 → 加核帮不上,需改增量计算") + print(f" [计算] 排队中位 {q:.0f}ms · 纯计算中位 {i:.0f}ms" + f" · {verdict}(worker {self.workers} 个 / 币 {len(SYMS)} 个)", + flush=True) # 五分钟一根都没进来,说明管道断了。不喊一声就只能靠人翻日志 if self.n_bars == self._hb_last_bars: print(f" ⚠ [停滞] 距上次心跳未处理任何 K 线" diff --git a/research/live/shadow_signal.py b/research/live/shadow_signal.py index e1c72fa..2299ab6 100644 --- a/research/live/shadow_signal.py +++ b/research/live/shadow_signal.py @@ -31,6 +31,7 @@ Hummingbot 的 asyncio 循环里会把行情处理一起卡住,所以必须隔 from __future__ import annotations import os +import time import warnings warnings.filterwarnings("ignore") @@ -153,8 +154,27 @@ def _rebuild(rows) -> "object": def compute_packed(payload: tuple) -> dict: - """ProcessPoolExecutor 的入口:收 (l_rows, h_rows, entry_px)。""" - l_rows, h_rows, entry_px = payload + """ProcessPoolExecutor 的入口:收 (l_rows, h_rows, entry_px[, t_submit])。 + + 返回里带上 `queue_ms` 与 `inner_ms`,把父进程看到的墙钟时间拆开: + + compute_ms(父进程测)= queue_ms + 反序列化 + inner_ms + 回传 + + 这个拆分决定「加核有没有用」。排队占大头说明 worker 数不够(币数多于 + worker 数时,同一秒收盘的币只能排队),加核直接见效;纯计算占大头说明 + 单核性能受限,加核帮不上,得从算法上改成增量更新。 + 两者混在一个数里就只能靠猜。 + """ + t_start = time.time() + l_rows, h_rows, entry_px, *rest = payload + t_submit = rest[0] if rest else None + + t0 = time.perf_counter() df_l = _rebuild(l_rows) df_h = _rebuild(h_rows) if h_rows else None - return compute(df_l, df_h, entry_px) + out = compute(df_l, df_h, entry_px) + out["inner_ms"] = int((time.perf_counter() - t0) * 1000) + # 同一台机器,父子进程时钟一致,可直接相减 + out["queue_ms"] = int((t_start - t_submit) * 1000) \ + if t_submit is not None else None + return out