""" 特征工程编排器 — 串联全部特征组,输出完整特征 DataFrame """ import pandas as pd from datetime import date from core.scoring.features.validator import load_candidates, DataContext from core.scoring.features.v3_2_features import calculate_features_v3_2 from core.scoring.features.v3_3_features import calculate_features_v3_3 from core.scoring.features.v3_4_features import calculate_features_v3_4 from core.scoring.features.negative_features import calculate_negative_features from core.scoring.features.independence_features import calculate_independence_features from core.scoring.features.sector_features import calculate_sector_independence_features from core.scoring.features.emotion_features import calculate_relaxed_emotion_features from core.logger import LogLevel, PrintLog class FeaturePipeline: """ 特征工程管道 — 串联 8 组特征计算,输出完整特征矩阵。 Usage: pipeline = FeaturePipeline(trade_date=date.today()) feature_df = pipeline.run() # DataFrame indexed by stock_code """ def __init__(self, trade_date: date): self.trade_date = trade_date self.ctx: DataContext = None def run(self) -> pd.DataFrame: """ 执行完整特征工程管道。 返回: DataFrame indexed by stock_code, columns = 全部特征 """ # Stage 0: 加载候选股和数据 PrintLog(LogLevel.INFO, '[pipeline] Stage 0: 加载候选股...') self.ctx = load_candidates(self.trade_date) if not self.ctx.candidates: PrintLog(LogLevel.WARNING, '[pipeline] 无候选股通过过滤') return pd.DataFrame() PrintLog(LogLevel.INFO, f'[pipeline] 候选股: {len(self.ctx.candidates)}, ' f'排除: {len(self.ctx.excluded)}') # Stage 1: v3.2 基础特征 (20维) PrintLog(LogLevel.INFO, '[pipeline] Stage 1: v3.2 基础特征 (20维)...') df = calculate_features_v3_2(self.ctx) PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # Stage 2: v3.3 扩展特征 (16维) PrintLog(LogLevel.INFO, '[pipeline] Stage 2: v3.3 扩展特征 (16维)...') df_v33 = calculate_features_v3_3(self.ctx) df = df.join(df_v33, how='inner', rsuffix='_v33') PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # Stage 3: v3.4 扩展特征 (16维) PrintLog(LogLevel.INFO, '[pipeline] Stage 3: v3.4 扩展特征 (16维)...') df_v34 = calculate_features_v3_4(self.ctx) df = df.join(df_v34, how='inner', rsuffix='_v34') PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # Stage 4: 负向指标 (5维) — 注意 #53 覆盖 #39 PrintLog(LogLevel.INFO, '[pipeline] Stage 4: 负向指标 (5维)...') df_neg = calculate_negative_features(self.ctx) # 使用 update 模式: 负向指标的 trend_consistency_20d 覆盖 v3.4 版本 common_cols = set(df.columns) & set(df_neg.columns) for col in common_cols: df[col] = df_neg[col] # 覆盖 new_cols = set(df_neg.columns) - common_cols for col in new_cols: df[col] = df_neg[col] PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # Stage 5: 大盘独立性 (3维) PrintLog(LogLevel.INFO, '[pipeline] Stage 5: 大盘独立性 (3维)...') df_ind = calculate_independence_features(self.ctx) df = df.join(df_ind, how='left') df[df_ind.columns] = df[df_ind.columns].fillna(0) PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # Stage 6: 行业独立性 (3维) PrintLog(LogLevel.INFO, '[pipeline] Stage 6: 行业独立性 (3维)...') df_sec = calculate_sector_independence_features(self.ctx) df = df.join(df_sec, how='left') df[df_sec.columns] = df[df_sec.columns].fillna(0) PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # Stage 7: 情绪弹性 (4维) PrintLog(LogLevel.INFO, '[pipeline] Stage 7: 情绪弹性 (4维)...') df_emo = calculate_relaxed_emotion_features(self.ctx) df = df.join(df_emo, how='left') df[df_emo.columns] = df[df_emo.columns].fillna(1.0) PrintLog(LogLevel.INFO, f'[pipeline] → {len(df)} stocks, {len(df.columns)} features') # 添加 latest_close 列 (Meta Ranker 输入) df['latest_close'] = 0.0 for code in df.index: g = self.ctx.kline[self.ctx.kline['stock_code'] == code] if not g.empty: g_sorted = g.sort_values('trade_date') df.at[code, 'latest_close'] = float(g_sorted['close'].iloc[-1]) PrintLog(LogLevel.INFO, f'[pipeline] 完成: {len(df)} 只股票, {len(df.columns)} 维特征') return df