import numpy as np import pandas as pd import talib from util import calculate_brar, calc_kdj_tdx # ===================================================================== # 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 df_tmp = calc_kdj_tdx(self.df) self.df["K"] = df_tmp['K'] self.df["D"] = df_tmp['D'] self.df["J"] = df_tmp['J'] cci = talib.CCI(self.df['high'], self.df['low'], self.df['close'], timeperiod=14) self.df["cci"] = cci df_tmp = calculate_brar(self.df) self.df["ar"] = df_tmp["ar"] self.df['br'] = df_tmp["br"] 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] prev_2 = self.df.iloc[-3] # 如果 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"] features["prev_ma5"] = prev["ma5"] if "ma5" in prev else prev["close"] features["K"] = latest["K"] features["D"] = latest["D"] features["J"] = latest["J"] features["prev_K"] = prev["K"] features["prev_D"] = prev["D"] features["prev_J"] = prev["J"] features["ar"] = latest["ar"] features["prev_ar"] = prev["ar"] features["br"] = latest["br"] features["prev_br"] = prev["br"] features["cci"] = latest["cci"] features["prev_cci"] = prev["cci"] features["prev_cci_2"] = prev_2["cci"] 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"]) or ((f["prev_ar"] >= 150 and f["ar"] < 150) and (f["prev_br"] >= 300 and f["br"] < 300))) ): if (f["close"] < f["ma10"] and f["prev_close"] > f["prev_ma10"]): return ( "STATE_4_WAVE_END", "⚠️ 主升浪确认结束!清空做T仓,准备重新激活做T。(股价跌破10日均线)", ) elif ((f["prev_ar"] >= 150 and f["ar"] < 150) and (f["prev_br"] >= 300 and f["br"] < 300)): return ( "STATE_4_WAVE_END", "⚠️ 主升浪确认结束!清空做T仓,准备重新激活做T。(ARBR跌破警戒线)", ) # 【判断是否爆发主升浪】 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(">>> 🔬 【执行动作】:地量筑底阶段,允许小仓位底部分批低吸。")