Files
Chan/research/live/shadow_signal.py
T
jackandCursor 54792fe015 把 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>
2026-08-28 03:02:06 +08:00

181 lines
7.8 KiB
Python
Raw 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.
"""影子交易器的信号函数——在子进程里跑,不碰事件循环。
单次调用约 0.26s 的纯 CPU,且 chanlun 是纯 Python 受 GIL 限制,放进
Hummingbot 的 asyncio 循环里会把行情处理一起卡住,所以必须隔离到独立进程。
## 口径必须与预算同源,缺一项数就不可比
预算(`lib/shadow_budget.BUDGET_BP`)算在 step42 的这套滤网上,
本文件逐行对齐 `step42_exit_tp_1m.run_one`
同向 h1_agree == 1
中枢阶梯 多头要求当前中枢整体高于前一个(zd > 前 zg),空头反之
ATR 门控 atr_pct ≥ ATR_GATE_BP(当前 8bp
早先这里只有 h1_agree。缺阶梯与门控测的就不是我们要交易的那批信号,
而这一项改常数解决不了——必须改信号路径本身。
门控阈值是**费率的函数**不是市场常数(低 ATR 信号的毛质量反而最好,
断崖只在扣费后出现),换 VIP 档或换交易所要重扫,不要抄 8bp。
## 与回测的两点差别
其一,这里只关心**最后一根已收盘 K 线**上有没有信号——实盘只能在当下下单。
其二,`atr_pct` 的分母取次根开盘价,与 `exit_model.walk_exits` 一致,
所以调用方必须把次根开盘价传进来。
未通过滤网的信号也一并返回并打上标志:过滤后样本很稀(门控后 8 个币
合计约 38 笔/周),未过滤的可作提前读数。但**统计主口径只能用 pass_all**
在我们根本不会下单的根上测滑点会把判据算宽。
"""
from __future__ import annotations
import os
import time
import warnings
warnings.filterwarnings("ignore")
for _v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS"):
os.environ.setdefault(_v, "1")
def compute(df_l, df_h, entry_px: float | None = None) -> dict:
"""在 df_l 的最后一根上找信号。df_l/df_h 都只含已收盘 K 线。
entry_px 是次根开盘价(回测 entry_delay=1 的成交价),用作 atr_pct 的
分母。取不到时退回用信号根收盘价,并在返回里标 atr_ref="close"。
返回 dict
last_idx 最后一根在 chanlun 处理后 dataframe 里的下标
n_bars 实际参与计算的根数
atr_pct 信号根 ATR / 次根开盘价
hits 命中列表,每项含方向与三个滤网标志、pass_all
error 出错时的说明,正常为 None
"""
import numpy as np
import pandas as pd
try:
from chanlun import TF_DF
from lib.fast_bsp3 import find_fast_bsp3
from lib.fx_signal import extract_fx_signals, signals_to_frame
from lib.nested_bsp import attach_htf_context, htf_fx_timeline
from lib.nested_level import build_htf_zones
from lib.shadow_budget import ATR_GATE_BP
chan_l = TF_DF(df_l, 1, "1m")
cdf = chan_l.dataframe
last = len(cdf) - 1
base = {"last_idx": last, "n_bars": int(len(df_l)), "hits": [],
"atr_pct": None, "atr_ref": None, "error": None}
# ATR 门控。分母与 exit_model.walk_exits 一致,取次根开盘价。
# 放在任何早退之前——无信号的根也要记,才能在线看到门控的真实刷除率
atr = float(cdf["atr"].to_numpy(dtype=float)[last]) \
if "atr" in cdf.columns else float("nan")
ref = entry_px if (entry_px and np.isfinite(entry_px)) else \
float(cdf["close"].to_numpy(dtype=float)[last])
atr_pct = atr / ref if (np.isfinite(atr) and ref) else float("nan")
gate_ok = bool(np.isfinite(atr_pct) and atr_pct * 1e4 >= ATR_GATE_BP)
base["atr_pct"] = None if not np.isfinite(atr_pct) else float(atr_pct)
base["atr_ref"] = "next_open" if (entry_px and np.isfinite(entry_px)) \
else "close"
zones = build_htf_zones(cdf, "1m", chan=chan_l).reset_index(drop=True)
if zones.empty:
return base
# 中枢阶梯:当前中枢是否整体脱离前一个。与 step42 同一算法
z = zones.copy()
prev_zg, prev_zd = z["zg"].shift(), z["zd"].shift()
z["z_above"], z["z_below"] = z["zd"] > prev_zg, z["zg"] < prev_zd
z["zone_i"] = np.arange(len(z))
sig = find_fast_bsp3(cdf, zones)
if sig is None or sig.empty:
return base
sig = sig.merge(z[["zone_i", "z_above", "z_below"]],
on="zone_i", how="left")
# 5m 同向。算不出时 h1_agree 记 0,该信号自然不会通过 pass_all
if df_h is not None and len(df_h) > 0:
chan_h = TF_DF(df_h, 1, "5m")
hdf = chan_h.dataframe
tl = htf_fx_timeline(
signals_to_frame(extract_fx_signals(chan_h, hdf)), hdf)
sig = attach_htf_context(sig, cdf, tl, "h1")
else:
sig["h1_agree"] = 0
cur = sig[sig["entry_idx"].astype(int) == last]
if cur.empty:
return base
hits = []
for _, r in cur.iterrows():
d = int(r["direction"])
push = r["z_above"] if d == 1 else r["z_below"]
ladder_ok = bool(pd.notna(push) and bool(push))
# attach_htf_context 在入场时刻之前没有大级别分型时写 NaN。
# 不能写成 `int(x or 0)`——NaN 是真值,会走到 int(nan) 抛异常,
# 整根的信号就此丢掉,只留一行报错
raw = r.get("h1_agree", 0)
agree = int(raw) if pd.notna(raw) else 0
hits.append({"direction": d, "h1_agree": agree,
"ladder_ok": int(ladder_ok), "gate_ok": int(gate_ok),
"pass_all": int(agree == 1 and ladder_ok and gate_ok)})
base["hits"] = hits
return base
except Exception as e: # 子进程里异常必须带回主进程,否则只见超时不见原因
import traceback
return {"last_idx": -1, "n_bars": int(len(df_l)) if df_l is not None else 0,
"hits": [], "atr_pct": None, "atr_ref": None,
"error": f"{type(e).__name__}: {e}",
"traceback": traceback.format_exc()}
NUM_COLS = ("timestamp", "open", "high", "low", "close", "volume")
def _rebuild(rows) -> "object":
"""只传数值列,date 在这里按 lib/data.py 的同一规则重建。
跨进程传 tz-aware 的 datetime 既慢又容易在字符串往返中丢时区,
而时区若与回测不一致,chanlun 的 K 线标签就会错位。
"""
import pandas as pd
df = pd.DataFrame(rows, columns=list(NUM_COLS))
df["timestamp"] = df["timestamp"].astype("int64")
date = pd.to_datetime(df["timestamp"], unit="ms", utc=True) \
.dt.tz_convert("Asia/Shanghai")
df.insert(1, "date", date)
return df
def compute_packed(payload: tuple) -> dict:
"""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
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