210 lines
8.8 KiB
Python
210 lines
8.8 KiB
Python
import numpy as np
|
|
import pandas as pd
|
|
import talib
|
|
|
|
|
|
# =====================================================================
|
|
# 1. 指标引擎:未来所有新发掘的指标,全写在这里
|
|
# =====================================================================
|
|
class FeatureEngine:
|
|
|
|
def __init__(self, daily_df, min5_df):
|
|
"""
|
|
daily_df: 日线数据 DataFrame
|
|
min5_df: 当天 5分钟 级别 K 线数据 DataFrame (需要包含 'close', 'volume')
|
|
"""
|
|
self.df = daily_df.copy()
|
|
self.min5_df = min5_df.copy()
|
|
self._extract_basic_arrays()
|
|
|
|
def _extract_basic_arrays(self):
|
|
self.close = self.df["close"].astype(float).values
|
|
self.high = self.df["high"].astype(float).values
|
|
self.low = self.df["low"].astype(float).values
|
|
self.amount = self.df["amount"].astype(float).values
|
|
|
|
def calculate_all_features(self):
|
|
features = {}
|
|
|
|
# =====================================================================
|
|
# 1. 🔄 【核心修改】:用 5分钟 K 线精确估算今日分时均价
|
|
# =====================================================================
|
|
if not self.min5_df.empty:
|
|
# 计算每 5 分钟的成交金额(收盘价 * 成交量)
|
|
# 注意:如果你的 5分钟接口直接自带 'amount'(成交额) 列,请直接用 self.min5_df['amount']
|
|
m5_close = self.min5_df["close"].astype(float)
|
|
m5_volume = self.min5_df["volume"].astype(float)
|
|
|
|
total_amount = (m5_close * m5_volume).sum()
|
|
total_volume = m5_volume.sum()
|
|
|
|
if total_volume > 0:
|
|
# 算出截止到当前(14:48)的 A 股全天分时均价
|
|
vwap_today = total_amount / total_volume
|
|
# 当前最新价格(最后一根 5分钟线的收盘价,最接近实时现价)
|
|
current_price = m5_close.iloc[-1]
|
|
|
|
# 计算偏离度:(最新价 - 分时均价) / 分时均价
|
|
features["min_bias"] = (current_price - vwap_today) / vwap_today
|
|
else:
|
|
features["min_bias"] = 0.0
|
|
else:
|
|
features["min_bias"] = 0.0
|
|
|
|
# =====================================================================
|
|
# 2. 日线常规指标计算(保持你的经典布林带逻辑不变)
|
|
# =====================================================================
|
|
up, mid, low = talib.BBANDS(
|
|
self.close, timeperiod=20, nbdevup=2, nbdevdn=2, matype=0
|
|
)
|
|
self.df["bb_up"], self.df["bb_mid"], self.df["bb_low"] = up, mid, low
|
|
|
|
self.df["percent_b"] = (self.df["close"] - low) / (up - low)
|
|
self.df["bandwidth"] = (up - low) / mid
|
|
self.df["amount_avg_20d"] = (
|
|
self.df["amount"].rolling(20).mean().shift(1)
|
|
)
|
|
self.df["price_max_60d"] = self.df["close"].rolling(60).max().shift(1)
|
|
self.df["low_min_10d"] = self.df["low"].rolling(10).min().shift(1)
|
|
|
|
latest = self.df.iloc[-1]
|
|
prev = self.df.iloc[-2]
|
|
|
|
# 如果 5分钟线数据拿到了,现价以 5分钟最新收盘价为准(更接近 14:48 真实盘面)
|
|
# 如果没拿到,退化使用日线昨日收盘(做测试用)
|
|
features["close"] = (
|
|
current_price if not self.min5_df.empty else latest["close"]
|
|
)
|
|
|
|
features["prev_close"] = prev["close"]
|
|
features["amount"] = latest["amount"]
|
|
features["amount_avg_20d"] = latest["amount_avg_20d"]
|
|
features["ma5"] = latest["ma5"] if "ma5" in latest else latest["close"]
|
|
features["ma10"] = latest["ma10"] if "ma10" in latest else latest["close"] # 👈 核心:补上这一行!
|
|
features["ma20"] = latest["bb_mid"]
|
|
features["ma60"] = latest["ma60"] if "ma60" in latest else latest["close"] # 👈 顺便把ma60也安全带上
|
|
features["percent_b"] = latest["percent_b"]
|
|
features["bandwidth"] = latest["bandwidth"]
|
|
features["is_price_60d_max"] = (
|
|
latest["close"] >= latest["price_max_60d"]
|
|
)
|
|
features["is_not_new_low_10d"] = latest["close"] > latest["low_min_10d"]
|
|
|
|
history_bw = self.df["bandwidth"].iloc[-250:]
|
|
features["bw_quantile"] = (history_bw < latest["bandwidth"]).mean()
|
|
|
|
return features
|
|
|
|
# =====================================================================
|
|
# 2. 状态规则书:这里只根据指标数据定义状态切换门槛
|
|
# =====================================================================
|
|
class StateRuleBook:
|
|
|
|
@staticmethod
|
|
def evaluate_next_state(current_state, f):
|
|
"""f 传入的是 FeatureEngine 计算出来的最新特征字典"""
|
|
|
|
# 计算主升浪基础多头条件
|
|
is_ma_bull = (f["ma5"] > f["ma10"] > f["ma20"]) and (
|
|
f["ma5"] > f["prev_ma5"]
|
|
)
|
|
|
|
# 核心决策流:利用解耦后的字典 f 进行条件拆解
|
|
# 【判断是否从主升浪跌破】
|
|
if (
|
|
current_state == "STATE_3_MAIN_WAVE"
|
|
and f["close"] < f["ma10"]
|
|
and f["prev_close"] > f["prev_ma10"]
|
|
):
|
|
return (
|
|
"STATE_4_WAVE_END",
|
|
"⚠️ 主升浪确认结束!清空做T仓,准备重新激活做T。",
|
|
)
|
|
|
|
# 【判断是否爆发主升浪】
|
|
if (
|
|
is_ma_bull
|
|
and f["is_price_60d_max"]
|
|
and (f["amount"] > f["amount_avg_20d"] * 1.8)
|
|
):
|
|
return (
|
|
"STATE_3_MAIN_WAVE",
|
|
"🚀 主升浪开启 / 放量向上突破!做T脚本自动拉闸休眠,锁仓死拿!",
|
|
)
|
|
|
|
# 【判断是否向下破位】
|
|
if f["percent_b"] < 0.0: # 跌破布林下轨
|
|
return "STATE_2_DOWN_BREAK", "❌ 向下破位!停止低吸做T,转为观望。"
|
|
|
|
# 【判断是否低量筑底】
|
|
is_low_volume = f["amount"] < (f["amount_avg_20d"] * 0.5)
|
|
is_ma_converge = abs(f["ma5"] - f["ma20"]) / f["ma20"] < 0.03
|
|
if f["is_not_new_low_10d"] and is_low_volume and is_ma_converge:
|
|
return (
|
|
"STATE_5_BOTTOMING",
|
|
"🌱 下跌结束,正在地量筑底。允许开始尝试轻仓做T。",
|
|
)
|
|
|
|
# 【判断是否处于变盘前夜】
|
|
if f["bw_quantile"] < 0.12:
|
|
return (
|
|
"STATE_1_OSCILLATION_SQUEEZE",
|
|
"🎚️ 变盘前夜:带宽极度压缩。保持做T,但防范单边突破。",
|
|
)
|
|
|
|
# 默认返回横盘震荡
|
|
return (
|
|
"STATE_1_OSCILLATION",
|
|
"☕ 正常横盘震荡期。激活做T脚本,正常执行高抛低吸。",
|
|
)
|
|
|
|
|
|
# =====================================================================
|
|
# 3. 策略状态机:调度核心
|
|
# =====================================================================
|
|
class QuantStateMachine:
|
|
|
|
def __init__(self, initial_state="STATE_1_OSCILLATION"):
|
|
self.current_state = initial_state
|
|
|
|
def run_daily_diagnostic(self, daily_df, min_df=None):
|
|
# 1. 扔给传感器计算指标
|
|
engine = FeatureEngine(daily_df, min_df)
|
|
features = engine.calculate_all_features()
|
|
|
|
# 2. 扔给规则书判定状态
|
|
next_state, comment = StateRuleBook.evaluate_next_state(
|
|
self.current_state, features
|
|
)
|
|
|
|
# 3. 更新并固化状态
|
|
self.current_state = next_state
|
|
|
|
# 4. 根据最终状态,分发当天的交易指令
|
|
self._execute_trading_action(features, comment)
|
|
|
|
def _execute_trading_action(self, f, comment):
|
|
print(f"\n[当前系统状态]: {self.current_state}")
|
|
print(f"[状态诊断提示]: {comment}")
|
|
|
|
# 具体的交易动作分发
|
|
if self.current_state in [
|
|
"STATE_1_OSCILLATION",
|
|
"STATE_1_OSCILLATION_SQUEEZE",
|
|
]:
|
|
# 只有在震荡期,才读取 %B 或分时执行做T
|
|
if f["percent_b"] >= 1.0 or f["min_bias"] > 0.035:
|
|
print(">>> 💰 【执行动作】:尾盘高抛,卖出 20% 网格仓。")
|
|
elif f["percent_b"] <= 0.0 or f["min_bias"] < -0.035:
|
|
print(">>> 🛒 【执行动作】:尾盘低吸,接回 20% 网格仓。")
|
|
else:
|
|
print(">>> ☕ 【执行动作】:未触及极端做T边界,长线持股观望。")
|
|
|
|
elif self.current_state == "STATE_3_MAIN_WAVE":
|
|
print(">>> 🔒 【执行动作】:主升浪锁仓护航中,禁止任何人乱动做T筹码。")
|
|
|
|
elif self.current_state == "STATE_2_DOWN_BREAK":
|
|
print(">>> 🛑 【执行动作】:市场向下破位,做T有被套风险,禁止低吸!")
|
|
|
|
elif self.current_state == "STATE_5_BOTTOMING":
|
|
print(">>> 🔬 【执行动作】:地量筑底阶段,允许小仓位底部分批低吸。") |