feat: MNF 诊断系统 + 全局对齐两遍扫描法

【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 参数透传链
This commit is contained in:
duxin
2026-07-30 12:43:49 +08:00
parent 314d947d34
commit b19cb427b9
2 changed files with 262 additions and 23 deletions

View File

@ -576,7 +576,8 @@ class WaterQualityModelingBatch:
cv_folds: int = 5, scoring: str = 'neg_mean_squared_error', cv_folds: int = 5, scoring: str = 'neg_mean_squared_error',
test_size: float = 0.2, random_state: int = 42, test_size: float = 0.2, random_state: int = 42,
split_method: str = "random", split_method: str = "random",
preprocess_method: str = "None") -> Dict: preprocess_method: str = "None",
mnf_n_components=0.95) -> Dict:
""" """
训练单个回归模型Pipeline 化preprocess_method 字符串内部构造 Pipeline 训练单个回归模型Pipeline 化preprocess_method 字符串内部构造 Pipeline
scaler/MSC.mean_spectrum_ 等状态被绑定在 best_model 上CV 与 test 评估 scaler/MSC.mean_spectrum_ 等状态被绑定在 best_model 上CV 与 test 评估
@ -589,6 +590,8 @@ class WaterQualityModelingBatch:
model_name: 模型名称 model_name: 模型名称
preprocess_method: 预处理方法字符串(如 'None' / 'SS' / 'MSC' preprocess_method: 预处理方法字符串(如 'None' / 'SS' / 'MSC'
训练结束后 best_model 字段即为 sklearn Pipeline。 训练结束后 best_model 字段即为 sklearn Pipeline。
mnf_n_components: MNFTransformer 成分数(仅 DualStream_MNF 使用);
float(0~1)=自动按累计方差截断int=固定数量。
其余参数cv_folds / scoring / test_size / random_state / split_method 其余参数cv_folds / scoring / test_size / random_state / split_method
含义保持不变。 含义保持不变。
@ -641,7 +644,8 @@ class WaterQualityModelingBatch:
_wl_list = None _wl_list = None
if preprocess_method in ("DualStream_MNF", "Physical_Only"): if preprocess_method in ("DualStream_MNF", "Physical_Only"):
_wl_list = self._extract_train_wavelengths(X_raw.columns) _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([ pipeline = Pipeline([
('imputer', SimpleImputer(strategy='median')), ('imputer', SimpleImputer(strategy='median')),
('preproc', preproc), ('preproc', preproc),
@ -737,6 +741,12 @@ class WaterQualityModelingBatch:
print(f" RMSE: {result['test_rmse']:.4f}") print(f" RMSE: {result['test_rmse']:.4f}")
print(f" MAE: {result['test_mae']:.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 return result
def save_model(self, model, target_column_name: str, preprocess_method: str, model_name: str, 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" filename = f"{safe_target_name}_{preprocess_method}_{model_name}.joblib"
filepath = self.artifacts_dir / filename 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 = { save_data = {
'model': model, 'model': model,
'target_column_name': target_column_name, 'target_column_name': target_column_name,
'preprocess_method': preprocess_method, 'preprocess_method': preprocess_method,
'model_name': model_name, 'model_name': model_name,
'metadata': metadata or {} 'metadata': metadata
} }
joblib.dump(save_data, filepath) joblib.dump(save_data, filepath)
print(f"模型已保存: {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 纯量矩阵 # ★ C++/Rust 部署导出:提取 MNF + SVR 纯量矩阵
# ═══════════════════════════════════════════════════════════ # ═══════════════════════════════════════════════════════════
@ -785,6 +818,11 @@ class WaterQualityModelingBatch:
'mnf_W': np.asarray(_mnf.W_mnf_), 'mnf_W': np.asarray(_mnf.W_mnf_),
'mnf_n_components': int(getattr(_mnf, 'n_components_', _mnf.n_components)), '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_ # SVR 部署矩阵rbf kernel 需要 support_vectors_
_deploy['svr_dual_coef'] = np.asarray(_svr.dual_coef_) _deploy['svr_dual_coef'] = np.asarray(_svr.dual_coef_)
_deploy['svr_intercept'] = np.asarray(_svr.intercept_) _deploy['svr_intercept'] = np.asarray(_svr.intercept_)
@ -914,41 +952,98 @@ class WaterQualityModelingBatch:
all_results = {} 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(): for target_column_name, y in y_dict.items():
print(f"\n{'='*80}") print(f"\n{'='*80}")
print(f"开始训练目标列: {target_column_name}") print(f"开始训练目标列: {target_column_name}")
if _has_dual_stream:
print(f"[MNF 全局对齐] 统一使用 {global_max_mnf_components} 个 MNF 成分")
print(f"{'='*80}") print(f"{'='*80}")
# 创建该目标列的子目录 # 创建该目标列的子目录
target_artifacts_dir = self.artifacts_dir / target_column_name target_artifacts_dir = self.artifacts_dir / target_column_name
target_artifacts_dir.mkdir(parents=True, exist_ok=True) target_artifacts_dir.mkdir(parents=True, exist_ok=True)
# 临时更改artifacts_dir # 临时更改artifacts_dir
original_artifacts_dir = self.artifacts_dir original_artifacts_dir = self.artifacts_dir
self.artifacts_dir = target_artifacts_dir self.artifacts_dir = target_artifacts_dir
try: try:
# 去除该目标列的空值 # 去除该目标列的空值
mask = ~y.isna() mask = ~y.isna()
if mask.sum() == 0: if mask.sum() == 0:
print(f"目标列 '{target_column_name}' 无有效数据,跳过") print(f"目标列 '{target_column_name}' 无有效数据,跳过")
continue continue
X_clean = X_raw[mask] X_clean = X_raw[mask]
y_clean = y[mask] y_clean = y[mask]
print(f"有效样本数: {len(y_clean)}") print(f"有效样本数: {len(y_clean)}")
# 训练该目标列的所有模型组合 # 训练该目标列的所有模型组合
target_results = self.train_models_single_target( target_results = self.train_models_single_target(
X_clean, y_clean, target_column_name, X_clean, y_clean, target_column_name,
preprocessing_methods, model_names, split_methods, 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 all_results[target_column_name] = target_results
except Exception as e: except Exception as e:
print(f"训练目标列 '{target_column_name}' 时出错: {e}") print(f"训练目标列 '{target_column_name}' 时出错: {e}")
continue continue
@ -958,13 +1053,14 @@ class WaterQualityModelingBatch:
# 保存所有结果的汇总 # 保存所有结果的汇总
self._save_batch_results_summary(all_results) self._save_batch_results_summary(all_results)
return all_results return all_results
def train_models_single_target(self, X_raw: pd.DataFrame, y: pd.Series, target_column_name: str, 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, 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, X_raw, y, model_name,
cv_folds, scoring, test_size, random_state, split_method, cv_folds, scoring, test_size, random_state, split_method,
preprocess_method=preprocess_method, preprocess_method=preprocess_method,
mnf_n_components=mnf_n_components,
) )
if result is not None: if result is not None:
@ -1205,6 +1302,89 @@ class WaterQualityModelingBatch:
print(f"批量训练汇总结果已保存: {batch_summary_path}") print(f"批量训练汇总结果已保存: {batch_summary_path}")
print(f"简化结果已保存: {simple_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): def load_model(self, preprocess_method: str, model_name: str):
""" """
加载保存的模型 加载保存的模型

View File

@ -409,19 +409,75 @@ class MNFTransformer(TransformerMixin, BaseEstimator, _ArrayAsFloat64):
eigvecs_w = eigvecs_w[:, _order] eigvecs_w = eigvecs_w[:, _order]
Wp = eigvecs_w 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: 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(self.eigvals_cumsum_, self.n_components)) + 1
self.n_components_ = int(np.searchsorted(_cumvar, self.n_components)) + 1
self.n_components_ = max(2, min(self.n_components_, n_features)) self.n_components_ = max(2, min(self.n_components_, n_features))
else: else:
self.n_components_ = int(self.n_components) self.n_components_ = int(self.n_components)
# 6) 最终变换矩阵 # 7) 最终变换矩阵
self.W_mnf_ = Wn @ Wp self.W_mnf_ = Wn @ Wp
return self 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): def transform(self, X):
X = self._to_ndarray(X) X = self._to_ndarray(X)
X = np.nan_to_num(X, nan=0.0, posinf=0.0, neginf=0.0) 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 实例。 """根据预处理方法名返回 sklearn 兼容的 Transformer 实例。
- method 为 "None" 或 None返回 IdentityTransformer - method 为 "None" 或 None返回 IdentityTransformer
- method 为 "MMS"/"SS":直接返回 sklearn 自带 MinMaxScaler/StandardScaler - method 为 "MMS"/"SS":直接返回 sklearn 自带 MinMaxScaler/StandardScaler
- method 为 "DualStream_MNF":返回 FeatureUnion并行连接 - method 为 "DualStream_MNF":返回 FeatureUnion并行连接
PhysicalFeatureExtractor(3 维物理指数) + MNFTransformer(10 维降维特征) PhysicalFeatureExtractor(物理指数) + MNFTransformer(降维特征)
- method 不识别:返回 IdentityTransformer + 打印警告 - method 不识别:返回 IdentityTransformer + 打印警告
Args: Args:
method: 预处理方法名 method: 预处理方法名
wavelengths: DualStream_MNF 时需要传入波长列表(float) wavelengths: DualStream_MNF 时需要传入波长列表(float)
供 PhysicalFeatureExtractor 定位波段列 供 PhysicalFeatureExtractor 定位波段列
mnf_n_components: MNFTransformer 成分数float(0~1)=自动按累计方差截断,
int=固定数量。DualStream_MNF 专用,其他 method 忽略。
Returns: Returns:
sklearn 兼容的 Transformer 实例(可直接放入 Pipeline sklearn 兼容的 Transformer 实例(可直接放入 Pipeline
@ -554,7 +613,7 @@ def get_preprocessing_transformer(method: str, wavelengths=None):
from sklearn.pipeline import FeatureUnion from sklearn.pipeline import FeatureUnion
return FeatureUnion([ return FeatureUnion([
('physical', PhysicalFeatureExtractor(wavelengths=wavelengths)), ('physical', PhysicalFeatureExtractor(wavelengths=wavelengths)),
('mnf', MNFTransformer()), ('mnf', MNFTransformer(n_components=mnf_n_components)),
]) ])
if method == "Physical_Only": if method == "Physical_Only":
from sklearn.pipeline import Pipeline as _Pipeline from sklearn.pipeline import Pipeline as _Pipeline