把 compute_ms 拆成排队与纯计算,据此否掉换机器这个方向
compute_ms 一直是「提交进程池到拿到结果」的墙钟时间,排队和纯计算混在一个 数里,所以「加核有没有用」只能靠猜——这也是原先打算在 AWS 开第二台比 CPU 的依据。 worker 内部自己计时,连同父进程传入的提交时刻一起回传,拆出 queue_ms 与 inner_ms。实测中位 3ms / 695ms:2 个 worker 跑 3 个币并不排队,因为三个币的 收盘消息错峰到达。瓶颈全在单线程,加核压不到。 顺带把 README 里的内存数据从臆测的 1.5GB 改成实测 410MiB,并注明跨站点比 CPU 收益有限。 Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -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 线"
|
||||
|
||||
Reference in New Issue
Block a user