"""Step 52:持仓中放量就平仓 —— 把成交量当出场信号。 用户提出:既然 step50 证明「入场根放量 = 给走完这一冲的人接盘」,那持仓过程中 出现放量根,是不是也说明这一波被走完了,该直接平掉? 这和 step50 是同一个机制的延伸,但**方向未知**:放量既可能是衰竭(该走), 也可能是突破续势的起点(走了就砍在起涨点)。step50 只证明了「入场时撞上放量 不好」,推不出「持仓时撞上放量该跑」——入场是你在接别人的盘,持仓时那根量 可能正是把你送上去的那批资金。 出场腿是 taker(收盘市价),成本按 TIME 计。触发优先级:同一根内止损(盘中) > 目标(盘中)> 放量平仓(收盘)。止损优先是保守侧。 ⚠️ 本步自己写了逐根模拟器,不走 walk_exits。**基线必须与 walk_exits 逐笔 相等才继续**——否则后面所有对比都是在和一个错的基线比。 """ from __future__ import annotations import argparse import os import sys 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", 340) SL, SCALE_AT, RUNNER, MAXB = 2.0, 3.0, 8.0, 48 GATE_BP = 8.0 THRESHOLDS = (3.0, 5.0, 8.0) OUT = HERE / "out" / "step52_volume_exit.feather" IS_START = pd.Timestamp("2026-01-30", tz="Asia/Shanghai") TP, SL_, TIME = 0, 1, 2 def simulate(cdf, sig, vr): """逐根前推。返回每笔在基线与各放量出场变体下的 (毛收益, 原因, 是否分批)。 runner 止损位与初始止损同为 2 ATR(HANDOFF §3.5 定的 rstop=SL), 所以全程止损线不动,不需要分段处理。 """ high = cdf["high"].to_numpy(float) low = cdf["low"].to_numpy(float) open_ = cdf["open"].to_numpy(float) close = cdf["close"].to_numpy(float) atr = cdf["atr"].to_numpy(float) n = len(cdf) variants = ["base"] for t in THRESHOLDS: variants += [f"v{t:g}", f"v{t:g}w"] rows = [] for s, d in zip(sig["entry_idx"].astype(int), sig["direction"].astype(int)): e = s + 1 if e >= n - 1: continue a = atr[s] if not np.isfinite(a) or a <= 0: continue entry = open_[e] end = min(e + MAXB, n - 1) row = {"sig_idx": s, "atr_pct": a / entry} for var in variants: thr = None if var == "base" else float(var[1:].rstrip("w")) only_win = var.endswith("w") scaled = False g = r = None xb = end for j in range(e, end + 1): adv = (high[j] - entry) / a if d == 1 else (entry - low[j]) / a ret = (entry - low[j]) / a if d == 1 else (high[j] - entry) / a # ① 止损(盘中)。同根内优先于目标,保守侧 if ret >= SL: hit = -SL * a / entry g = 0.5 * (SCALE_AT * a / entry) + 0.5 * hit if scaled else hit r, xb = SL_, j break # ② 目标(盘中限价) if not scaled and adv >= SCALE_AT: # 同根内可能既到 3 ATR 又到 8 ATR,按先减仓后续跑处理 scaled = True if adv >= RUNNER: g = 0.5 * (SCALE_AT * a / entry) + 0.5 * (RUNNER * a / entry) r, xb = TP, j break elif scaled and adv >= RUNNER: g = 0.5 * (SCALE_AT * a / entry) + 0.5 * (RUNNER * a / entry) r, xb = TP, j break # ③ 放量平仓(收盘市价) if thr is not None and np.isfinite(vr[j]) and vr[j] >= thr: px = d * (close[j] - entry) / entry if not only_win or px > 0: g = 0.5 * (SCALE_AT * a / entry) + 0.5 * px if scaled else px r, xb = TIME, j break if g is None: # 超时:末根收盘市价 px = d * (close[end] - entry) / entry g = 0.5 * (SCALE_AT * a / entry) + 0.5 * px if scaled else px r, xb = TIME, end row[f"{var}_g"], row[f"{var}_r"] = g, r row[f"{var}_c"] = int(scaled) row[f"{var}_b"] = xb - e + 1 rows.append(row) return pd.DataFrame(rows) def collect(sym: str, rows: int): 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 chanlun.analysis.fast_bsp import ( add_zone_ladder, attach_htf_agree, attach_zone_ladder, build_htf_zones, find_fast_bsp3, htf_fx_timeline, ) from lib.data import fetch_ohlcv from lib.exit_model import cfg_name, walk_exits try: df = fetch_ohlcv(f"{sym}/USDT:USDT", "1m", rows) if df is None or len(df) < 50_000: return None chan = TF_DF(df, 1, "1m", lean=True) cdf = chan.dataframe zones = build_htf_zones(cdf, "1m", chan=chan) if zones.empty: return None zl = add_zone_ladder(zones.reset_index(drop=True)) sig = find_fast_bsp3(cdf, zl) if sig.empty: return None dh = fetch_ohlcv(f"{sym}/USDT:USDT", "5m", 10 ** 9) ch = TF_DF(dh, 1, "5m", lean=True) sig = attach_zone_ladder( attach_htf_agree(sig, cdf, htf_fx_timeline(ch, ch.dataframe)), zl) vr = (cdf["volume"] / cdf["volume"].rolling(60, min_periods=10).mean() ).to_numpy(float) res = simulate(cdf, sig, vr) # 对拍:基线必须与已验证的 walk_exits 逐笔相等 ref = walk_exits(cdf, sig, [SL], [RUNNER], [MAXB], scale_at=SCALE_AT, runners=(RUNNER,), runner_stops=(SL,)) cfg = cfg_name(SL, RUNNER, MAXB, SL) if len(ref) != len(res): print(f" {sym} 对拍失败:笔数 {len(ref)} vs {len(res)}", flush=True) return None dg = np.abs(ref[f"{cfg}_g"].to_numpy() - res["base_g"].to_numpy()) dr = (ref[f"{cfg}_r"].to_numpy() != res["base_r"].to_numpy()).sum() if dg.max() > 1e-12 or dr: print(f" {sym} 对拍失败:毛收益最大差 {dg.max():.3e},原因分歧 {dr} 笔", flush=True) return None idx = sig["entry_idx"].to_numpy().astype(int) keep = np.isin(idx, res["sig_idx"].to_numpy()) idx = idx[keep] atr = cdf["atr"].to_numpy(float)[idx] close = cdf["close"].to_numpy(float)[idx] out = pd.DataFrame({ "sym": sym, "date": cdf["date"].to_numpy()[idx], "atr_bp": atr / close * 1e4, "htf": sig["htf_agree"].to_numpy()[keep], "lad": sig["ladder_ok"].to_numpy()[keep], "vr60_entry": vr[idx], }) for c in res.columns: if c not in ("sig_idx",): out[c] = res[c].to_numpy() return out except Exception as e: print(f" {sym} 失败: {e!r}", flush=True) return None def stat(d: pd.DataFrame, var: str, label: str) -> dict: from lib.exit_model import fee_of, taker_notional g = d[f"{var}_g"].to_numpy() r, c = d[f"{var}_r"].to_numpy(), d[f"{var}_c"].to_numpy() tn = taker_notional(r, c) net = g - fee_of(r, c) risk = SL * d.atr_pct.to_numpy() R = net / risk w, o = net[net > 0].sum(), -net[net <= 0].sum() return {"方案": label, "笔数": len(d), "胜率": f"{(net > 0).mean()*100:.1f}%", "毛R": round((g / risk).mean(), 3), "净均R": round(R.mean(), 3), "R夏普": round(R.mean() / R.std(ddof=1), 3), "PF": round(w / o, 2) if o > 0 else np.inf, "余量bp": round(net.mean() / tn.mean() * 1e4, 2), "均持仓": round(d[f"{var}_b"].mean(), 1), "触发率": f"{(r == TIME).mean()*100:.0f}%"} def report(d: pd.DataFrame, label: str) -> None: rows = [stat(d, "base", "基线(不看量)")] for t in THRESHOLDS: rows.append(stat(d, f"v{t:g}", f"vr60≥{t:g} 就平")) rows.append(stat(d, f"v{t:g}w", f"vr60≥{t:g} 且浮盈才平")) print(f"\n--- {label}({len(d)} 笔)---") print(pd.DataFrame(rows).to_string(index=False)) def main() -> None: ap = argparse.ArgumentParser() ap.add_argument("--symbols", default="BTC,BNB,ETH,SOL,LINK,LTC,AVAX,XRP,DOGE,ADA") ap.add_argument("--rows", type=int, default=800_000) ap.add_argument("--workers", type=int, default=3) ap.add_argument("--reuse", action="store_true") args = ap.parse_args() if args.reuse and OUT.exists(): d = pd.read_feather(OUT) else: syms = [s.strip() for s in args.symbols.split(",")] print(f"[放量出场] {len(syms)} 币 × {args.rows} 根 1m\n", flush=True) parts = [] with ProcessPoolExecutor(max_workers=args.workers) as ex: fut = {ex.submit(collect, s, args.rows): s for s in syms} for i, f in enumerate(as_completed(fut), 1): r = f.result() print(f" [{i}/{len(syms)}] {fut[f]} {0 if r is None else len(r)}" f"{' ⚠对拍未过' if r is None else ' 对拍通过'}", flush=True) if r is not None: parts.append(r) if not parts: print("无结果") return d = pd.concat(parts, ignore_index=True) d.to_feather(OUT) d["date"] = pd.to_datetime(d["date"]) d = d[(d.htf == 1.0) & d.lad & (d.atr_bp >= GATE_BP)].copy() print(f"\n实盘口径 {len(d)} 笔 {d.date.min():%Y-%m-%d} ~ {d.date.max():%Y-%m-%d}") print("\n" + "=" * 118) print("########## 放量出场 vs 基线 ##########") report(d[d.date < IS_START], "样本外") report(d[d.date >= IS_START], "发现期") report(d, "全样本") if __name__ == "__main__": main()