全局修正

This commit is contained in:
DXC
2026-06-29 16:16:55 +08:00
parent 2788fb3fe1
commit e337f01312
7 changed files with 437 additions and 104 deletions

View File

@ -19,6 +19,7 @@ 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 joblib import parallel_backend
# 第三方模型导入
# try:
@ -44,7 +45,7 @@ import os
is_frozen_env = getattr(sys, 'frozen', False)
safe_n_jobs = 1 if is_frozen_env else -1
from src.preprocessing.spectral_Preprocessing import Preprocessing
from src.preprocessing.spectral_Preprocessing import Preprocessing, get_preprocessing_transformer
from src.core.utils.split_methods import spxy, ks
@ -454,25 +455,28 @@ class WaterQualityModelingBatch:
else:
raise ValueError(f"不支持的划分方法: {method}. 支持的方法: {self.split_methods}")
def train_single_model(self, X: np.ndarray, y: pd.Series, model_name: str,
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") -> Dict:
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: 特征数据
X_raw: 原始特征数据(未经预处理,DataFrame 形态方便 Pipeline 内部转换)
y: 目标值数据
model_name: 模型名称
cv_folds: 交叉验证折数
scoring: 评分指标
test_size: 测试集比例
random_state: 随机种子
split_method: 数据划分方法
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}")
@ -483,18 +487,18 @@ class WaterQualityModelingBatch:
print(f"模型 {model_name} 不可用,请安装相应的库")
return None
print(f"开始训练模型: {model_name}")
print(f"开始训练模型: {model_name} (预处理: {preprocess_method})")
# 使用指定方法分割训练集和测试集
# 使用指定方法分割训练集和测试集(用原始 X_raw,Pipeline 内置 transform 处理)
X_train, X_test, y_train, y_test = self.split_data(
X, y, method=split_method, test_size=test_size, random_state=random_state
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:
@ -506,12 +510,25 @@ class WaterQualityModelingBatch:
elif model_name == 'LightGBM':
base_model.set_params(verbose=-1)
# 随机搜索 —— 替代穷举式 GridSearchCV,大幅降低寻优时间
# ============ 关键:把预处理器塞进 Pipeline ============
preproc = get_preprocessing_transformer(preprocess_method)
pipeline = Pipeline([
('preproc', preproc),
('model', base_model),
])
# RandomizedSearchCV 需要以「步骤名__参数名」的格式索引参数网格;
# 我们原有的 config['params'] 是模型层的(无 __),统一加 model__ 前缀。
prefixed_params = {
f"model__{k}": v for k, v in config['params'].items()
}
# 随机搜索:直接对 Pipeline 调优(scaler 仅在 train fold 上 fit)
cv_strategy = KFold(n_splits=cv_folds, shuffle=True, random_state=random_state)
grid_search = RandomizedSearchCV(
base_model,
config['params'],
pipeline,
prefixed_params,
n_iter=10,
cv=cv_strategy,
scoring=scoring,
@ -522,20 +539,22 @@ class WaterQualityModelingBatch:
grid_search.fit(X_train, y_train)
# 获取最佳模型
# 获取最佳模型(已是 Pipeline)
best_model = grid_search.best_estimator_
# 交叉验证评估(在训练集上)
cv_scores = cross_val_score(best_model, X_train, y_train, cv=cv_strategy, scoring=scoring)
# 交叉验证评估(在训练集上):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)
@ -562,7 +581,10 @@ class WaterQualityModelingBatch:
# 数据分割信息
'train_size': X_train.shape[0],
'test_size': X_test.shape[0],
'split_method': split_method
'split_method': split_method,
# Pipeline 信息(用于诊断 / metadata)
'preprocess_method': preprocess_method,
'is_pipeline': isinstance(best_model, Pipeline),
}
print(f"模型 {model_name} 训练完成:")
@ -714,21 +736,21 @@ class WaterQualityModelingBatch:
print(f"{'-' * 60}")
try:
# 数据预处理
X_processed = self.preprocess_data(X_raw, preprocess_method)
# 训练模型
result = self.train_single_model(X_processed, y, model_name,
cv_folds, scoring, test_size, random_state, split_method)
# 不再外部 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:
# 保存模型
# 保存模型(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_processed.shape,
'data_shape': X_raw.shape,
'target_range': [float(y.min()), float(y.max())],
'train_r2': result['train_r2'],
'train_rmse': result['train_rmse'],
@ -738,10 +760,13 @@ class WaterQualityModelingBatch:
'test_mae': result['test_mae'],
'train_size': result['train_size'],
'test_size': result['test_size'],
'split_method': result['split_method']
'split_method': result['split_method'],
# Pipeline 标记(便于旧 inference 路径兼容/诊断)
'preprocess_method': preprocess_method,
'is_pipeline': result.get('is_pipeline', False),
}
self.save_model(result['model'], target_column_name,
self.save_model(result['model'], target_column_name,
f"{split_method}_{preprocess_method}",
model_name, metadata)

View File

@ -38,6 +38,10 @@ from typing import Any, Callable, Dict, List, Optional, Tuple
import numpy as np
import pandas as pd
# sklearn Pipeline + 预处理 Transformer(避免 AutoML 训练时数据泄露)
from sklearn.pipeline import Pipeline
from src.preprocessing.spectral_Preprocessing import get_preprocessing_transformer
# ============================================================
# 常量
@ -189,8 +193,14 @@ def _get_search_space(model_name: str, trial) -> Dict[str, Any]:
def _make_objective(model_name: str, X: np.ndarray, y: np.ndarray,
cv_folds: int, random_state: int):
"""构造 Optuna objective(5 折 CV R²)。"""
cv_folds: int, random_state: int,
preproc_transformer=None):
"""构造 Optuna objective(5 折 CV R²)。
Pipeline 化:把 preproc_transformer 与 builder 装成 Pipeline,
cross_val_score 会 clone Pipeline 后每 fold 重 fit_transform,
从而避免「scaler 在 full X 上 fit → split → CV」造成的数据泄露。
"""
from sklearn.model_selection import KFold, cross_val_score
def objective(trial):
@ -199,9 +209,14 @@ def _make_objective(model_name: str, X: np.ndarray, y: np.ndarray,
builder = _build_model(model_name, random_state=random_state)
if builder is None:
return -1.0
model = builder(**params)
base = builder(**params)
pipe = Pipeline([
('preproc', preproc_transformer if preproc_transformer is not None
else get_preprocessing_transformer('None')),
('model', base),
])
kf = KFold(n_splits=cv_folds, shuffle=True, random_state=random_state)
scores = cross_val_score(model, X, y, cv=kf, scoring="r2", n_jobs=1)
scores = cross_val_score(pipe, X, y, cv=kf, scoring="r2", n_jobs=1)
return float(np.mean(scores))
except Exception:
return -1.0
@ -210,14 +225,20 @@ def _make_objective(model_name: str, X: np.ndarray, y: np.ndarray,
def _refit_full(model_name: str, best_params: Dict[str, Any],
X: np.ndarray, y: np.ndarray, random_state: int):
"""用 best params 在**全量数据**上 refit。"""
X: np.ndarray, y: np.ndarray, random_state: int,
preproc_transformer=None):
"""用 best params 在**全量数据**上 refit,保存为 Pipeline(含 scaler 等状态)。"""
builder = _build_model(model_name, random_state=random_state)
if builder is None:
return None
model = builder(**best_params)
model.fit(X, y)
return model
base = builder(**best_params)
pipe = Pipeline([
('preproc', preproc_transformer if preproc_transformer is not None
else get_preprocessing_transformer('None')),
('model', base),
])
pipe.fit(X, y)
return pipe
# ============================================================
@ -360,18 +381,10 @@ def train_with_automl(
feat_cols = [c for c in df.columns if c not in y_cols]
X_all = df[feat_cols].values.astype(np.float64)
# ---- 3) 预处理(仅第一项) ----
if preproc != "None":
try:
from src.preprocessing.spectral_Preprocessing import Preprocessing
processed = Preprocessing(preproc, df[feat_cols])
if isinstance(processed, pd.DataFrame):
X_all = processed.values.astype(np.float64)
else:
X_all = np.asarray(processed, dtype=np.float64)
except Exception as e:
notify("warning", f"预处理 {preproc} 失败: {e!r},改用 None")
preproc = "None"
# ---- 3) 预处理(Pipeline 化:不再手工 Preprocessing,把 transformer 透传给 _make_objective/_refit_full) ----
# 用一个 sklearn 兼容的 transformer 实例,后续每 trial 会被 clone 进 Pipeline.fit,
# 保证 scaler 只在 train fold 上 fit(无数据泄露)。
preproc_transformer = get_preprocessing_transformer(preproc)
# ---- 4) 检查 Optuna 是否可用 ----
try:
@ -427,7 +440,7 @@ def train_with_automl(
sampler=optuna.samplers.TPESampler(seed=random_state),
)
study.optimize(
_make_objective(model_name, X_sub, y_sub, cv_folds, random_state),
_make_objective(model_name, X_sub, y_sub, cv_folds, random_state, preproc_transformer),
n_trials=n_trials,
timeout=per_model_timeout,
show_progress_bar=False,
@ -437,8 +450,8 @@ def train_with_automl(
notify("warning", f"{tgt}/{model_name}: 全部 trial 失败(CV 全部 <= -1)")
continue
# refit on FULL
final_model = _refit_full(model_name, study.best_params, X_t, y_t, random_state)
# refit on FULL(Pipeline 化:scaler 用全量 X_t 拟合一次)
final_model = _refit_full(model_name, study.best_params, X_t, y_t, random_state, preproc_transformer)
if final_model is None:
continue
@ -453,6 +466,7 @@ def train_with_automl(
"model_name": model_name,
"metadata": {
"automl": True,
"is_pipeline": isinstance(final_model, Pipeline),
"best_params": study.best_params,
"cv_score": float(study.best_value),
"n_trials_done": len(study.trials),

View File

@ -12,7 +12,7 @@ warnings.filterwarnings('ignore')
import sys
import os
from src.preprocessing.spectral_Preprocessing import Preprocessing
from src.preprocessing.spectral_Preprocessing import Preprocessing, get_preprocessing_transformer
from src.core.utils.split_methods import spxy, ks
# try:
@ -22,6 +22,7 @@ from src.core.utils.split_methods import spxy, ks
# 机器学习相关导入
from sklearn.model_selection import train_test_split
from sklearn.pipeline import Pipeline
class WaterQualityInference:
@ -423,33 +424,20 @@ class WaterQualityInference:
from src.utils.water_index import WaterQualityIndexCalculator
calc = WaterQualityIndexCalculator()
# 提取纯计算方法(排除 find_closest_wavelength 和 calculate_all_indices,
# 以及不返回 Series 的辅助方法)
algorithm_methods = []
for m in dir(calc):
if m.startswith('_'):
continue
if m in ['find_closest_wavelength', 'calculate_all_indices']:
continue
attr = getattr(calc, m)
if callable(attr):
algorithm_methods.append(m)
original_col_count = spectra.shape[1]
for algo_name in algorithm_methods:
try:
algo_func = getattr(calc, algo_name)
result = algo_func(spectra)
# 只追加返回 Series 且长度为样本数的合法结果
if isinstance(result, pd.Series) and len(result) == len(spectra):
spectra[algo_name] = result.values
else:
spectra[algo_name] = np.nan
except Exception:
spectra[algo_name] = np.nan
print(f"[特征补全] 完成!光谱列已扩充至 {spectra.shape[1]} 列"
f"(追加了 {spectra.shape[1] - original_col_count} 个 WQI 指数)")
# CSV 驱动的 WaterQualityIndexCalculator:所有公式名通过 list_available() 拿;
# 一次性 calculate_many() 批量计算。彻底摆脱 dir(calc) 反射扫描 + 单 algo_func
# 调用这种碎片化写法(Calculator 早已重构为公式驱动,不再有独立公式方法)。
formulas = calc.list_available()
if not formulas:
print("[特征补全] Calculator 未持有任何公式,跳过补全")
else:
results_df = calc.calculate_many(formulas, spectra)
# results_df 是列对齐的 WQI 计算结果(每列一个公式,行数=样本数)
if isinstance(results_df, pd.DataFrame) and not results_df.empty:
original_col_count = spectra.shape[1]
spectra = pd.concat([spectra, results_df], axis=1)
print(f"[特征补全] 完成!光谱列已扩充至 {spectra.shape[1]} 列"
f"(追加了 {spectra.shape[1] - original_col_count} 个 WQI 指数)")
except Exception as e:
print(f"[特征补全] 失败,将使用原始光谱特征: {e}")
@ -471,6 +459,14 @@ class WaterQualityInference:
print(f"[特征对齐] 最终输入维度: {spectra.shape}")
# ---- Pipeline 化分支:模型内置 scaler/MSC.mean_spectrum_ 等状态时,跳过手动 Preprocessing ----
if isinstance(model, Pipeline):
print(f"[Pipeline] 检测到模型是 sklearn Pipeline,"
f"其内置预处理步骤({list(model.named_steps.keys())[0]})将处理原始光谱,"
f"无需外部 Preprocessing")
return spectra.values
# ---- 兼容路径:旧 .joblib(裸模型 + preprocess_method 字符串)回退手动 Preprocessing ----
try:
# 应用预处理
spectra_processed = Preprocessing(actual_preprocess_method, spectra)
@ -479,7 +475,8 @@ class WaterQualityInference:
if isinstance(spectra_processed, pd.DataFrame):
spectra_processed = spectra_processed.values
print(f"预处理后数据形状: {spectra_processed.shape}")
print(f" [Legacy] 旧裸模型 + 手动 Preprocessing({actual_preprocess_method}) 完成,"
f"数据形状: {spectra_processed.shape}")
return spectra_processed

View File

@ -24,6 +24,7 @@ from sklearn.linear_model import LinearRegression, Ridge, Lasso, ElasticNet
from sklearn.model_selection import GridSearchCV, 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.pipeline import Pipeline
from src.core.utils.split_methods import spxy, ks
# 第三方模型导入
@ -46,7 +47,7 @@ CB_AVAILABLE = False # 注释掉catboost
import sys
import os
from src.preprocessing.spectral_Preprocessing import Preprocessing
from src.preprocessing.spectral_Preprocessing import Preprocessing, get_preprocessing_transformer
class WaterQualityScatterBatch:
@ -626,12 +627,18 @@ class WaterQualityScatterBatch:
best_model_data = self.load_model(artifacts_path, model_file_prefix, model_name, folder_name)
best_model = best_model_data['model']
# 应用相同的数据预处理
X_processed = self.preprocess_data(X_raw, actual_preprocess_method)
# 应用相同的数据预处理:Pipeline 模型自带 scaler,直接喂 raw;
# 旧模型(裸模型)才需要手动 Preprocessing
if isinstance(best_model, Pipeline):
X_pred_input = X_raw
print(f" [Pipeline] 模型自带预处理器,跳过外部 Preprocessing")
else:
X_pred_input = self.preprocess_data(X_raw, actual_preprocess_method)
print(f" [Legacy] 旧裸模型,回退到手动 Preprocessing({actual_preprocess_method})")
# 使用相同的数据分割方法
X_train, X_test, y_train, y_test = self.split_data(
X_processed, y_true, method=split_method,
X_pred_input, y_true, method=split_method,
test_size=test_size, random_state=random_state
)