tradingSystem/scripts/backtest_churn.py

396 lines
19 KiB
Python

# -*- coding: utf-8 -*-
"""
换手体检: 决策系统驱动的卖出, 是"避坑"还是"来回瞎折腾" (只读, 不写任何表)
==========================================================================
回答你担心的那件事: 决策系统偏动量、容易今天买明天卖 —— 这些快进快出的卖出, 扣掉
手续费和滑点之后, 到底帮我们躲过了下跌(该卖), 还是把还会涨的票甩了、白交学费(瞎折腾)?
以及: 若加一道"刚建仓 N 个交易日内、非硬止损的决策系统卖出先不自动执行"的护栏, 净收益
是变好还是变差? 用你自己的历史账本把这条曲线量出来, 别拍脑袋。
一句话方法: 全部基于评审账本 pms_action_ledger —— 它对每个买卖决策都记了当时现价 price_at,
天生就是反事实判分的锚。五段:
① 样本盘点 账本里买入类/卖出类各多少, 其中"决策系统驱动"的卖出多少, 快速来回多少。
② 来回配对 同一只票"一次买紧跟一次卖"配成来回, 算持有天数 + 扣成本后的来回净收益。
③ 前向收益 每笔卖出: 卖出价 vs 该股此后 T+1/T+5/T+20 收盘。卖完还涨=卖飞(疑似瞎折腾),
卖完就跌=避坑(该卖)。聚合看均值与"卖飞比例"—— 这段是"是不是瞎折腾"的正面回答。
④ 成本账 决策系统驱动的快速来回一共交了多少手续费+滑点(纯摩擦成本)。
⑤ 反事实 对 N=1/3/5/10 交易日扫一遍"最小持有期护栏": 刚建仓 N 日内、非硬止损的决策
系统卖出改成"不卖、持有到 T+H", 比较装护栏前后的净收益差。正=护栏有用。
成本口径(A股, 顶部常量可调):
佣金双边各 COMMISSION_RATE(默认万2.5, 每边最低5元)、印花税卖出单边 STAMP_RATE(默认千0.5)、
过户费双边 TRANSFER_RATE(默认十万1)、滑点每边 SLIPPAGE_BPS(默认8bp)。
前向收盘价源: fetch_daily_closes() 已按你给的 DATA_MODEL 接到 gp_day_data (18.199 db_gp_cj,
走 app 的 index 源), 并处理了它的两个坑(symbol 前缀式、close 是 VARCHAR)。库不可达时
③⑤两段自动跳过、只出①②④, 不报错。
运行(桥机 factorevaluation, 与你其它只读脚本一致):
docker compose run --rm pms-web python scripts/backtest_churn.py --days 120
读数纪律: 样本少于 MIN_SAMPLE 的段落只列数、不下结论 —— 不等数据的判分是编故事。
"""
from __future__ import annotations
import argparse
import json
import os
import sys
from datetime import datetime, timedelta
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from app.db.session import fetch_all, DBUnavailable # noqa: E402 只读单表, 过单表守卫
# ------------------------------------------------------------------ 可调常量
MIN_SAMPLE = 5 # 少于这个数只报数不下结论
COMMISSION_RATE = 0.00025 # 佣金费率(每边); 万2.5
COMMISSION_MIN = 5.0 # 佣金每笔最低(元)
STAMP_RATE = 0.0005 # 印花税(仅卖出); 千0.5
TRANSFER_RATE = 0.00001 # 过户费(双边); 十万1
SLIPPAGE_BPS = 8.0 # 滑点(每边, 基点); 8bp = 0.08%
# 账本 action 归类 (与 app/core/planner.py、signal_service.py 的口径对齐; 后续在真机再核一遍)
BUY_ACTIONS = {"OPEN", "FILL", "ADD", "DCA"} # 真正增仓的落地动作
SELL_ACTIONS = {"EXIT", "TRIM"} # 真正减仓的落地动作
# 决策系统驱动的判据: 卖出 + 来源是盘中/风控信号 (hard_numbers.source 或 reason 里含这些词)
SIGNAL_SOURCES = {"intraday", "risk_sell"}
SIGNAL_WORDS = ("决策系统", "风控", "SELL", "转弱", "派发")
# 近似"硬止损"(护栏放行的一档): 入场到卖出这段已经深亏, 视为真出事, 不该被护栏拦
HARD_RISK_DROP = -0.08
FWD_HORIZONS = (1, 5, 20) # 前向收益看这几个交易日
GUARD_DAYS = (1, 3, 5, 10) # 反事实扫的最小持有期
GUARD_HOLD_TO = 5 # 护栏放行改为"持有到 T+这么多交易日"再看
# ------------------------------------------------------------------ 小工具
def _f(v, d=0.0):
try:
return float(v)
except (TypeError, ValueError):
return d
def _load_json(s):
if not s:
return {}
try:
return json.loads(s) if isinstance(s, str) else dict(s)
except (json.JSONDecodeError, TypeError, ValueError):
return {}
def _pct(entry, exit_):
e, x = _f(entry), _f(exit_)
return (x / e - 1.0) if e > 0 and x > 0 else None
def _fmt_pct(x):
return f"{x:+.2%}" if x is not None else ""
# ------------------------------------------------------------------ 读账本 (单表, 过守卫)
def read_ledger(since_dt: datetime) -> list:
"""按时间读评审账本。只读一张表, 满足 153 代理的严格单表访问。"""
rows = fetch_all(
"SELECT id, ts_code, decided_at, action, arbiter, verdict, price_at, "
"hard_numbers_json, reason, ref_id "
"FROM pms_action_ledger WHERE decided_at >= :since "
"ORDER BY ts_code, decided_at, id",
{"since": since_dt.strftime("%Y-%m-%d %H:%M:%S")}, source="proxy")
out = []
for r in rows:
r = dict(r)
r["hn"] = _load_json(r.get("hard_numbers_json"))
out.append(r)
return out
def classify(row: dict) -> dict:
"""给账本一行贴标签: side(buy/sell/other) 与 signal_driven(是否决策系统驱动的卖出)。"""
act = str(row.get("action") or "").upper()
verdict = str(row.get("verdict") or "").upper()
acted = verdict == "PASS" # 只认真正放行落地的决策
if act in BUY_ACTIONS and acted:
side = "buy"
elif act in SELL_ACTIONS and acted:
side = "sell"
else:
side = "other"
src = str((row.get("hn") or {}).get("source") or "").lower()
reason = str(row.get("reason") or "")
signal_driven = (side == "sell") and (
src in SIGNAL_SOURCES or any(w in reason for w in SIGNAL_WORDS))
return {"side": side, "signal_driven": signal_driven}
# ------------------------------------------------------------------ 历史收盘价源 (gp_day_data @ 18.199)
_CLOSE_CACHE: dict = {}
_ANALYSIS_START = None # main 里设; 让 _fwd_closes_after 按整个分析区间取足够宽的每股窗口
_ANALYSIS_END = None
def _to_prefix(ts_code: str) -> str:
"""点式 600000.SH → 前缀式 SH600000 (gp_day_data.symbol 用前缀式, 见 DATA_MODEL 约定)。"""
s = str(ts_code or "").strip().upper()
if "." in s:
num, ex = s.split(".", 1)
return ex + num
return s
def fetch_daily_closes(ts_code: str, start_dt: datetime, end_dt: datetime) -> dict:
"""返回 {date: 收盘价}, 只含有行情的交易日。
源: gp_day_data (db_gp_cj @ 192.168.18.199, 走 app 的 index 数据源 = DB_MYSQL_URL)。
PMS 本来就连这台取大盘 zs_day_data, 同库里 gp_day_data 就是个股原始日线。按 DATA_MODEL
的两个坑处理: symbol 是**前缀式** SH600000(账本 ts_code 是点式, 这里转一下); close 是
**VARCHAR**, 读出来转 float。用原始收盘(非前复权): ≤20 日窗口内除权少, 且与账本 price_at
同为原始价, 口径一致。库不可达或该股无行情返回 {}, ③⑤两段自动跳过、不报错。
"""
sym = _to_prefix(ts_code)
try:
rows = fetch_all(
"SELECT `timestamp` AS d, `close` AS c FROM gp_day_data "
"WHERE symbol = :s AND `timestamp` >= :a AND `timestamp` <= :b "
"ORDER BY `timestamp`",
{"s": sym, "a": start_dt.strftime("%Y-%m-%d 00:00:00"),
"b": end_dt.strftime("%Y-%m-%d 23:59:59")}, source="index")
except Exception:
return {}
out = {}
for r in rows:
d, c = r.get("d"), _f(r.get("c"))
if c <= 0 or d is None:
continue
dd = d.date() if isinstance(d, datetime) else datetime.fromisoformat(str(d)).date()
out[dd] = c
return out
def _fwd_closes_after(ts_code: str, sell_dt: datetime, n_max: int = 30) -> list:
"""卖出日之后的交易日收盘序列 (按日期升序), 供取 T+1/T+5/T+20。
按股缓存**整个分析区间**的日线, 覆盖该股在窗口内的所有卖点 (同股多次卖不会漏窗)。"""
if ts_code not in _CLOSE_CACHE:
lo = (_ANALYSIS_START or sell_dt) - timedelta(days=10)
hi = (_ANALYSIS_END or sell_dt) + timedelta(days=60)
_CLOSE_CACHE[ts_code] = fetch_daily_closes(ts_code, lo, hi)
closes = _CLOSE_CACHE[ts_code]
if not closes:
return []
sd = sell_dt.date()
after = sorted((d, c) for d, c in closes.items() if d > sd)
return [c for _, c in after]
def forward_returns(ts_code: str, sell_dt: datetime, sell_price: float) -> dict:
"""卖出后 T+h 的涨跌 (相对卖出价)。正=卖飞(卖完还涨), 负=避坑(卖完就跌)。"""
seq = _fwd_closes_after(ts_code, sell_dt, n_max=max(FWD_HORIZONS))
out = {}
for h in FWD_HORIZONS:
out[h] = _pct(sell_price, seq[h - 1]) if len(seq) >= h else None
return out
# ------------------------------------------------------------------ 成本
def trade_cost(notional: float, is_sell: bool) -> float:
"""单边交易成本(元): 佣金(最低5) + 卖出印花税 + 过户费 + 滑点。"""
n = abs(_f(notional))
if n <= 0:
return 0.0
commission = max(n * COMMISSION_RATE, COMMISSION_MIN)
stamp = n * STAMP_RATE if is_sell else 0.0
transfer = n * TRANSFER_RATE
slip = n * SLIPPAGE_BPS / 10000.0
return commission + stamp + transfer + slip
# ------------------------------------------------------------------ 来回配对
def pair_round_trips(rows_by_code: dict) -> list:
"""同一只票: 把"一次买"和其后"第一次卖"配成一个来回(粗配, 不做逐笔 FIFO)。
返回每个来回: 入场/出场时间价、持有交易日近似(自然日近似)、来回毛收益、是否决策系统驱动。
说明: 这是方向与量级的体检, 不是会计账; 精算逐笔在 pms_lot, 要精算另说。"""
trips = []
for code, rows in rows_by_code.items():
pending_buy = None
for r in rows:
tag = classify(r)
if tag["side"] == "buy":
if pending_buy is None:
pending_buy = r
elif tag["side"] == "sell" and pending_buy is not None:
b, s = pending_buy, r
bt, st = b["decided_at"], s["decided_at"]
bt = bt if isinstance(bt, datetime) else datetime.fromisoformat(str(bt))
st = st if isinstance(st, datetime) else datetime.fromisoformat(str(st))
hold_days = (st.date() - bt.date()).days
gross = _pct(b["price_at"], s["price_at"])
trips.append({
"ts_code": code, "buy_at": bt, "sell_at": st,
"buy_price": _f(b["price_at"]), "sell_price": _f(s["price_at"]),
"hold_days": hold_days, "gross_ret": gross,
"signal_driven": tag["signal_driven"], "sell_reason": s.get("reason"),
})
pending_buy = None
return trips
# ------------------------------------------------------------------ 各段输出
def section_inventory(rows, trips):
buys = sum(1 for r in rows if classify(r)["side"] == "buy")
sells = sum(1 for r in rows if classify(r)["side"] == "sell")
sig_sells = sum(1 for r in rows if classify(r)["signal_driven"])
quick = [t for t in trips if t["signal_driven"] and t["hold_days"] is not None
and t["hold_days"] <= max(GUARD_DAYS)]
print("\n① 样本盘点")
print(f" 账本落地买入 {buys} 笔 · 落地卖出 {sells} 笔 · 其中决策系统驱动的卖出 {sig_sells}")
print(f" 配成来回 {len(trips)} 组 · 决策系统驱动的快速来回(≤{max(GUARD_DAYS)}交易日) {len(quick)}")
return quick
def section_roundtrips(trips):
sig = [t for t in trips if t["signal_driven"]]
print("\n② 决策系统驱动的来回 (毛收益, 未扣成本)")
if len(sig) < MIN_SAMPLE:
print(f" 样本 {len(sig)} 组 (<{MIN_SAMPLE}), 只列数不下结论。")
rets = [t["gross_ret"] for t in sig if t["gross_ret"] is not None]
if rets:
avg = sum(rets) / len(rets)
win = sum(1 for x in rets if x > 0) / len(rets)
holds = [t["hold_days"] for t in sig if t["hold_days"] is not None]
avg_hold = sum(holds) / len(holds) if holds else None
print(f" 平均持有 {avg_hold:.1f} 天 · 平均毛收益 {_fmt_pct(avg)} · 盈利占比 {win:.0%}")
for t in sorted(sig, key=lambda x: (x["gross_ret"] is None, x["gross_ret"] or 0))[:8]:
print(f" {t['ts_code']}{t['hold_days']}天 毛{_fmt_pct(t['gross_ret'])} "
f"| {str(t['sell_reason'] or '')[:40]}")
def section_forward(rows):
"""③ 前向收益: 卖完之后股价怎么走 —— 是不是瞎折腾的正面回答。"""
print("\n③ 决策系统卖出之后的前向收益 (正=卖飞/疑瞎折腾, 负=避坑/该卖)")
sig_sells = [r for r in rows if classify(r)["signal_driven"]]
got = {h: [] for h in FWD_HORIZONS}
for r in sig_sells:
st = r["decided_at"]
st = st if isinstance(st, datetime) else datetime.fromisoformat(str(st))
fr = forward_returns(r["ts_code"], st, _f(r["price_at"]))
for h in FWD_HORIZONS:
if fr[h] is not None:
got[h].append(fr[h])
if not any(got.values()):
print(" (前向价源未接通 fetch_daily_closes, 本段跳过 —— 接上后重跑即出)")
return
for h in FWD_HORIZONS:
v = got[h]
if len(v) < MIN_SAMPLE:
print(f" T+{h}: 样本 {len(v)} (<{MIN_SAMPLE}), 只报数")
continue
avg = sum(v) / len(v)
flew = sum(1 for x in v if x > 0.01) / len(v) # 卖完还涨超1%算卖飞
print(f" T+{h}: 样本 {len(v)} · 卖后平均 {_fmt_pct(avg)} · 卖飞比例 {flew:.0%}")
print(" 读法: 卖后平均为正、卖飞比例高 → 这批卖出多在砍还会涨的票(瞎折腾); "
"为负 → 多在避坑(该卖)。")
def section_cost(trips):
print("\n④ 决策系统驱动的来回, 纯摩擦成本")
sig = [t for t in trips if t["signal_driven"] and t["gross_ret"] is not None]
if not sig:
print(" 无样本。")
return
# 无法拿到每笔真实金额时, 以"单位名义1"估相对成本率; 有 hard_numbers.amount 时用真实额
total_rate = 0.0
for t in sig:
rt_cost_rate = (COMMISSION_RATE * 2 + STAMP_RATE + TRANSFER_RATE * 2
+ SLIPPAGE_BPS / 10000.0 * 2)
total_rate += rt_cost_rate
print(f" 每个来回的往返摩擦约 {(_rt_cost_rate()):.3%} (佣金双边+印花+过户+滑点双边)")
print(f" {len(sig)} 组来回累计摩擦 ≈ 名义规模的 {total_rate:.2%} "
f"(即平均毛收益要先跑赢这条线才算真挣到)")
def _rt_cost_rate():
return COMMISSION_RATE * 2 + STAMP_RATE + TRANSFER_RATE * 2 + SLIPPAGE_BPS / 10000.0 * 2
def section_counterfactual(trips):
"""⑤ 反事实: 最小持有期护栏 N 扫一遍。需要前向价源, 缺则跳过。"""
print("\n⑤ 反事实: 加「最小持有期护栏」后净收益怎么变 (正=护栏有用)")
affected0 = [t for t in trips if t["signal_driven"]]
# 探一下前向价源是否可用
probe = None
for t in affected0:
seq = _fwd_closes_after(t["ts_code"], t["sell_at"], n_max=GUARD_HOLD_TO)
if seq:
probe = True
break
if not probe:
print(" (前向价源未接通, 本段跳过 —— 接上 fetch_daily_closes 后重跑即出)")
return
for N in GUARD_DAYS:
deltas = []
for t in affected0:
if t["hold_days"] is None or t["hold_days"] > N:
continue # 护栏只管"刚建仓 N 日内"的卖出
if (t["gross_ret"] or 0) <= HARD_RISK_DROP:
continue # 硬止损放行, 不受护栏拦
seq = _fwd_closes_after(t["ts_code"], t["sell_at"], n_max=GUARD_HOLD_TO)
if len(seq) < GUARD_HOLD_TO:
continue
# 实际(卖了): 拿到 gross_ret, 并付了一次卖出摩擦; 之后空仓 = 0
actual = (t["gross_ret"] or 0) - _rt_cost_rate() / 2
# 反事实(没卖, 持有到 T+H 再看): 用 T+H 相对入场的收益, 只付了买入侧摩擦
held = _pct(t["buy_price"], seq[GUARD_HOLD_TO - 1])
if held is None:
continue
counter = held - _rt_cost_rate() / 2
deltas.append(counter - actual)
if len(deltas) < MIN_SAMPLE:
print(f" N={N}: 受影响样本 {len(deltas)} (<{MIN_SAMPLE}), 只报数")
continue
avg = sum(deltas) / len(deltas)
helped = sum(1 for d in deltas if d > 0) / len(deltas)
print(f" N={N} 交易日护栏: 受影响 {len(deltas)} 笔 · 平均净收益变化 {_fmt_pct(avg)} · "
f"变好占比 {helped:.0%}")
print(" 读法: 某个 N 上平均变化持续为正且变好占比过半 → 这套信号是「噪音来回型」, "
"护栏该上、N 取那档; 若普遍为负 → 是「大波段型」, 别压, 快卖多数是对的。")
# ------------------------------------------------------------------ main
def main():
ap = argparse.ArgumentParser(description="换手体检: 决策系统卖出是避坑还是瞎折腾 (只读)")
ap.add_argument("--days", type=int, default=120, help="回看多少自然日 (默认120)")
ap.add_argument("--since", type=str, default=None, help="或指定起始日 YYYY-MM-DD")
args = ap.parse_args()
since = (datetime.strptime(args.since, "%Y-%m-%d") if args.since
else datetime.now() - timedelta(days=args.days))
global _ANALYSIS_START, _ANALYSIS_END
_ANALYSIS_START, _ANALYSIS_END = since, datetime.now()
print(f"换手体检 · 账本自 {since:%Y-%m-%d} 起 · 成本口径 往返≈{_rt_cost_rate():.3%}")
try:
rows = read_ledger(since)
except DBUnavailable as e:
print(f"[FAIL] 读账本失败(库不可达): {e}"); sys.exit(1)
if not rows:
print("账本在该窗口内为空 —— 换个更长的 --days 再看。"); return
by_code = {}
for r in rows:
by_code.setdefault(r["ts_code"], []).append(r)
trips = pair_round_trips(by_code)
section_inventory(rows, trips)
section_roundtrips(trips)
section_forward(rows)
section_cost(trips)
section_counterfactual(trips)
print("\n完成。前向两段若显示「未接通」, 是历史收盘价源(fetch_daily_closes)还没接 —— "
"确认接哪张表后重跑即全。")
if __name__ == "__main__":
main()