Files

107 lines
4.9 KiB
Python

"""
特征工程编排器 — 串联全部特征组,输出完整特征 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