Files
Chan/research/step62_trend_end.py
jackyu66gitandCursor 5695d8e983 research: 趋势末端识别(step62),ext_run 单调区分但 PF 仍不过 1
用户指出很多一二类实际在趋势中途被识别而非末期,若真在末期即使有延迟也该走出
行情。用 step60 的线段顶点当标签找实时可算的区分特征。

ext_run(极值越过中枢边界几个 ATR)单调有效:5m 上四分位的命中率是
12.0/22.1/37.0/43.8%,PF 0.12/0.11/0.22/0.46。短延伸那批就是趋势中途被识别的,
占一半且 PF 仅 0.11。最佳组合 ext_run≥P75 且 div≥中位:命中率 47.9%、PF 0.54,
相对基准 28.7%/0.22 精度接近翻倍。

两个反直觉结果:背驰越强反而越差(div Q1 命中 17.0%/PF 0.13,Q4 37.5%/0.28),
是对 MACD 面积判据的直接证伪;趋势级数无区分力(命中率 28.5/30.6/28.6/25.5%
基本持平)。

另修正一个我先前的猜测:以为引擎漏了缠论「趋势 vs 盘整」前提,实测该条件在
2504 笔上恒为 True——B1 要求 enter_bi.dir == leave_bi.dir == DOWN,中枢向下进
向下出本身就定义了它嵌在下跌趋势里,引擎已隐含强制,过滤器无从添加。
zs_count 也不可用,它是全局中枢序号而非趋势内序号。

结论:识别可优化且幅度不小,但不是瓶颈,瓶颈是入场时点。

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-28 22:00:52 +08:00

