Files
WQ_GUI/src/core/modeling/modeling_batch.py
2026-07-29 15:46:14 +08:00

1292 lines
54 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.preprocessing import StandardScaler
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 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 = {}
@staticmethod
def _is_wavelength_column(col_name: str) -> bool:
"""判断列名是否为波长值(纯数字字符串,如 '374.285004'"""
try:
float(str(col_name))
return True
except (ValueError, TypeError):
return False
@staticmethod
def _is_wqi_column(col_name: str) -> bool:
"""判断列名是否为 WQI 水质指数列('WQI_' 前缀)"""
return str(col_name).startswith('WQI_')
@staticmethod
def _extract_train_wavelengths(columns) -> List[float]:
"""从列名列表中提取波长值float 列表)
遍历列名,将所有可转为 float 的列名提取为波长列表。
用于写入模型 metadata['train_wavelengths'],供推理端光谱重采样。
"""
wl_list = []
for c in columns:
try:
wl_list.append(float(str(c)))
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
# 保留波长列
if self._is_wavelength_column(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") -> 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。
其余参数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_column(c)]
X_raw = X_raw[_spec_cols]
print(f"[{preprocess_method}] 精简为纯光谱: {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)
# ============ 关键:把预处理器塞进 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)
pipeline = Pipeline([
('imputer', SimpleImputer(strategy='median')),
('preproc', preproc),
('cleaner', _SafeFiniteTransformer()),
('scaler', StandardScaler()),
('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__ 前缀。
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}")
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
# 保存模型和元数据
save_data = {
'model': model,
'target_column_name': target_column_name,
'preprocess_method': preprocess_method,
'model_name': model_name,
'metadata': metadata or {}
}
joblib.dump(save_data, filepath)
print(f"模型已保存: {filepath}")
# ═══════════════════════════════════════════════════════════
# ★ 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)),
}
# StandardScaler 参数C++ 端需复现 (X - mean) / scale
_scaler = model.named_steps.get('scaler')
if _scaler is not None and hasattr(_scaler, 'mean_'):
_deploy['scaler_mean'] = np.asarray(_scaler.mean_, dtype=np.float64)
_deploy['scaler_scale'] = np.asarray(_scaler.scale_, dtype=np.float64)
# 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 not None and model_step.__class__.__name__ == 'SVR':
return model_step
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 = {}
# 对每个目标列进行训练
for target_column_name, y in y_dict.items():
print(f"\n{'='*80}")
print(f"开始训练目标列: {target_column_name}")
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
)
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) -> 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,
)
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}")
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()