fast_bsp3 改用 tol=-1 + require_touch=False,信号滞后从 5.8 根降到 2.2 根。 滞后与收益严格单调(年化 370% -> 906%,同一份数据同一套成本), 这是本轮提升的主因,也意味着实盘延迟会直接侵蚀收益。 新增 step31~39 验证策略能否落地: - 跨品种样本外——8 个未参与调参的币,PF 2.73 / t 28.5,无一为负 - 时点重建——只喂到信号那一根重算,同根命中 100%,确认无未来函数; 1m 在 2000 根窗口即饱和,计算耗时 0.20s - 偏差审计——多空对称、中枢生效时刻零回退、滑点稳健至 30bp、持仓几乎不重叠 - 消融——alpha 来自缠论中枢的上下文定位,而非「收盘转强」这个触发动作 补 research/HANDOFF.md:记录确切口径与参数、已排除的偏差、 已验证无效因而不必重做的方向,以及下一步用影子交易器实测执行滑点的方案。 清理 step1~20 的输出:早期方法论已被推翻(存在未来函数偏差), 其结论不再被引用;脚本保留,需要时可重跑。 Co-authored-by: Cursor <cursoragent@cursor.com>
247 lines
10 KiB
Python
247 lines
10 KiB
Python
"""Step 39:1m 的时点重建 + 窗口扫描——实盘该带多长历史。
|
||
|
||
step38 只验证了 15m/30m,且窗口固定 4000 根。对 15m 那是 41 天,
|
||
对 1m 只有 2.8 天,中枢的左边界效应完全不是一个量级。
|
||
|
||
1m 那条腿的毛收益 +0.0991%/笔、滑点预算只有 3.9bp,而窗口长度同时
|
||
牵动两头:
|
||
窗口太短 -> 中枢被截断,信号与回测口径不符(回测收益拿不到)
|
||
窗口太长 -> 每根算得慢,下单延迟大,滑点吃掉全部预算
|
||
|
||
本步对同一批信号在多个窗口下重建,同时计时,找命中率饱和且耗时可接受
|
||
的那个窗口,直接作为影子交易器的参数。
|
||
|
||
同根命中 全量口径的信号,在该窗口的时点重建下同一根也出现
|
||
假阳性 随机非信号点,时点重建却在当根给出信号(实盘会下、回测没有)
|
||
耗时 单次 TF_DF + 中枢 + 信号 的墙钟时间,即下单延迟的下限
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import argparse
|
||
import os
|
||
import sys
|
||
import time
|
||
import warnings
|
||
from concurrent.futures import ProcessPoolExecutor, as_completed
|
||
from pathlib import Path
|
||
|
||
import numpy as np
|
||
import pandas as pd
|
||
|
||
warnings.filterwarnings("ignore")
|
||
for v in ("OMP_NUM_THREADS", "OPENBLAS_NUM_THREADS", "MKL_NUM_THREADS"):
|
||
os.environ.setdefault(v, "1")
|
||
|
||
HERE = Path(__file__).resolve().parent
|
||
sys.path.insert(0, str(HERE))
|
||
sys.path.insert(0, str(HERE.parent))
|
||
pd.set_option("display.width", 320)
|
||
|
||
LTF, HTF = "1m", "5m" # step23 里 1m 那条腿的配置:5m 分型同向
|
||
WINDOWS = (2000, 4000, 8000, 16000)
|
||
HTF_RATIO = 5 # 5m 窗口 = 1m 窗口 / 5
|
||
HTF_MIN = 800
|
||
TOLERANCE = 3
|
||
MAX_ROWS = 1_200_000
|
||
|
||
|
||
def _pit_once(df_slice: pd.DataFrame, ltf: str) -> tuple[set[int], float]:
|
||
"""只用切片重建,返回信号下标集合与耗时(秒)。"""
|
||
from chanlun import TF_DF
|
||
from lib.fast_bsp3 import find_fast_bsp3
|
||
from lib.nested_level import build_htf_zones
|
||
|
||
t0 = time.perf_counter()
|
||
chan = TF_DF(df_slice, 1, ltf)
|
||
cdf = chan.dataframe
|
||
zones = build_htf_zones(cdf, ltf, chan=chan)
|
||
idx: set[int] = set()
|
||
if not zones.empty:
|
||
sig = find_fast_bsp3(cdf, zones.reset_index(drop=True))
|
||
if not sig.empty:
|
||
idx = set(sig["entry_idx"].astype(int).tolist())
|
||
return idx, time.perf_counter() - t0
|
||
|
||
|
||
def _pit_agree(df_htf_slice: pd.DataFrame, htf: str) -> pd.DataFrame:
|
||
"""切片重建 5m 分型时间线,用于检查同向过滤在当时是否成立。"""
|
||
from chanlun import TF_DF
|
||
from lib.fx_signal import extract_fx_signals, signals_to_frame
|
||
from lib.nested_bsp import htf_fx_timeline
|
||
|
||
chan = TF_DF(df_htf_slice, 1, htf)
|
||
s = signals_to_frame(extract_fx_signals(chan, chan.dataframe))
|
||
return htf_fx_timeline(s, chan.dataframe)
|
||
|
||
|
||
def run_one(task: tuple) -> dict | None:
|
||
import warnings as _w
|
||
_w.filterwarnings("ignore")
|
||
sys.path.insert(0, str(HERE))
|
||
sys.path.insert(0, str(HERE.parent))
|
||
|
||
from chanlun import TF_DF
|
||
from lib.data import fetch_ohlcv
|
||
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
|
||
|
||
sym, n_probe = task
|
||
pair = f"{sym}/USDT:USDT"
|
||
try:
|
||
df_l = fetch_ohlcv(pair, LTF, MAX_ROWS)
|
||
df_h = fetch_ohlcv(pair, HTF, 10 ** 9)
|
||
if df_l is None or len(df_l) < 50_000:
|
||
return None
|
||
|
||
chan_l = TF_DF(df_l, 1, LTF)
|
||
cdf = chan_l.dataframe
|
||
zones = build_htf_zones(cdf, LTF, chan=chan_l).reset_index(drop=True)
|
||
if zones.empty:
|
||
return None
|
||
|
||
chan_h = TF_DF(df_h, 1, HTF)
|
||
hdf = chan_h.dataframe
|
||
tl = htf_fx_timeline(signals_to_frame(extract_fx_signals(chan_h, hdf)), hdf)
|
||
|
||
# 全量口径:与 step23 的 1m 数字同一套筛选(仅 h1_agree)
|
||
full = attach_htf_context(find_fast_bsp3(cdf, zones), cdf, tl, "h1")
|
||
fin = full[full["h1_agree"] == 1]
|
||
idx_all = fin["entry_idx"].astype(int).to_numpy()
|
||
idx_all = idx_all[idx_all >= max(WINDOWS)]
|
||
if len(idx_all) < 20:
|
||
return None
|
||
|
||
rng = np.random.default_rng(11)
|
||
probe = (rng.choice(idx_all, size=min(n_probe, len(idx_all)), replace=False)
|
||
if len(idx_all) > n_probe else idx_all)
|
||
pool = np.setdiff1d(np.arange(max(WINDOWS), len(cdf) - 1), idx_all)
|
||
fp_probe = rng.choice(pool, size=min(len(probe), len(pool)), replace=False)
|
||
|
||
ltf_ts = cdf["timestamp"].to_numpy()
|
||
htf_ts = hdf["timestamp"].to_numpy()
|
||
|
||
rows, fps, times = [], [], []
|
||
for w in WINDOWS:
|
||
hw = max(w // HTF_RATIO, HTF_MIN)
|
||
for i in sorted(probe):
|
||
sl = cdf.iloc[i - w + 1: i + 1].reset_index(drop=True)
|
||
pit, dt = _pit_once(sl, LTF)
|
||
times.append({"窗口": w, "秒": dt})
|
||
last = len(sl) - 1
|
||
near = [p - last for p in pit if abs(p - last) <= TOLERANCE]
|
||
# 同向过滤也做时点重建:5m 只喂到不晚于该 1m 根的部分
|
||
h_end = int(np.searchsorted(htf_ts, ltf_ts[i], side="right"))
|
||
agree_ok = np.nan
|
||
if last in pit and h_end >= hw:
|
||
hsl = hdf.iloc[h_end - hw: h_end].reset_index(drop=True)
|
||
tl_p = _pit_agree(hsl, HTF)
|
||
one = pd.DataFrame({"entry_idx": [last],
|
||
"direction": [1]}) # 方向下面覆盖
|
||
d_full = int(fin.loc[fin["entry_idx"] == i, "direction"].iloc[0])
|
||
one["direction"] = d_full
|
||
got = attach_htf_context(one, sl, tl_p, "h1")
|
||
agree_ok = float(got["h1_agree"].iloc[0] == 1)
|
||
rows.append({"窗口": w, "idx": int(i), "exact": last in pit,
|
||
"near": bool(near),
|
||
"shift": min(near, key=abs) if near else np.nan,
|
||
"agree_ok": agree_ok})
|
||
for i in sorted(fp_probe):
|
||
sl = cdf.iloc[i - w + 1: i + 1].reset_index(drop=True)
|
||
pit, _ = _pit_once(sl, LTF)
|
||
fps.append({"窗口": w, "fp": (len(sl) - 1) in pit})
|
||
|
||
return {"task": f"{sym} {LTF}", "sym": sym, "n_full": len(idx_all),
|
||
"probe": pd.DataFrame(rows), "fp": pd.DataFrame(fps),
|
||
"times": pd.DataFrame(times)}
|
||
except Exception as e:
|
||
return {"task": f"{sym} {LTF}", "error": repr(e)[:250]}
|
||
|
||
|
||
def main() -> None:
|
||
ap = argparse.ArgumentParser()
|
||
ap.add_argument("--symbols", default="BTC,ETH,SOL")
|
||
ap.add_argument("--probe", type=int, default=60)
|
||
ap.add_argument("--workers", type=int, default=3)
|
||
args = ap.parse_args()
|
||
|
||
syms = [s.strip() for s in args.symbols.split(",")]
|
||
tasks = [(s, args.probe) for s in syms]
|
||
print(f"[1m 时点重建] {len(tasks)} 币 × 每币抽 {args.probe} 信号 × "
|
||
f"窗口 {WINDOWS}\n", flush=True)
|
||
|
||
res = []
|
||
with ProcessPoolExecutor(max_workers=args.workers) as ex:
|
||
futs = {ex.submit(run_one, t): t for t in tasks}
|
||
for i, f in enumerate(as_completed(futs), 1):
|
||
r = f.result()
|
||
if r is None or "error" in (r or {}):
|
||
print(f" [{i}] 跳过 {(r or {}).get('error', '')}", flush=True)
|
||
continue
|
||
res.append(r)
|
||
p = r["probe"]
|
||
best = p[p["窗口"] == max(WINDOWS)]
|
||
print(f" [{i}/{len(tasks)}] {r['task']} 全量信号 {r['n_full']},"
|
||
f"最大窗口同根命中 {best['exact'].mean() * 100:.0f}%", flush=True)
|
||
if not res:
|
||
print("无结果")
|
||
return
|
||
|
||
allp = pd.concat([r["probe"] for r in res], ignore_index=True)
|
||
allf = pd.concat([r["fp"] for r in res], ignore_index=True)
|
||
allt = pd.concat([r["times"] for r in res], ignore_index=True)
|
||
|
||
print("\n" + "=" * 116)
|
||
print("########## 1. 窗口 × 同根命中率(1m 结构在当时是否已成型)##########")
|
||
rows = []
|
||
for w, g in allp.groupby("窗口"):
|
||
f = allf[allf["窗口"] == w]
|
||
t = allt[allt["窗口"] == w]["秒"]
|
||
ag = g["agree_ok"].dropna()
|
||
rows.append({
|
||
"窗口(根)": w, "覆盖天数": f"{w / 1440:.1f}",
|
||
"抽检": len(g),
|
||
"同根命中": f"{g['exact'].mean() * 100:.1f}%",
|
||
f"±{TOLERANCE}根内": f"{g['near'].mean() * 100:.1f}%",
|
||
"完全消失": f"{(~g['near']).mean() * 100:.1f}%",
|
||
"同向过滤也成立": f"{ag.mean() * 100:.1f}%" if len(ag) else "—",
|
||
"假阳性": f"{f['fp'].mean() * 100:.1f}%",
|
||
"耗时中位": f"{t.median():.2f}s",
|
||
"耗时P95": f"{t.quantile(0.95):.2f}s",
|
||
})
|
||
tb = pd.DataFrame(rows)
|
||
print(tb.to_string(index=False))
|
||
print(" 同根命中率随窗口饱和的位置 = 实盘至少要带的历史长度。")
|
||
print(" 耗时是下单延迟的下限,1m 上每 100ms 都在吃那 3.9bp 预算。")
|
||
|
||
print("\n########## 2. 分币种(最大窗口口径)##########")
|
||
rows = []
|
||
for r in res:
|
||
p = r["probe"]
|
||
p = p[p["窗口"] == max(WINDOWS)]
|
||
rows.append({"品种": r["sym"], "全量信号": r["n_full"], "抽检": len(p),
|
||
"同根命中": f"{p['exact'].mean() * 100:.1f}%",
|
||
"完全消失": f"{(~p['near']).mean() * 100:.1f}%"})
|
||
print(pd.DataFrame(rows).to_string(index=False))
|
||
|
||
print("\n########## 3. 偏移分布(没命中同一根的,偏了几根)##########")
|
||
for w, g in allp.groupby("窗口"):
|
||
sh = g["shift"].dropna()
|
||
nz = sh[sh != 0]
|
||
print(f" 窗口 {w}: 有偏移 {len(nz)}/{len(g)} 笔"
|
||
+ (f",中位 {nz.median():+.0f} 根,范围 [{nz.min():+.0f}, {nz.max():+.0f}]"
|
||
if len(nz) else ""))
|
||
|
||
tb.to_csv(HERE / "out" / "step39_pit_1m.csv", index=False)
|
||
print("\n########## 结论 ##########")
|
||
top = allp[allp["窗口"] == max(WINDOWS)]
|
||
print(f" 最大窗口 {max(WINDOWS)} 根下:同根命中 {top['exact'].mean() * 100:.1f}%,"
|
||
f"完全消失 {(~top['near']).mean() * 100:.1f}%")
|
||
print(" 命中率若在某个窗口就饱和,用它跑影子交易器;")
|
||
print(" 若到最大窗口仍未饱和,说明 1m 回测口径本身就拿不到,比滑点更要紧。")
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|