From b19cb427b98ee19b79ff866392e8214fa113f552 Mon Sep 17 00:00:00 2001 From: duxin Date: Thu, 30 Jul 2026 12:43:49 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20MNF=20=E8=AF=8A=E6=96=AD=E7=B3=BB?= =?UTF-8?q?=E7=BB=9F=20+=20=E5=85=A8=E5=B1=80=E5=AF=B9=E9=BD=90=E4=B8=A4?= =?UTF-8?q?=E9=81=8D=E6=89=AB=E6=8F=8F=E6=B3=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 【MNFTransformer 增强】(spectral_Preprocessing.py) - fit() 存储 eigvals_w_/eigvals_ratio_/eigvals_cumsum_ 实例属性 - _print_band_selection() 精简 3 行控制台输出 (总成分/选中数/Top3 贡献率) - get_diagnostics_dict() 返回结构化诊断字典供 JSON 导出 【模型保存增强】(modeling_batch.py) - save_model: joblib.dump 前注入 metadata['mnf_diagnostics'] - save_model: 并行导出 *_mnf_diag.json 独立诊断文件 - 部署 npz 新增 mnf_eigvals/_ratio/_cumsum 三个特征值数组 【批量训练报告】(modeling_batch.py) - _generate_mnf_report() 生成 batch_mnf_eigenvalue_report.md - 按目标列+模型分块,每个选中成分按特征值降序列出具体数值 - 容错:无 MNF 模型时标注警告,不崩溃 【MNF 全局对齐 — 两遍扫描法】(modeling_batch.py) - 第一遍预扫描: 对每个目标 fit MNFTransformer(0.95),记录所需成分数 - 取全局最大值 global_max_mnf_components (兜底 50) - 第二遍正式训练: 所有 Pipeline 统一使用固定整型成分数 - 仅当启用 DualStream_MNF 时触发,非 MNF 场景透明跳过 - get_preprocessing_transformer 新增 mnf_n_components 参数透传链 --- src/core/modeling/modeling_batch.py | 212 ++++++++++++++++++-- src/preprocessing/spectral_Preprocessing.py | 73 ++++++- 2 files changed, 262 insertions(+), 23 deletions(-) diff --git a/src/core/modeling/modeling_batch.py b/src/core/modeling/modeling_batch.py index 4e8cda9..fc1d917 100644 --- a/src/core/modeling/modeling_batch.py +++ b/src/core/modeling/modeling_batch.py @@ -576,7 +576,8 @@ class WaterQualityModelingBatch: 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: + preprocess_method: str = "None", + mnf_n_components=0.95) -> Dict: """ 训练单个回归模型(Pipeline 化:preprocess_method 字符串内部构造 Pipeline, scaler/MSC.mean_spectrum_ 等状态被绑定在 best_model 上,CV 与 test 评估 @@ -589,6 +590,8 @@ class WaterQualityModelingBatch: 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) 含义保持不变。 @@ -641,7 +644,8 @@ class WaterQualityModelingBatch: _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) + preproc = get_preprocessing_transformer(preprocess_method, wavelengths=_wl_list, + mnf_n_components=mnf_n_components) pipeline = Pipeline([ ('imputer', SimpleImputer(strategy='median')), ('preproc', preproc), @@ -737,6 +741,12 @@ class WaterQualityModelingBatch: 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, @@ -757,18 +767,41 @@ class WaterQualityModelingBatch: 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 or {} + 'metadata': metadata } joblib.dump(save_data, filepath) 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 纯量矩阵 # ═══════════════════════════════════════════════════════════ @@ -785,6 +818,11 @@ class WaterQualityModelingBatch: '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_) @@ -914,41 +952,98 @@ class WaterQualityModelingBatch: 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 内的逻辑一致) + _spec_cols = [c for c in X_clean.columns + if self._is_wavelength_column(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 + 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 @@ -958,13 +1053,14 @@ class WaterQualityModelingBatch: # 保存所有结果的汇总 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], + preprocessing_methods: List[str], model_names: List[str], split_methods: List[str], cv_folds: int, scoring: str, - test_size: float, random_state: int) -> Dict: + test_size: float, random_state: int, + mnf_n_components=0.95) -> Dict: """ 训练单个目标列的所有模型组合 """ @@ -985,6 +1081,7 @@ class WaterQualityModelingBatch: 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: @@ -1205,6 +1302,89 @@ class WaterQualityModelingBatch: 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): """ 加载保存的模型 diff --git a/src/preprocessing/spectral_Preprocessing.py b/src/preprocessing/spectral_Preprocessing.py index 078c4c9..7fc696f 100644 --- a/src/preprocessing/spectral_Preprocessing.py +++ b/src/preprocessing/spectral_Preprocessing.py @@ -409,19 +409,75 @@ class MNFTransformer(TransformerMixin, BaseEstimator, _ArrayAsFloat64): eigvecs_w = eigvecs_w[:, _order] Wp = eigvecs_w - # 5) 自动确定成分数 + # 5) 存储特征值(降序排列,用于波段选择诊断) + self.eigvals_w_ = np.maximum(eigvals_w, 0) + self.eigvals_ratio_ = self.eigvals_w_ / np.sum(self.eigvals_w_) + self.eigvals_cumsum_ = np.cumsum(self.eigvals_ratio_) + + # 6) 自动确定成分数 if isinstance(self.n_components, float) and 0 < self.n_components <= 1: - _cumvar = np.cumsum(np.maximum(eigvals_w, 0)) / np.sum(np.maximum(eigvals_w, 0)) - self.n_components_ = int(np.searchsorted(_cumvar, self.n_components)) + 1 + self.n_components_ = int(np.searchsorted(self.eigvals_cumsum_, self.n_components)) + 1 self.n_components_ = max(2, min(self.n_components_, n_features)) else: self.n_components_ = int(self.n_components) - # 6) 最终变换矩阵 + # 7) 最终变换矩阵 self.W_mnf_ = Wn @ Wp return self + def _print_band_selection(self): + """终端输出 MNF 波段选择诊断信息(精简 3 行格式)。""" + n_total = len(self.eigvals_w_) + n_sel = self.n_components_ + cum_var = self.eigvals_cumsum_[n_sel - 1] * 100 + + # Top 3 主成分贡献率 + top_parts = [] + for i in range(min(3, n_sel)): + top_parts.append(f"MNF{i + 1} ({self.eigvals_ratio_[i] * 100:.2f}%)") + top_str = ", ".join(top_parts) + + print(f"[MNF 诊断] 总成分: {n_total} | 选中: {n_sel} (前 {cum_var:.2f}% 方差贡献) | 舍弃: {n_total - n_sel}") + print(f"[MNF 诊断] Top 3 主成分贡献率: {top_str}") + print(f"[MNF 诊断] 详细数据请查看 *_mnf_diag.json 文件") + + def get_diagnostics_dict(self, target_name: str = "") -> dict: + """返回结构化 MNF 诊断字典,供 JSON 导出和 metadata 嵌入。 + + Parameters + ---------- + target_name : str + 目标指数名称(如 "WQI_Chla") + + Returns + ------- + dict + 包含 total_components / selected_count / cumulative_variance_ratio / + selected_bands / components_detail 等字段的结构化诊断数据。 + """ + n_total = len(self.eigvals_w_) + n_sel = self.n_components_ + + components_detail = [] + for i in range(n_sel): + components_detail.append({ + "rank": i + 1, + "name": f"MNF{i + 1}", + "eigenvalue": float(self.eigvals_w_[i]), + "variance_ratio": float(self.eigvals_ratio_[i]), + "cum_ratio": float(self.eigvals_cumsum_[i]), + }) + + return { + "target_name": target_name, + "total_components": n_total, + "selected_count": n_sel, + "cumulative_variance_ratio": float(self.eigvals_cumsum_[n_sel - 1]), + "selected_bands": [f"MNF{i + 1}" for i in range(n_sel)], + "components_detail": components_detail, + } + def transform(self, X): X = self._to_ndarray(X) X = np.nan_to_num(X, nan=0.0, posinf=0.0, neginf=0.0) @@ -531,19 +587,22 @@ _PREPROCESSING_TRANSFORMERS = { } -def get_preprocessing_transformer(method: str, wavelengths=None): +def get_preprocessing_transformer(method: str, wavelengths=None, + mnf_n_components=0.95): """根据预处理方法名返回 sklearn 兼容的 Transformer 实例。 - method 为 "None" 或 None:返回 IdentityTransformer - method 为 "MMS"/"SS":直接返回 sklearn 自带 MinMaxScaler/StandardScaler - method 为 "DualStream_MNF":返回 FeatureUnion,并行连接 - PhysicalFeatureExtractor(3 维物理指数) + MNFTransformer(10 维降维特征) + PhysicalFeatureExtractor(物理指数) + MNFTransformer(降维特征) - method 不识别:返回 IdentityTransformer + 打印警告 Args: method: 预处理方法名 wavelengths: DualStream_MNF 时需要传入波长列表(float), 供 PhysicalFeatureExtractor 定位波段列 + mnf_n_components: MNFTransformer 成分数;float(0~1)=自动按累计方差截断, + int=固定数量。DualStream_MNF 专用,其他 method 忽略。 Returns: sklearn 兼容的 Transformer 实例(可直接放入 Pipeline) @@ -554,7 +613,7 @@ def get_preprocessing_transformer(method: str, wavelengths=None): from sklearn.pipeline import FeatureUnion return FeatureUnion([ ('physical', PhysicalFeatureExtractor(wavelengths=wavelengths)), - ('mnf', MNFTransformer()), + ('mnf', MNFTransformer(n_components=mnf_n_components)), ]) if method == "Physical_Only": from sklearn.pipeline import Pipeline as _Pipeline