272 lines
11 KiB
Python
Raw Permalink 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.
"""怎么把「趋势末端的一类」从「趋势中途的一类」里挑出来。
§3.397 的关键数字:一类里命中线段顶点(真反转)的只有 30%,那 30% 即使带着
8~9 根滞后也有 PF 0.70~0.83**没命中的 70% 是 PF 0.08**。
所以亏损几乎全部来自被误识别在趋势中途的那批 —— 用户的判断。
于是问题变成:有没有**实时可算**的特征能把两批分开。
**首要候选来自缠论本身**:一类买点要求的是**趋势背驰**,而趋势的定义是
「至少两个同向连续的中枢」。引擎的 `find_all_bsp` 对**任意**中枢都发信号,
完全没查这个前提 —— 单个盘整中枢上的「背驰」只是盘整背驰,本就不该当一类用。
`fast_bsp.add_zone_ladder` 早就实现了这个判定(B4 上把 PF 2.72 提到 3.41),
一类这边却没接。
测的特征全部只用信号时刻及之前的数据:
ladder 本中枢相对前一中枢是否同向推进(下降趋势要求 zg < 前一个 zd)
zs_count 该中枢在本段里的序号,越大趋势越成熟
div 离开段 MACD 面积 / 进入段面积,越小背驰越强
ext_run 极值越过中枢边界多少个 ATR,越大越延伸
atr_z 极值处波动率(§3.398)
评判分两层:**能否提高命中线段顶点的概率**(检测器精度),
以及**能否提高实际收益**(可交易性)。前者好后者不好也没用。
"""
from __future__ import annotations
import argparse
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")
HERE = Path(__file__).resolve().parent
sys.path.insert(0, str(HERE))
sys.path.insert(0, str(HERE.parent))
OUT = HERE / "out" / "step62_trend_end.feather"
SL, SCALE_AT, RUNNER, RSTOP, MAXB = 2.0, 3.0, 8.0, 2.0, 48
BASE_WIN, TOL = 200, 2
def collect(sym: str, tf: str, rows: int) -> pd.DataFrame | None:
from chanlun import TF_DF
from chanlun.core.ChanEnum import Chan_BSP_TYPE, Chan_SEG_DIR
from lib.data import fetch_ohlcv
from lib.exit_model import cfg_name, walk_exits
try:
df = fetch_ohlcv(f"{sym}/USDT:USDT", tf, rows)
if df is None or len(df) < 5_000:
return None
chan = TF_DF(df, 1, tf, lean=False)
cdf = chan.dataframe
bz = chan.cal_bi_zs_list_pure(chan.bi_list)
if not bz:
return None
bsp = chan.find_all_bsp(chan.bi_list, bz) or []
dser = pd.to_datetime(cdf["date"])
if dser.dt.tz is not None:
dser = dser.dt.tz_localize(None)
didx = pd.DatetimeIndex(dser)
n = len(cdf)
def to_i(ts) -> int:
t = pd.Timestamp(ts)
return int(didx.searchsorted(t.tz_localize(None) if t.tz else t))
atr = cdf["atr"].to_numpy(float)
cl = cdf["close"].to_numpy(float)
base = (pd.Series(atr).rolling(BASE_WIN, min_periods=50)
.median().shift(1).to_numpy())
# 标准答案:线段终点(未来函数,只当标签用,不进入任何过滤器)
seg_bot, seg_top = [], []
for sg in getattr(chan, "seg_list", []) or []:
if sg.end_time is None:
continue
i = to_i(sg.end_time)
if 0 <= i < n:
(seg_bot if sg.dir == Chan_SEG_DIR.DOWN
else seg_top).append(i)
if not seg_bot or not seg_top:
return None
truth = {1: np.array(sorted(seg_bot)), -1: np.array(sorted(seg_top))}
def near(i: int, d: int) -> bool:
a = truth[d]
k = int(np.searchsorted(a, i))
return any(0 <= j < len(a) and abs(int(a[j]) - i) <= TOL
for j in (k - 1, k))
# 中枢阶梯:按可用顺序排好,才谈得上「相对前一个」
zs_seq = sorted(bz, key=lambda z: to_i(z.bi_list[0].start_time))
pos = {id(z): k for k, z in enumerate(zs_seq)}
want = {Chan_BSP_TYPE.B1: ("B1", 1), Chan_BSP_TYPE.S1: ("S1", -1)}
rec = []
for b in bsp:
tag = want.get(b.type)
if tag is None or b.sure_time is None or b.zs is None:
continue
name, d = tag
i_ext, i_sure = to_i(b.klc.end_time), to_i(b.sure_time)
if not (0 <= i_ext < n and 0 <= i_sure < n):
continue
a = atr[i_ext]
if not np.isfinite(a) or a <= 0 or not np.isfinite(base[i_ext]):
continue
zs = b.zs
k = pos.get(id(zs))
# 趋势成熟度:本中枢往前数,连续同向推进的中枢有几个。
# 单看「相对前一个是否同向」没有区分力 —— B1 要求
# enter_bi.dir == leave_bi.dir == DOWN,即中枢向下进、向下出,
# 这本身就定义了它嵌在下跌趋势里,连续纯中枢自然逐级下移,
# 实测该条件在 2504 笔上恒为 True。**引擎已隐含强制了「趋势」前提。**
# 有区分力的是「连了几级」,那才是趋势成熟度。
ladder_n = 0
if k is not None:
j = k
while j > 0:
cur, prv = zs_seq[j], zs_seq[j - 1]
ok = (float(cur.zg) < float(prv.zd) if d == 1
else float(cur.zd) > float(prv.zg))
if not ok:
break
ladder_n += 1
j -= 1
enter_bi = zs.bi_list[0].pre if zs.bi_list else None
ea = abs(float(enter_bi.macd_hist)) if enter_bi is not None else np.nan
la = abs(float(b.bi.macd_hist))
edge = float(zs.zd) if d == 1 else float(zs.zg)
ext = float(b.klc.low if d == 1 else b.klc.high)
rec.append({
"sym": sym, "tf": tf, "type": name, "dir": d,
"i_ext": i_ext, "i_sure": i_sure,
"lag_bars": i_sure - i_ext,
"hit": near(i_ext, d),
"ladder_n": int(ladder_n),
"div": la / ea if (ea and np.isfinite(ea) and ea > 0) else np.nan,
"ext_run": abs(ext - edge) / a,
"atr_z": a / base[i_ext],
})
if not rec:
return None
r = pd.DataFrame(rec)
r = r[(r.i_sure < n - 2) & np.isfinite(atr[r.i_sure.values])
& (atr[r.i_sure.values] > 0)].reset_index(drop=True)
if r.empty:
return None
cfg = cfg_name(SL, RUNNER, MAXB, RSTOP)
res = walk_exits(cdf, pd.DataFrame({
"entry_idx": r.i_sure.values, "direction": r.dir.values}),
[SL], [RUNNER], [MAXB], scale_at=SCALE_AT, runners=(RUNNER,),
runner_stops=(RSTOP,))
if len(res) != len(r):
return None
for c in ("g", "r", "c", "b"):
r[c] = res[f"{cfg}_{c}"].to_numpy()
r["atr_pct"] = atr[r.i_sure.values] / cl[r.i_sure.values]
return r
except Exception as e: # noqa: BLE001
print(f" {sym} {tf} 失败: {type(e).__name__}: {e}", flush=True)
return None
def perf(g: pd.DataFrame) -> dict:
from lib.exit_model import fee_of, taker_notional
net = g.g.values - fee_of(g.r.values, g.c.values)
gR = g.g.values / (SL * g.atr_pct.values)
tn = taker_notional(g.r.values, g.c.values)
w, o = net[net > 0].sum(), -net[net <= 0].sum()
return {
"笔数": len(g), "命中率": f"{g.hit.mean()*100:.1f}%",
"胜率": f"{(net > 0).mean()*100:.1f}%",
"毛R": round(gR.mean(), 3),
"PF": round(w / o, 2) if o > 0 else np.inf,
"余量bp": round(net.mean() / tn.mean() * 1e4, 2),
"t值": round(gR.mean() / (gR.std(ddof=1) / np.sqrt(len(g))), 2),
}
def report(d: pd.DataFrame) -> None:
for tf, x in d.groupby("tf"):
print("\n" + "#" * 96)
print(f"########## {tf} · {len(x)} 笔一类 ##########")
print("\n【一】单特征对「命中线段顶点」的区分力(命中率基准 "
f"{x.hit.mean()*100:.1f}%")
rows = []
for nm, col, qs in [("背驰div", "div", 4), ("延伸ext_run", "ext_run", 4),
("波动atr_z", "atr_z", 4), ("趋势级数ladder_n", "", 0)]:
if not col:
q = pd.cut(x["ladder_n"], [-1, 1, 2, 3, 999],
labels=["级数≤1", "=2", "=3", "≥4"])
for k, g in x.groupby(q, observed=True):
if len(g) >= 30:
rows.append({"分组": str(k), **perf(g)})
continue
y = x.dropna(subset=[col])
if len(y) < 100:
continue
q = pd.qcut(y[col], qs,
labels=[f"{nm}Q{i+1}" for i in range(qs)],
duplicates="drop")
for k, g in y.groupby(q, observed=True):
if len(g) >= 30:
rows.append({"分组": str(k), **perf(g)})
print(pd.DataFrame(rows).to_string(index=False))
print("\n【二】叠加过滤:延伸是唯一单调的特征,看叠加还能不能推上去")
rows = [{"过滤器": "无(现状)", **perf(x)}]
e75 = x["ext_run"].quantile(.75)
e50 = x["ext_run"].quantile(.50)
for nm, g in [
(f"ext_run≥P50({e50:.1f})", x[x["ext_run"] >= e50]),
(f"ext_run≥P75({e75:.1f})", x[x["ext_run"] >= e75]),
(f"ext_run≥P75 且 级数≥3",
x[(x["ext_run"] >= e75) & (x["ladder_n"] >= 3)]),
(f"ext_run≥P75 且 div≥中位",
x[(x["ext_run"] >= e75) & (x["div"] >= x["div"].median())]),
]:
if len(g) >= 30:
rows.append({"过滤器": nm, **perf(g)})
print(pd.DataFrame(rows).to_string(index=False))
print("\n判读:命中率若被显著抬高,说明特征确实在区分「趋势末端 vs 中途」。"
"\n但 PF 才是能不能做的判据 —— 命中率上去而 PF 不过 1,"
"说明滞后仍然吃掉了全部。")
def main() -> None:
ap = argparse.ArgumentParser()
ap.add_argument("--symbols", default="BTC,ETH,SOL,LINK,DOGE")
ap.add_argument("--tfs", default="5m,15m")
ap.add_argument("--rows", type=int, default=200_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():
report(pd.read_feather(OUT))
return
syms = [s.strip() for s in args.symbols.split(",")]
tfs = [t.strip() for t in args.tfs.split(",")]
parts = []
with ProcessPoolExecutor(max_workers=args.workers) as ex:
fut = {ex.submit(collect, s, t, args.rows): (s, t)
for s in syms for t in tfs}
for i, f in enumerate(as_completed(fut), 1):
r = f.result()
s, t = fut[f]
print(f" [{i}/{len(fut)}] {s} {t} "
f"{0 if r is None else len(r)}", flush=True)
if r is not None:
parts.append(r)
if not parts:
print("无结果")
return
d = pd.concat(parts, ignore_index=True)
OUT.parent.mkdir(exist_ok=True)
d.to_feather(OUT)
report(d)
if __name__ == "__main__":
main()