Files
WQ_GUI/src/core/modeling/modeling_batch.py
duxin f97527598d feat: 各步骤产物改为原子化落盘,复用前校验 .done 完成标记
跳过条件从 Path.exists() 改为 is_file_complete(),产物成功落盘后统一补 mark_file_complete()。涉及 step1 水域掩膜、step3 去耀斑、SHP 栅格化缓存三处快路径。

step3 去耀斑新增清理分支:无 .done 标记的旧文件一律视为半成品,先删除再重跑(Kutser/Goodman/Hedley/SUGAR 四种方法一致)。

模型落盘改为 atomic_filepath:modeling_batch 批量建模与 automl_trainer 的 AutoML 结果均先写 .__wip 再原子替换,避免中断产生半个 joblib 文件。
2026-09-15 10:34:34 +08:00

1511 lines
66 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import numpy as np
import pandas as pd
import joblib
import os
from pathlib import Path
from typing import List, Dict, Union, Tuple, Optional
import warnings
warnings.filterwarnings('ignore')
# 机器学习模型导入 - 改为回归模型
from sklearn.base import BaseEstimator, TransformerMixin
from sklearn.svm import SVR
from sklearn.ensemble import RandomForestRegressor
from sklearn.neighbors import KNeighborsRegressor
from sklearn.linear_model import LinearRegression, Ridge, Lasso, ElasticNet
from sklearn.model_selection import GridSearchCV, RandomizedSearchCV, cross_val_score, KFold, train_test_split
from sklearn.metrics import mean_squared_error, mean_absolute_error, r2_score
from sklearn.cross_decomposition import PLSRegression
from sklearn.ensemble import GradientBoostingRegressor, AdaBoostRegressor, ExtraTreesRegressor
from sklearn.tree import DecisionTreeRegressor
from sklearn.neural_network import MLPRegressor
from sklearn.pipeline import Pipeline
from sklearn.impute import SimpleImputer
from sklearn.compose import TransformedTargetRegressor
from sklearn.preprocessing import StandardScaler
from joblib import parallel_backend
# 第三方模型导入
# try:
# import lightgbm as lgb
# LGB_AVAILABLE = True
# except ImportError:
# LGB_AVAILABLE = False
LGB_AVAILABLE = False # 注释掉lightgbm
# try:
# import catboost as cb
# CB_AVAILABLE = True
# except ImportError:
# CB_AVAILABLE = False
CB_AVAILABLE = False # 注释掉catboost
# 导入预处理模块
# 动态导入预处理模块
import sys
import os
# PyInstaller 打包环境感知EXE 模式下强制单核,防止 Windows 派生无限重启
is_frozen_env = getattr(sys, 'frozen', False)
safe_n_jobs = 1 if is_frozen_env else -1
from src.preprocessing.spectral_Preprocessing import Preprocessing, get_preprocessing_transformer
from src.core.utils.split_methods import spxy, ks
class _SafeFiniteTransformer(BaseEstimator, TransformerMixin):
"""Pipeline 安全网np.nan_to_num + clip确保无 inf/NaN 进入下游模型。
放在 Pipeline 的 preproc 与 model 之间无论上游MNF/比值除法)
产生何种极端值SVR 等严格校验的模型都不会因 inf 拒绝输入。
"""
def fit(self, X, y=None):
return self
def transform(self, X):
result = np.nan_to_num(np.asarray(X, dtype=np.float64),
nan=0.0, posinf=0.0, neginf=0.0)
result = np.clip(result, -1e15, 1e15)
return result
class WaterQualityModelingBatch:
"""水质参数反演批量建模类"""
def __init__(self, artifacts_dir: str = "models/artifacts"):
"""
初始化批量建模类
Args:
artifacts_dir: 模型保存目录
"""
self.artifacts_dir = Path(artifacts_dir)
self.artifacts_dir.mkdir(parents=True, exist_ok=True)
# 定义支持的回归模型及其参数网格
self.model_configs = {
'SVR': {
'model': SVR,
'params': {
'C': [0.1, 1, 10, 100],
'gamma': ['scale', 'auto', 0.001, 0.01, 0.1, 1],
'kernel': ['rbf', 'poly', 'sigmoid'],
'epsilon': [0.01, 0.1, 0.2]
},
'available': True
},
'RF': {
'model': RandomForestRegressor,
'params': {
'n_estimators': [50, 100, 200],
'max_depth': [None, 10, 20, 30],
'min_samples_split': [2, 5, 10],
'min_samples_leaf': [1, 2, 4]
},
'available': True
},
'KNN': {
'model': KNeighborsRegressor,
'params': {
'n_neighbors': [3, 5, 7, 9, 11],
'weights': ['uniform', 'distance'],
'metric': ['euclidean', 'manhattan', 'minkowski']
},
'available': True
},
'LinearRegression': {
'model': LinearRegression,
'params': {
'fit_intercept': [True, False]
},
'available': True
},
'Ridge': {
'model': Ridge,
'params': {
'alpha': [0.01, 0.1, 1, 10, 100],
'fit_intercept': [True, False]
},
'available': True
},
'Lasso': {
'model': Lasso,
'params': {
'alpha': [0.01, 0.1, 1, 10, 100],
'fit_intercept': [True, False],
'max_iter': [1000, 2000]
},
'available': True
},
'ElasticNet': {
'model': ElasticNet,
'params': {
'alpha': [0.01, 0.1, 1, 10],
'l1_ratio': [0.1, 0.3, 0.5, 0.7, 0.9],
'fit_intercept': [True, False],
'max_iter': [1000, 2000]
},
'available': True
},
'XGBoost': {
'model': None, # xgboost is removed, so set to None
'params': {
'n_estimators': [50, 100, 200],
'max_depth': [3, 6, 9],
'learning_rate': [0.01, 0.1, 0.2],
'subsample': [0.8, 0.9, 1.0]
},
'available': False
},
'LightGBM': {
'model': lgb.LGBMRegressor if LGB_AVAILABLE else None,
'params': {
'n_estimators': [50, 100, 200],
'max_depth': [3, 6, 9],
'learning_rate': [0.01, 0.1, 0.2],
'num_leaves': [31, 50, 100]
},
'available': LGB_AVAILABLE
},
'CatBoost': {
'model': cb.CatBoostRegressor if CB_AVAILABLE else None,
'params': {
'iterations': [50, 100, 200],
'depth': [3, 6, 9],
'learning_rate': [0.01, 0.1, 0.2],
'l2_leaf_reg': [1, 3, 5]
},
'available': CB_AVAILABLE
},
'PLS': {
'model': PLSRegression,
'params': {
'n_components': [2, 3, 5, 7, 10]
},
'available': True
},
'GradientBoosting': {
'model': GradientBoostingRegressor,
'params': {
'n_estimators': [50, 100, 200],
'learning_rate': [0.01, 0.1, 0.2],
'max_depth': [3, 5, 7],
'subsample': [0.8, 0.9, 1.0],
'min_samples_split': [2, 5, 10],
'min_samples_leaf': [1, 2, 4]
},
'available': True
},
'AdaBoost': {
'model': AdaBoostRegressor,
'params': {
'n_estimators': [50, 100, 200],
'learning_rate': [0.01, 0.1, 0.2],
'loss': ['linear', 'square', 'exponential']
},
'available': True
},
'DecisionTree': {
'model': DecisionTreeRegressor,
'params': {
'max_depth': [None, 5, 10, 20, 30],
'min_samples_split': [2, 5, 10],
'min_samples_leaf': [1, 2, 4],
'max_features': ['auto', 'sqrt', 'log2']
},
'available': True
},
'MLP': {
'model': MLPRegressor,
'params': {
'hidden_layer_sizes': [(50,), (100,), (50, 50), (100, 50)],
'activation': ['relu', 'tanh', 'logistic'],
'solver': ['adam', 'sgd'],
'alpha': [0.0001, 0.001, 0.01],
'learning_rate': ['constant', 'invscaling', 'adaptive'],
'max_iter': [1000, 2000]
},
'available': True
},
'ExtraTrees': {
'model': ExtraTreesRegressor,
'params': {
'n_estimators': [50, 100, 200],
'max_depth': [None, 10, 20, 30],
'min_samples_split': [2, 5, 10],
'min_samples_leaf': [1, 2, 4],
'max_features': ['auto', 'sqrt', 'log2']
},
'available': True
}
}
# 预处理方法列表
self.preprocessing_methods = [
"None", "MMS", "SS", "CT", "SNV", "MA", "SG", "MSC", "D1", "D2", "DT", "WVAE",
"DualStream_MNF", "Physical_Only"
]
# 样本划分方法列表
self.split_methods = ["random", "spxy", "ks"]
self.results = {}
self.best_models = {}
# ★ 黄金光谱区间:统一收缩到 400-1000nm消除边缘噪声和跨传感器不一致
WAVELENGTH_RANGE = (400.0, 1000.0)
@staticmethod
def _is_wavelength_column(col_name: str) -> bool:
"""判断列名是否为波长值(纯数字字符串,如 '374.285004'"""
try:
float(str(col_name))
return True
except (ValueError, TypeError):
return False
@classmethod
def _is_wavelength_in_range(cls, col_name: str) -> bool:
"""判断列名是否为波长值且在黄金区间 [400, 1000] nm 内。"""
try:
wl = float(str(col_name))
lo, hi = cls.WAVELENGTH_RANGE
return lo <= wl <= hi
except (ValueError, TypeError):
return False
@staticmethod
def _is_wqi_column(col_name: str) -> bool:
"""判断列名是否为 WQI 水质指数列('WQI_' 前缀)"""
return str(col_name).startswith('WQI_')
@classmethod
def _extract_train_wavelengths(cls, columns) -> List[float]:
"""从列名列表中提取波长值float 列表),自动过滤到黄金区间。
遍历列名,将所有可转为 float 且落在 [400, 1000] nm 的列名提取为波长列表。
用于写入模型 metadata['train_wavelengths'],供推理端光谱重采样。
"""
lo, hi = cls.WAVELENGTH_RANGE
wl_list = []
for c in columns:
try:
wl = float(str(c))
if lo <= wl <= hi:
wl_list.append(wl)
except (ValueError, TypeError):
pass
return wl_list
def _extract_feature_columns(self, data: pd.DataFrame,
feature_start_column: Union[int, str, None] = None
) -> Tuple[pd.DataFrame, List[int]]:
"""从 DataFrame 中提取特征列(基于列名语义,而非位置索引)
策略(优先级从高到低):
1) 如果 feature_start_column 是列名str定位该列并取之后的列
2) 如果 feature_start_column 是整数,兼容旧逻辑(位置索引)
3) 如果 feature_start_column 为 None自动识别
- 保留所有纯数字列名(波长列)
- 保留所有 WQI_ 前缀列(水质指数列)
- 跳过坐标/元数据列
Returns:
X: 特征 DataFrame保留列名
feature_col_indices: 特征列在原 data 中的位置索引列表
"""
all_cols = list(data.columns)
if feature_start_column is not None:
# ── 兼容旧逻辑:按列名或索引位置截取 ──
if isinstance(feature_start_column, str):
if feature_start_column not in data.columns:
raise ValueError(
f"指定的特征开始列 '{feature_start_column}' 不存在于数据中"
)
start_idx = data.columns.get_loc(feature_start_column)
print(f"[特征提取] 按列名 '{feature_start_column}' 定位 → 索引 {start_idx}")
else:
start_idx = int(feature_start_column)
print(f"[特征提取] 按位置索引 {start_idx} 截取")
X = data.iloc[:, start_idx:]
feature_indices = list(range(start_idx, len(all_cols)))
else:
# ── 智能识别:按列名语义过滤 ──
# 黑名单:坐标列、元数据列(不会被误判为特征)
_meta_patterns = {
'x_coord', 'y_coord', 'pixel_x', 'pixel_y',
'longitude', 'latitude', 'lon', 'lat',
'id', 'station', 'sample_id',
}
feature_indices = []
for i, col in enumerate(all_cols):
col_lower = str(col).lower().strip()
# 跳过元数据列
if col_lower in _meta_patterns:
continue
# 保留波长列(仅黄金区间 400-1000nm
if self._is_wavelength_in_range(col):
feature_indices.append(i)
# 保留 WQI 列
elif self._is_wqi_column(col):
feature_indices.append(i)
# 其他:跳过(可能是目标列或其他非特征列)
if not feature_indices:
raise ValueError(
"智能特征提取失败:未找到任何波长列或 WQI 列。"
f"CSV 列名: {all_cols[:10]}..."
)
X = data.iloc[:, feature_indices]
print(f"[特征提取] 智能识别: {len(feature_indices)} 个特征列 "
f"(波长列 + WQI 列)")
print(f"[特征提取] 特征数据形状: {X.shape}")
return X, feature_indices
def load_data_batch(self, csv_path: str,
feature_start_column: Union[int, str, None] = None
) -> Tuple[pd.DataFrame, Dict[str, pd.Series]]:
"""批量加载 CSV 数据,自动识别特征列与目标列
改造要点v2
- 特征列按列名语义提取(波长数字 / WQI_ 前缀),不再按硬编码位置一刀切
- 目标列 = 不在特征列中、且非系统保留列ID/坐标等)的数值列
- X 保持为 DataFrame列名保留波长信息供后续提取 train_wavelengths
Args:
csv_path: CSV 文件路径
feature_start_column: (可选)旧版兼容参数:
- str: 特征起始列名
- int: 特征起始列索引
- None: 自动智能识别
Returns:
X: 特征数据 (DataFrame保留列名)
y_dict: 目标值数据字典,键为列名
"""
# 读取 CSV 数据,处理空字符串和缺失值
try:
data = pd.read_csv(csv_path, na_values=['', ' ', 'NaN', 'nan', 'NULL', 'null'])
except pd.errors.EmptyDataError:
raise ValueError(f"CSV文件 '{csv_path}' 为空或不存在")
except Exception as e:
raise ValueError(f"读取CSV文件 '{csv_path}' 时出错: {e}")
print("数据清理...")
original_shape = data.shape
# 将空字符串替换为 NaN
data = data.replace(r'^\s*$', np.nan, regex=True)
# 对于数值列,将无法转换为数字的字符串替换为 NaN
for col in data.columns:
try:
data[col] = pd.to_numeric(data[col], errors='coerce')
except Exception:
pass
cleaned_shape = data.shape
if cleaned_shape != original_shape:
print(f"数据清理完成: {original_shape[0]}{original_shape[1]}"
f"-> {cleaned_shape[0]}{cleaned_shape[1]}")
print(f"数据加载完成,总列数: {data.shape[1]}")
print(f"所有列名: {list(data.columns)[:10]}..." if len(data.columns) > 10
else f"所有列名: {list(data.columns)}")
# ── 提取特征列(基于列名语义) ──
X, feature_indices = self._extract_feature_columns(
data, feature_start_column
)
# ── 提取目标列(不在特征列中 + 非系统保留列 + 数值类型) ──
feature_idx_set = set(feature_indices)
ignore_cols = {
'ID', 'id', 'Id',
'Longitude', 'Latitude', 'Lon', 'Lat',
'longitude', 'latitude', 'lon', 'lat',
'Station', 'station', 'sample_id',
'x_coord', 'y_coord', 'pixel_x', 'pixel_y',
}
y_dict = {}
for i, col_name in enumerate(data.columns):
if i in feature_idx_set:
continue # 已在特征矩阵中
col_str = str(col_name).strip()
if col_str in ignore_cols:
print(f" 跳过 '{col_name}': 系统保留列")
continue
y_series = data[col_name]
if not pd.api.types.is_numeric_dtype(y_series):
print(f" 跳过 '{col_name}': 非数值类型")
continue
if y_series.isna().all():
print(f" 跳过 '{col_name}': 所有值为空")
continue
y_dict[col_name] = y_series
print(f" 目标列 '{col_name}': {y_series.count()} 个非空值, "
f"范围: {y_series.min():.4f} ~ {y_series.max():.4f}")
print(f"特征数据形状: {X.shape}")
print(f"有效目标列数量: {len(y_dict)}")
return X, y_dict
def load_data_single(self, csv_path: str, target_column_name: str,
feature_start_column: Union[int, str, None] = None
) -> Tuple[pd.DataFrame, pd.Series]:
"""
加载单个目标列的CSV数据v2基于列名语义提取特征
Args:
csv_path: CSV文件路径
target_column_name: 目标列名
feature_start_column: 特征起始位置可选None=智能识别)
Returns:
X: 特征数据 (DataFrame保留列名)
y: 目标值数据
"""
data = pd.read_csv(csv_path)
# 检查目标列是否存在
if target_column_name not in data.columns:
raise ValueError(f"目标列 '{target_column_name}' 不存在于数据中")
# 提取目标值
y = data[target_column_name]
# 去除 y 值为空的行
mask = ~y.isna()
data_cleaned = data[mask]
y = data_cleaned[target_column_name]
# 提取特征(基于列名语义,排除目标列自身)
X, _ = self._extract_feature_columns(
data_cleaned, feature_start_column
)
# 确保目标列不在 X 中
if target_column_name in X.columns:
X = X.drop(columns=[target_column_name])
print(f"目标列 '{target_column_name}' 数据加载完成:")
print(f" 样本数量: {X.shape[0]}")
print(f" 特征数量: {X.shape[1]}")
print(f" 目标值范围: {y.min():.4f} ~ {y.max():.4f}")
print(f" 目标值均值: {y.mean():.4f}")
return X, y
def preprocess_data(self, X: pd.DataFrame, method: str) -> np.ndarray:
"""
数据预处理
Args:
X: 原始特征数据
method: 预处理方法
Returns:
预处理后的数据
"""
print(f"应用预处理方法: {method}")
# 如果方法为None直接返回原始数据
if method == "None" or method is None:
print("跳过预处理,使用原始数据")
return X.values
try:
X_processed = Preprocessing(method, X)
# 确保返回的是numpy数组
if isinstance(X_processed, pd.DataFrame):
X_processed = X_processed.values
print(f"预处理完成,数据形状: {X_processed.shape}")
return X_processed
except Exception as e:
print(f"预处理失败: {e}")
print("使用原始数据")
return X.values
def random(self, data, label, test_ratio=0.2, random_state=123):
"""
随机划分数据集
Args:
data: shape (n_samples, n_features)
label: shape (n_sample, )
test_ratio: 测试集比例,默认: 0.2
random_state: 随机种子,默认: 123
Returns:
X_train: (n_samples, n_features)
X_test: (n_samples, n_features)
y_train: (n_sample, )
y_test: (n_sample, )
"""
X_train, X_test, y_train, y_test = train_test_split(
data, label, test_size=test_ratio, random_state=random_state
)
return X_train, X_test, y_train, y_test
def spxy(self, data, label, test_size=0.2):
"""SPXY算法划分数据集委托至 src.core.utils.split_methods.spxy"""
return spxy(data, label, test_size=test_size)
def ks(self, data, label, test_size=0.2):
"""Kennard-Stone算法划分数据集委托至 src.core.utils.split_methods.ks"""
return ks(data, label, test_size=test_size)
def split_data(self, X: np.ndarray, y: pd.Series, method: str = "random",
test_size: float = 0.2, random_state: int = 42) -> Tuple[np.ndarray, np.ndarray, np.ndarray, np.ndarray]:
"""
根据指定方法划分数据集
Args:
X: 特征数据
y: 目标值数据
method: 划分方法 ("random", "spxy", "ks")
test_size: 测试集比例
random_state: 随机种子仅对random方法有效
Returns:
X_train, X_test, y_train, y_test
"""
print(f"使用 {method} 方法划分数据集")
if method == "random":
return self.random(X, y, test_ratio=test_size, random_state=random_state)
elif method == "spxy":
return self.spxy(X, y, test_size=test_size)
elif method == "ks":
return self.ks(X, y, test_size=test_size)
else:
raise ValueError(f"不支持的划分方法: {method}. 支持的方法: {self.split_methods}")
def train_single_model(self, X_raw: pd.DataFrame, y: pd.Series, model_name: str,
cv_folds: int = 5, scoring: str = 'neg_mean_squared_error',
test_size: float = 0.2, random_state: int = 42,
split_method: str = "random",
preprocess_method: str = "None",
mnf_n_components=0.95) -> Dict:
"""
训练单个回归模型Pipeline 化preprocess_method 字符串内部构造 Pipeline
scaler/MSC.mean_spectrum_ 等状态被绑定在 best_model 上CV 与 test 评估
都在「只拟合训练 fold 的 scaler」之上避免传统「X_full → scaler.fit →
split → CV」造成的数据泄露
Args:
X_raw: 原始特征数据未经预处理DataFrame 形态方便 Pipeline 内部转换)
y: 目标值数据
model_name: 模型名称
preprocess_method: 预处理方法字符串(如 'None' / 'SS' / 'MSC'
训练结束后 best_model 字段即为 sklearn Pipeline。
mnf_n_components: MNFTransformer 成分数(仅 DualStream_MNF 使用);
float(0~1)=自动按累计方差截断int=固定数量。
其余参数cv_folds / scoring / test_size / random_state / split_method
含义保持不变。
Returns:
训练结果字典;'model' 字段现在是 sklearn.pipeline.Pipeline含 scaler
"""
if model_name not in self.model_configs:
raise ValueError(f"不支持的模型: {model_name}")
config = self.model_configs[model_name]
if not config['available']:
print(f"模型 {model_name} 不可用,请安装相应的库")
return None
print(f"开始训练模型: {model_name} (预处理: {preprocess_method})")
# ═══════════════════════════════════════════════════════════
# ★ DualStream_MNF只取纯光谱列WQI 由 Pipeline 内 PhysicalExtractor 动态计算
# ═══════════════════════════════════════════════════════════
if preprocess_method in ("DualStream_MNF", "Physical_Only"):
_spec_cols = [c for c in X_raw.columns if self._is_wavelength_in_range(c)]
X_raw = X_raw[_spec_cols]
print(f"[{preprocess_method}] 精简为纯光谱(400-1000nm): {X_raw.shape[1]}"
f"({X_raw.columns[0]} ~ {X_raw.columns[-1]} nm)")
# 使用指定方法分割训练集和测试集
X_train, X_test, y_train, y_test = self.split_data(
X_raw, y, method=split_method, test_size=test_size, random_state=random_state
)
print(f"数据分割完成:")
print(f" 训练集样本数: {X_train.shape[0]}")
print(f" 测试集样本数: {X_test.shape[0]}")
# 构造 base_model
if callable(config['model']):
base_model = config['model']()
else:
base_model = config['model']
# 特殊处理某些模型
if model_name == 'CatBoost':
base_model.set_params(verbose=False)
elif model_name == 'LightGBM':
base_model.set_params(verbose=-1)
# ★ 致命修复 1SVR 对 y 尺度极度敏感。若 y 未缩放(如浊度 100~1000
# epsilon=0.1 相对 y 量级几乎为零 → 所有样本都变成支持向量 → 过拟合/躺平。
# TransformedTargetRegressor 在 fit 时自动 StandardScaler(y)
# predict 时自动 inverse_transform对调用方完全透明。
_is_svr = (config['model'] == SVR)
if _is_svr:
base_model = TransformedTargetRegressor(
regressor=base_model,
transformer=StandardScaler()
)
# ============ 关键:把预处理器塞进 Pipeline ============
# DualStream_MNF / Physical_Only 需要波长列表,供 PhysicalFeatureExtractor 定位波段
_wl_list = None
if preprocess_method in ("DualStream_MNF", "Physical_Only"):
_wl_list = self._extract_train_wavelengths(X_raw.columns)
preproc = get_preprocessing_transformer(preprocess_method, wavelengths=_wl_list,
mnf_n_components=mnf_n_components)
pipeline = Pipeline([
('imputer', SimpleImputer(strategy='median')),
('preproc', preproc),
('cleaner', _SafeFiniteTransformer()),
('model', base_model),
])
# ★ 入口清洗inf → NaN后续 SimpleImputer/cleaner 接力处理
X_train = np.asarray(X_train, dtype=np.float64)
X_test = np.asarray(X_test, dtype=np.float64)
X_train[~np.isfinite(X_train)] = np.nan
X_test[~np.isfinite(X_test)] = np.nan
# 以「步骤名__参数名」的格式索引参数网格
# config['params'] 是模型层的(无 __统一加 model__ 前缀。
# ★ SVR 被 TransformedTargetRegressor 包裹后,实际模型在 model.regressor_ 下,
# GridSearchCV 参数路径变为 model__regressor__{param}
if _is_svr:
prefixed_params = {
f"model__regressor__{k}": v for k, v in config['params'].items()
}
else:
prefixed_params = {
f"model__{k}": v for k, v in config['params'].items()
}
# 全量网格搜索SVR 超参组合仅 216 种4×6×3×3
# 穷举远优于 RandomizedSearchCV(n_iter=10) 的随机抽样
cv_strategy = KFold(n_splits=cv_folds, shuffle=True, random_state=random_state)
grid_search = GridSearchCV(
pipeline,
prefixed_params,
cv=cv_strategy,
scoring=scoring,
n_jobs=safe_n_jobs,
verbose=1,
)
grid_search.fit(X_train, y_train)
# 获取最佳模型(已是 Pipeline
best_model = grid_search.best_estimator_
# 交叉验证评估在训练集上cross_val_score 会对 Pipeline 重 clone
# 保证每个 fold 重 fit 预处理CV 评分反映「无泄露」真实泛化能力
cv_scores = cross_val_score(best_model, X_train, y_train, cv=cv_strategy,
scoring=scoring, n_jobs=safe_n_jobs)
# 计算训练集上的回归指标Pipeline 内 fit_transform 只发生一次,已 fit 完毕)
y_train_pred = best_model.predict(X_train)
train_mse = mean_squared_error(y_train, y_train_pred)
train_mae = mean_absolute_error(y_train, y_train_pred)
train_r2 = r2_score(y_train, y_train_pred)
train_rmse = np.sqrt(train_mse)
# 计算测试集上的回归指标(用训练集 fit 出的 scaler正确的 deploy-time 行为)
y_test_pred = best_model.predict(X_test)
test_mse = mean_squared_error(y_test, y_test_pred)
test_mae = mean_absolute_error(y_test, y_test_pred)
test_r2 = r2_score(y_test, y_test_pred)
test_rmse = np.sqrt(test_mse)
result = {
'model': best_model,
'best_params': grid_search.best_params_,
'best_score': grid_search.best_score_,
'cv_mean': cv_scores.mean(),
'cv_std': cv_scores.std(),
'cv_scores': cv_scores,
# 训练集指标
'train_mse': train_mse,
'train_mae': train_mae,
'train_rmse': train_rmse,
'train_r2': train_r2,
# 测试集指标
'test_mse': test_mse,
'test_mae': test_mae,
'test_rmse': test_rmse,
'test_r2': test_r2,
# 数据分割信息
'train_size': X_train.shape[0],
'test_size': X_test.shape[0],
'split_method': split_method,
# Pipeline 信息(用于诊断 / metadata
'preprocess_method': preprocess_method,
'is_pipeline': isinstance(best_model, Pipeline),
}
print(f"模型 {model_name} 训练完成:")
print(f" 最佳参数: {result['best_params']}")
print(f" 最佳得分: {result['best_score']:.4f}")
print(f" CV均值: {result['cv_mean']:.4f} ± {result['cv_std']:.4f}")
print(f" 训练集指标:")
print(f" R²: {result['train_r2']:.4f}")
print(f" RMSE: {result['train_rmse']:.4f}")
print(f" MAE: {result['train_mae']:.4f}")
print(f" 测试集指标:")
print(f" R²: {result['test_r2']:.4f}")
print(f" RMSE: {result['test_rmse']:.4f}")
print(f" MAE: {result['test_mae']:.4f}")
# ★ MNF 波段选择诊断:从 Pipeline 中提取 MNFTransformer 并打印波段信息
if isinstance(best_model, Pipeline):
_mnf = self._find_mnf_in_pipeline(best_model)
if _mnf is not None and hasattr(_mnf, 'eigvals_w_'):
_mnf._print_band_selection()
return result
def save_model(self, model, target_column_name: str, preprocess_method: str, model_name: str,
metadata: Dict = None):
"""
保存模型,使用目标列名作为文件名的一部分
Args:
model: 训练好的模型
target_column_name: 目标列名
preprocess_method: 预处理方法名称
model_name: 模型名称
metadata: 模型元数据
"""
# 清理目标列名,移除可能的特殊字符
safe_target_name = "".join(c for c in target_column_name if c.isalnum() or c in ('-', '_')).rstrip()
filename = f"{safe_target_name}_{preprocess_method}_{model_name}.joblib"
filepath = self.artifacts_dir / filename
# ═══════════════════════════════════════════════════════════
# ★ MNF 诊断:提前提取并注入 metadata在 joblib.dump 之前)
# ═══════════════════════════════════════════════════════════
if metadata is None:
metadata = {}
_mnf_for_meta = None
if isinstance(model, Pipeline):
_mnf_for_meta = self._find_mnf_in_pipeline(model)
if _mnf_for_meta is not None and hasattr(_mnf_for_meta, 'eigvals_w_'):
diag_dict = _mnf_for_meta.get_diagnostics_dict(target_name=target_column_name)
metadata['mnf_diagnostics'] = diag_dict
# 保存模型和元数据
save_data = {
'model': model,
'target_column_name': target_column_name,
'preprocess_method': preprocess_method,
'model_name': model_name,
'metadata': metadata
}
from src.utils.util import atomic_filepath
with atomic_filepath(str(filepath)) as _tmp:
joblib.dump(save_data, _tmp)
print(f"模型已保存: {filepath}")
# ═══════════════════════════════════════════════════════════
# ★ MNF 诊断 JSON 导出(独立文件,方便外部工具读取)
# ═══════════════════════════════════════════════════════════
if _mnf_for_meta is not None and hasattr(_mnf_for_meta, 'eigvals_w_'):
import json as _json
json_filename = f"{safe_target_name}_{preprocess_method}_{model_name}_mnf_diag.json"
json_path = self.artifacts_dir / json_filename
with open(str(json_path), 'w', encoding='utf-8') as _f:
_json.dump(diag_dict, _f, indent=2, ensure_ascii=False)
print(f"[MNF 诊断] JSON 已保存: {json_path}")
# ═══════════════════════════════════════════════════════════
# ★ C++/Rust 部署导出:提取 MNF + SVR 纯量矩阵
# ═══════════════════════════════════════════════════════════
if not isinstance(model, Pipeline):
return
_mnf = self._find_mnf_in_pipeline(model)
_svr = self._find_svr_in_pipeline(model)
_phys = self._find_physical_extractor_in_pipeline(model)
if _mnf is not None and _svr is not None:
_deploy = {
'mnf_mean': np.asarray(_mnf.mean_).ravel(),
'mnf_W': np.asarray(_mnf.W_mnf_),
'mnf_n_components': int(getattr(_mnf, 'n_components_', _mnf.n_components)),
}
# ★ 附加 MNF 特征值信息(用于波段选择诊断)
if hasattr(_mnf, 'eigvals_w_'):
_deploy['mnf_eigvals'] = np.asarray(_mnf.eigvals_w_)
_deploy['mnf_eigvals_ratio'] = np.asarray(_mnf.eigvals_ratio_)
_deploy['mnf_eigvals_cumsum'] = np.asarray(_mnf.eigvals_cumsum_)
# SVR 部署矩阵rbf kernel 需要 support_vectors_
_deploy['svr_dual_coef'] = np.asarray(_svr.dual_coef_)
_deploy['svr_intercept'] = np.asarray(_svr.intercept_)
if hasattr(_svr, 'support_vectors_'):
_deploy['svr_support_vectors'] = np.asarray(_svr.support_vectors_)
_deploy['svr_gamma'] = float(getattr(_svr, '_gamma', 1.0))
# kernel type for C++ dispatch
_kernel = str(getattr(_svr, 'kernel', 'rbf')).lower()
_deploy['svr_kernel'] = _kernel
if _kernel == 'poly':
_deploy['svr_degree'] = int(getattr(_svr, 'degree', 3))
_deploy['svr_coef0'] = float(getattr(_svr, 'coef0', 0.0))
# 物理特征索引C++ 端据此定位波段列)
if _phys is not None:
_deploy['physical_wavelengths'] = np.asarray(
_phys.wavelengths, dtype=np.float64
)
_feat_map = {}
for name, ftype, ia, ib in _phys._feat_cols_:
_feat_map[name] = {
'ftype': str(ftype),
'wl_a': float(_phys.wavelengths[ia]),
'wl_b': float(_phys.wavelengths[ib]),
'idx_a': int(ia),
'idx_b': int(ib),
}
_deploy['physical_feature_map'] = _feat_map
npz_path = str(filepath).replace('.joblib', '_deploy.npz')
np.savez_compressed(npz_path, **_deploy)
print(f"[部署导出] MNF+SVR 纯量矩阵已保存: {npz_path}")
elif _mnf is None and _svr is not None:
print("[部署导出] 仅检测到 SVR 未检测到 MNFTransformer跳过 .npz 导出")
# ═══════════════════════════════════════════════════════════
# 部署矩阵提取辅助方法:从 Pipeline 中定位特定 Transformer
# ═══════════════════════════════════════════════════════════
@staticmethod
def _find_transformer_in_pipeline(pipeline, target_class_name):
"""在 Pipeline 的 'preproc' 步骤中递归查找指定类型的 Transformer。
支持 FeatureUnion 嵌套:若 preproc 是 FeatureUnion
则遍历其 transformer_list 逐个子变压器查找。
"""
from sklearn.pipeline import FeatureUnion
preproc = pipeline.named_steps.get('preproc')
if preproc is None:
return None
if isinstance(preproc, FeatureUnion):
for name, trans in preproc.transformer_list:
found = WaterQualityModelingBatch._find_transformer_in_step(
trans, target_class_name)
if found is not None:
return found
return None
return WaterQualityModelingBatch._find_transformer_in_step(
preproc, target_class_name)
@staticmethod
def _find_transformer_in_step(step, target_class_name):
"""在单个 Pipeline 步骤(可能是 Pipeline 自身)中查找 Transformer。"""
if step.__class__.__name__ == target_class_name:
return step
# 处理 Pipeline 嵌套
if hasattr(step, 'named_steps'):
for sub_name, sub_step in step.named_steps.items():
found = WaterQualityModelingBatch._find_transformer_in_step(
sub_step, target_class_name)
if found is not None:
return found
return None
@staticmethod
def _find_mnf_in_pipeline(pipeline):
return WaterQualityModelingBatch._find_transformer_in_pipeline(
pipeline, 'MNFTransformer')
@staticmethod
def _find_svr_in_pipeline(pipeline):
model_step = pipeline.named_steps.get('model')
if model_step is None:
return None
# 直接就是 SVR旧模型 / 非 SVR 的其他回归器)
if model_step.__class__.__name__ == 'SVR':
return model_step
# ★ TransformedTargetRegressor 包裹的 SVR修复 1 引入)
if hasattr(model_step, 'regressor_'):
inner = model_step.regressor_
if inner.__class__.__name__ == 'SVR':
return inner
return None
@staticmethod
def _find_physical_extractor_in_pipeline(pipeline):
return WaterQualityModelingBatch._find_transformer_in_pipeline(
pipeline, 'PhysicalFeatureExtractor')
def train_models_batch(self, csv_path: str, feature_start_column: Union[int, str],
preprocessing_methods: Union[str, List[str]] = "None",
model_names: Union[str, List[str]] = "RF",
split_methods: Union[str, List[str]] = "random",
cv_folds: int = 5,
scoring: str = 'neg_mean_squared_error',
test_size: float = 0.2,
random_state: int = 42) -> Dict:
"""
批量训练多个目标列的模型
Args:
csv_path: 数据文件路径
feature_start_column: 特征开始列索引int或列名str
preprocessing_methods: 预处理方法列表
model_names: 模型名称列表
split_methods: 数据划分方法列表
cv_folds: 交叉验证折数
scoring: 评分指标(回归指标)
test_size: 测试集比例
random_state: 随机种子
Returns:
所有模型的训练结果
"""
# 转换为列表
if isinstance(preprocessing_methods, str):
preprocessing_methods = [preprocessing_methods]
if isinstance(model_names, str):
model_names = [model_names]
if isinstance(split_methods, str):
split_methods = [split_methods]
# 加载数据
X_raw, y_dict = self.load_data_batch(csv_path, feature_start_column)
all_results = {}
# ═══════════════════════════════════════════════════════════════
# ★ 第一遍预扫描Pre-scan—— 确定全局统一的 MNF 成分数
# ═══════════════════════════════════════════════════════════════
global_max_mnf_components = 50 # 安全兜底默认值
_has_dual_stream = any("DualStream_MNF" in str(pm)
for pm in preprocessing_methods)
if _has_dual_stream:
print(f"\n{'='*80}")
print(f"[MNF 全局对齐] 第一遍 — 预扫描各目标列 MNF 成分需求...")
print(f"{'='*80}")
_per_target_mnf_counts: Dict[str, int] = {}
from src.preprocessing.spectral_Preprocessing import MNFTransformer
for target_column_name, y in y_dict.items():
mask = ~y.isna()
if mask.sum() <= 1:
print(f" [{target_column_name}] 有效样本不足,跳过")
continue
X_clean = X_raw[mask]
# 只保留纯光谱列(与 DualStream_MNF Pipeline 内的逻辑一致400-1000nm
_spec_cols = [c for c in X_clean.columns
if self._is_wavelength_in_range(c)]
X_spec = X_clean[_spec_cols].values.astype(np.float64)
try:
_mnf = MNFTransformer(n_components=0.95)
_mnf.fit(X_spec)
_cnt = int(getattr(_mnf, 'n_components_', 0))
_per_target_mnf_counts[target_column_name] = _cnt
print(f" [{target_column_name}] 需要 {_cnt} 个 MNF 成分 "
f"(累计方差 {_mnf.eigvals_cumsum_[_cnt - 1] * 100:.2f}%)")
except Exception as _e:
print(f" [{target_column_name}] ⚠ MNF 预扫描失败: {_e}"
f"将使用兜底值")
if _per_target_mnf_counts:
_min_cnt = min(_per_target_mnf_counts.values())
_max_cnt = max(_per_target_mnf_counts.values())
global_max_mnf_components = _max_cnt
print(f"\n[MNF 全局对齐] 预扫描完成。"
f"各指标需求范围: {_min_cnt}~{_max_cnt}"
f"全局统一对齐特征数设定为: {global_max_mnf_components}")
else:
print(f"[MNF 全局对齐] 预扫描无有效结果,"
f"使用兜底值: {global_max_mnf_components}")
else:
print(f"\n[MNF 全局对齐] 未启用 DualStream_MNF跳过预扫描")
# ═══════════════════════════════════════════════════════════════
# ★ 第二遍:正式训练 — 所有目标使用统一的 MNF 成分数
# ═══════════════════════════════════════════════════════════════
for target_column_name, y in y_dict.items():
print(f"\n{'='*80}")
print(f"开始训练目标列: {target_column_name}")
if _has_dual_stream:
print(f"[MNF 全局对齐] 统一使用 {global_max_mnf_components} 个 MNF 成分")
print(f"{'='*80}")
# 创建该目标列的子目录
target_artifacts_dir = self.artifacts_dir / target_column_name
target_artifacts_dir.mkdir(parents=True, exist_ok=True)
# 临时更改artifacts_dir
original_artifacts_dir = self.artifacts_dir
self.artifacts_dir = target_artifacts_dir
try:
# 去除该目标列的空值
mask = ~y.isna()
if mask.sum() == 0:
print(f"目标列 '{target_column_name}' 无有效数据,跳过")
continue
X_clean = X_raw[mask]
y_clean = y[mask]
print(f"有效样本数: {len(y_clean)}")
# 训练该目标列的所有模型组合
target_results = self.train_models_single_target(
X_clean, y_clean, target_column_name,
preprocessing_methods, model_names, split_methods,
cv_folds, scoring, test_size, random_state,
mnf_n_components=global_max_mnf_components,
)
all_results[target_column_name] = target_results
except Exception as e:
print(f"训练目标列 '{target_column_name}' 时出错: {e}")
continue
finally:
# 恢复原始artifacts_dir
self.artifacts_dir = original_artifacts_dir
# 保存所有结果的汇总
self._save_batch_results_summary(all_results)
return all_results
def train_models_single_target(self, X_raw: pd.DataFrame, y: pd.Series, target_column_name: str,
preprocessing_methods: List[str], model_names: List[str],
split_methods: List[str], cv_folds: int, scoring: str,
test_size: float, random_state: int,
mnf_n_components=0.95) -> Dict:
"""
训练单个目标列的所有模型组合
"""
results = {}
# 遍历所有组合
for split_method in split_methods:
for preprocess_method in preprocessing_methods:
for model_name in model_names:
combo_key = f"{split_method}_{preprocess_method}_{model_name}"
print(f"\n{'-' * 60}")
print(f"训练组合: {combo_key}")
print(f"{'-' * 60}")
try:
# 不再外部 Preprocessing——改传给 train_single_model 由 Pipeline 处理
result = self.train_single_model(
X_raw, y, model_name,
cv_folds, scoring, test_size, random_state, split_method,
preprocess_method=preprocess_method,
mnf_n_components=mnf_n_components,
)
if result is not None:
# 提取训练波长列表(写入 metadata供推理端光谱重采样
train_wavelengths = self._extract_train_wavelengths(
X_raw.columns
)
print(f"[波长记忆] 提取到 {len(train_wavelengths)} 个训练波长: "
f"{train_wavelengths[0]:.2f} ~ {train_wavelengths[-1]:.2f} nm")
# 保存模型result['model'] 已是 sklearn Pipeline
metadata = {
'target_column_name': target_column_name,
'cv_mean': result['cv_mean'],
'cv_std': result['cv_std'],
'best_params': result['best_params'],
'data_shape': X_raw.shape,
'target_range': [float(y.min()), float(y.max())],
'train_r2': result['train_r2'],
'train_rmse': result['train_rmse'],
'train_mae': result['train_mae'],
'test_r2': result['test_r2'],
'test_rmse': result['test_rmse'],
'test_mae': result['test_mae'],
'train_size': result['train_size'],
'test_size': result['test_size'],
'split_method': result['split_method'],
'preprocess_method': preprocess_method,
'is_pipeline': result.get('is_pipeline', False),
# ★ 波长记忆:推理端可据此进行跨传感器光谱重采样
'train_wavelengths': train_wavelengths,
'train_columns': list(X_raw.columns),
}
self.save_model(result['model'], target_column_name,
f"{split_method}_{preprocess_method}",
model_name, metadata)
results[combo_key] = result
except Exception as e:
print(f"训练组合 {combo_key} 失败: {e}")
continue
# 保存该目标列的结果摘要
self._save_single_target_results_summary(target_column_name, results)
return results
def _save_single_target_results_summary(self, target_column_name: str, results: Dict):
"""保存单个目标列的结果摘要"""
if not results:
print(f"目标列 '{target_column_name}' 没有训练结果")
return
summary_data = []
for combo_key, result in results.items():
# 分离划分方法、预处理方法和建模方法
parts = combo_key.split('_', 2)
split_method = parts[0] if len(parts) > 0 else ''
preprocess_method = parts[1] if len(parts) > 1 else ''
model_method = parts[2] if len(parts) > 2 else ''
summary_data.append({
'划分方法': split_method,
'预处理方法': preprocess_method,
'建模方法': model_method,
'CV均值': result['cv_mean'],
'CV标准差': result['cv_std'],
'最佳得分': result['best_score'],
'训练集R²': result['train_r2'],
'训练集RMSE': result['train_rmse'],
'训练集MAE': result['train_mae'],
'训练集MSE': result['train_mse'],
'测试集R²': result['test_r2'],
'测试集RMSE': result['test_rmse'],
'测试集MAE': result['test_mae'],
'测试集MSE': result['test_mse'],
'训练样本数': result['train_size'],
'测试样本数': result['test_size'],
'最佳参数': str(result['best_params'])
})
summary_df = pd.DataFrame(summary_data)
# 按测试集R²降序排列R²越大越好
summary_df = summary_df.sort_values('测试集R²', ascending=False)
# 清理目标列名,移除可能的特殊字符
safe_target_name = "".join(c for c in target_column_name if c.isalnum() or c in ('-', '_')).rstrip()
# 保存详细结果CSV中文版
detailed_path = self.artifacts_dir / f"{safe_target_name}_detailed_results.csv"
summary_df.to_csv(detailed_path, index=False, encoding='utf-8-sig')
# 保存简化版本用于兼容性(英文版)
summary_data_simple = []
for combo_key, result in results.items():
summary_data_simple.append({
'combination': combo_key,
'cv_mean': result['cv_mean'],
'cv_std': result['cv_std'],
'best_score': result['best_score'],
'train_r2': result['train_r2'],
'train_rmse': result['train_rmse'],
'train_mae': result['train_mae'],
'test_r2': result['test_r2'],
'test_rmse': result['test_rmse'],
'test_mae': result['test_mae'],
'train_size': result['train_size'],
'test_size': result['test_size'],
'split_method': result.get('split_method', 'unknown'),
'best_params': str(result['best_params'])
})
summary_df_simple = pd.DataFrame(summary_data_simple)
summary_df_simple = summary_df_simple.sort_values('test_r2', ascending=False)
simple_summary_path = self.artifacts_dir / f"{safe_target_name}_training_summary.csv"
summary_df_simple.to_csv(simple_summary_path, index=False)
print(f"\n{'-' * 60}")
print(f"目标列 '{target_column_name}' 训练结果摘要:")
print(f"{'-' * 60}")
print(summary_df[
['划分方法', '预处理方法', '建模方法', '训练集R²', '测试集R²', '训练集RMSE', '测试集RMSE', 'CV均值']].to_string(
index=False))
print(f"\n详细结果已保存: {detailed_path}")
print(f"简化结果已保存: {simple_summary_path}")
def _save_batch_results_summary(self, all_results: Dict):
"""保存批量训练结果汇总"""
all_summary_data = []
for target_column_name, target_results in all_results.items():
for combo_key, result in target_results.items():
# 分离划分方法、预处理方法和建模方法
parts = combo_key.split('_', 2)
split_method = parts[0] if len(parts) > 0 else ''
preprocess_method = parts[1] if len(parts) > 1 else ''
model_method = parts[2] if len(parts) > 2 else ''
all_summary_data.append({
'目标列': target_column_name,
'划分方法': split_method,
'预处理方法': preprocess_method,
'建模方法': model_method,
'CV均值': result['cv_mean'],
'CV标准差': result['cv_std'],
'最佳得分': result['best_score'],
'训练集R²': result['train_r2'],
'训练集RMSE': result['train_rmse'],
'训练集MAE': result['train_mae'],
'训练集MSE': result['train_mse'],
'测试集R²': result['test_r2'],
'测试集RMSE': result['test_rmse'],
'测试集MAE': result['test_mae'],
'测试集MSE': result['test_mse'],
'训练样本数': result['train_size'],
'测试样本数': result['test_size'],
'最佳参数': str(result['best_params'])
})
if all_summary_data:
summary_df = pd.DataFrame(all_summary_data)
# 按目标列和测试集R²排序
summary_df = summary_df.sort_values(['目标列', '测试集R²'], ascending=[True, False])
# 保存详细结果CSV中文版
detailed_path = self.artifacts_dir / "batch_detailed_results.csv"
summary_df.to_csv(detailed_path, index=False, encoding='utf-8-sig')
# 保持原有的批量训练汇总结果(中文版)
batch_summary_path = self.artifacts_dir / "batch_training_summary.csv"
summary_df.to_csv(batch_summary_path, index=False, encoding='utf-8-sig')
# 创建简化版本用于兼容性(英文版)
all_summary_data_simple = []
for target_column_name, target_results in all_results.items():
for combo_key, result in target_results.items():
all_summary_data_simple.append({
'target_column': target_column_name,
'combination': combo_key,
'cv_mean': result['cv_mean'],
'cv_std': result['cv_std'],
'best_score': result['best_score'],
'train_r2': result['train_r2'],
'train_rmse': result['train_rmse'],
'train_mae': result['train_mae'],
'test_r2': result['test_r2'],
'test_rmse': result['test_rmse'],
'test_mae': result['test_mae'],
'train_size': result['train_size'],
'test_size': result['test_size'],
'split_method': result.get('split_method', 'unknown'),
'best_params': str(result['best_params'])
})
summary_df_simple = pd.DataFrame(all_summary_data_simple)
summary_df_simple = summary_df_simple.sort_values(['target_column', 'test_r2'], ascending=[True, False])
simple_summary_path = self.artifacts_dir / "batch_training_summary_simple.csv"
summary_df_simple.to_csv(simple_summary_path, index=False)
print(f"\n{'='*80}")
print("批量训练结果汇总:")
print(f"{'='*80}")
# 显示每个目标列的最佳模型
for target_col in summary_df['目标列'].unique():
target_data = summary_df[summary_df['目标列'] == target_col]
best_row = target_data.iloc[0] # 已经按R²降序排列
print(f"\n目标列 '{target_col}' 最佳模型:")
print(f" 组合: {best_row['划分方法']}_{best_row['预处理方法']}_{best_row['建模方法']}")
print(f" 测试集R²: {best_row['测试集R²']:.4f}")
print(f" 测试集RMSE: {best_row['测试集RMSE']:.4f}")
print(f" 最佳参数: {best_row['最佳参数']}")
print(f"\n详细结果已保存: {detailed_path}")
print(f"批量训练汇总结果已保存: {batch_summary_path}")
print(f"简化结果已保存: {simple_summary_path}")
# ★ 生成 MNF 特征筛选诊断报告
self._generate_mnf_report(all_results)
def _generate_mnf_report(self, all_results: Dict):
"""生成独立的 MNF 特征筛选诊断报告Markdown 格式)。
遍历所有训练结果,从 Pipeline 中提取 MNFTransformer 的诊断数据,
按目标指标 + 模型分块,每个成分按特征值降序列出具体数值。
Parameters
----------
all_results : Dict
train_models_batch() 返回的嵌套字典:
{target_column_name: {combo_key: result_dict}}
"""
report_path = self.artifacts_dir / "batch_mnf_eigenvalue_report.md"
lines = []
lines.append("# MNF 特征筛选诊断报告")
lines.append("")
import datetime as _dt
lines.append(f"> 生成时间: {_dt.datetime.now().strftime('%Y-%m-%d %H:%M:%S')}")
lines.append("")
lines.append("---")
lines.append("")
has_any_mnf = False
for target_name in sorted(all_results.keys()):
target_results = all_results[target_name]
for combo_key, result in target_results.items():
model = result.get('model')
if not isinstance(model, Pipeline):
continue
_mnf = self._find_mnf_in_pipeline(model)
if _mnf is None or not hasattr(_mnf, 'eigvals_w_'):
continue
has_any_mnf = True
# 解析组合键
parts = combo_key.split('_', 2)
split_method = parts[0] if len(parts) > 0 else ''
preprocess_method = parts[1] if len(parts) > 1 else ''
model_name = parts[2] if len(parts) > 2 else ''
diag = _mnf.get_diagnostics_dict(target_name=target_name)
n_sel = diag['selected_count']
cum_var = diag['cumulative_variance_ratio'] * 100
# ── 块标题 ──
lines.append(f"## 目标列: {target_name} | 模型: {model_name}")
lines.append("")
lines.append(f"- **划分方法**: {split_method}")
lines.append(f"- **预处理方法**: {preprocess_method}")
lines.append(f"- **总计选择数量**: {n_sel} 个主成分 "
f"(累计方差贡献: {cum_var:.2f}%)")
lines.append(f"- **选中成分详情 (按特征值降序排列)**:")
lines.append("")
# ── 成分列表(按特征值降序,逐个列出具体数值)──
for comp in diag['components_detail']:
lines.append(
f" {comp['rank']}. {comp['name']} "
f"(特征值: {comp['eigenvalue']:.6f}, "
f"方差占比: {comp['variance_ratio'] * 100:.2f}%)"
)
lines.append("")
lines.append("---")
lines.append("")
if not has_any_mnf:
lines.append("> ⚠ 未包含 MNF 诊断信息(所有模型均未使用 DualStream_MNF 预处理)")
lines.append("")
with open(str(report_path), 'w', encoding='utf-8') as f:
f.write('\n'.join(lines))
print(f"\n[MNF 报告] 已生成: {report_path}")
def load_model(self, preprocess_method: str, model_name: str):
"""
加载保存的模型
Args:
preprocess_method: 预处理方法名称
model_name: 模型名称
Returns:
加载的模型数据
"""
filename = f"{preprocess_method}_{model_name}.joblib"
filepath = self.artifacts_dir / filename
if not filepath.exists():
raise FileNotFoundError(f"模型文件不存在: {filepath}")
return joblib.load(filepath)
def get_best_model(self, metric: str = 'test_r2') -> Tuple[str, Dict]:
"""
获取最佳模型
Args:
metric: 评估指标默认使用测试集R²
可选:'test_r2', 'train_r2', 'test_rmse', 'test_mae',
'train_rmse', 'train_mae', 'cv_mean', 'best_score'
Returns:
最佳模型的组合名称和结果
"""
if not self.results:
raise ValueError("没有训练结果,请先训练模型")
# 对于回归指标R²和负MSE需要取最大值RMSE和MAE需要取最小值
if metric in ['test_r2', 'train_r2', 'cv_mean', 'best_score']:
best_combo = max(self.results.keys(),
key=lambda k: self.results[k][metric])
else: # rmse, mae等越小越好
best_combo = min(self.results.keys(),
key=lambda k: self.results[k][metric])
return best_combo, self.results[best_combo]
def main():
"""主函数示例 - 批量训练"""
# 创建批量建模实例
modeler = WaterQualityModelingBatch(r"D:\BaiduNetdiskDownload\yaobao\model")
# 批量训练多个目标列的模型
all_results = modeler.train_models_batch(
csv_path=r"D:\BaiduNetdiskDownload\yaobao\csv\yangdian_output.csv",
feature_start_column="374.285004", # 使用列名指定特征开始位置
preprocessing_methods=['None', 'MMS', 'SS', 'SNV', 'MA', 'SG', 'MSC', 'D1', 'D2', 'DT', 'CT'],#
model_names=['SVR', 'RF', 'Ridge', 'Lasso'],#, 'ElasticNet', 'XGBoost', 'LightGBM', 'CatBoost'
split_methods=['spxy', 'ks','random' ], #
cv_folds=5
)
print(f"\n批量训练完成,共训练了 {len(all_results)} 个目标列的模型")
# 显示每个目标列的最佳模型
for target_column_name, target_results in all_results.items():
if target_results:
best_combo = max(target_results.keys(),
key=lambda k: target_results[k]['test_r2'])
best_result = target_results[best_combo]
print(f"\n目标列 '{target_column_name}' 最佳模型:")
print(f" 组合: {best_combo}")
print(f" 测试集R²: {best_result['test_r2']:.4f}")
print(f" 测试集RMSE: {best_result['test_rmse']:.4f}")
if __name__ == "__main__":
main()