"""Exploratory candidate-ranking benchmark for ENGIN fault severity. EXPERIMENT STATUS: CANDIDATE SEARCH. Absolute end-to-end metrics from this script do not describe ``final_pipeline.py`` because its frozen label model is an older baseline. Its reusable transformers and winning severity candidate are imported by the final pipeline. Use ``ml_polish_benchmark.py`` for final end-to-end numbers. The label model is frozen as combined spectral features + Logistic Regression. Severity candidates compare raw spectra, engine-relative deviation features, ordinal models, and fault-specific models. Every experiment uses complete held- out engines and is evaluated on clean validation spectra and the same spectra with 5% of measurements masked. """ from __future__ import annotations import argparse import json from dataclasses import dataclass from pathlib import Path from typing import Iterable import numpy as np import pandas as pd from sklearn.base import BaseEstimator from sklearn.dummy import DummyClassifier from sklearn.ensemble import ExtraTreesClassifier from sklearn.linear_model import LogisticRegression, Ridge from sklearn.metrics import accuracy_score from sklearn.pipeline import make_pipeline as sklearn_pipeline from sklearn.preprocessing import StandardScaler from sklearn.svm import SVC from engin.features import FEATURE_SETS, SeverityFeatures from benchmark_grouped import ( FAULT_LABELS, FREQ_COLS, NOT_APPLICABLE, SEVERITIES, make_pipeline, make_splits, macro_f1, ml_points, raw_score, validate_data, ) from robustness_grouped import mask_spectral_cells, missing_scenario_name DEFAULT_SEEDS = [7, 21, 42, 77, 123] SEVERITY_TO_INT = {"male": 0, "srednie": 1, "duze": 2} INT_TO_SEVERITY = np.asarray(SEVERITIES, dtype=object) @dataclass(frozen=True) class Candidate: candidate_id: str feature_set: str estimator: str include_fault_type: bool = False per_fault: bool = False CANDIDATES = [ Candidate("raw_extra_trees", "raw", "extra_trees"), Candidate("deviation_extra_trees", "deviation", "extra_trees"), Candidate( "deviation_extra_trees_mf03", "deviation", "extra_trees_mf03", ), Candidate("all_extra_trees", "all", "extra_trees", include_fault_type=True), Candidate( "all_extra_trees_leaf2", "all", "extra_trees_leaf2", include_fault_type=True, ), Candidate("all_logistic", "all", "logistic", include_fault_type=True), Candidate("all_svc", "all", "svc", include_fault_type=True), Candidate("all_ordinal_ridge", "all", "ordinal_ridge", include_fault_type=True), Candidate( "all_ordinal_logistic", "all", "ordinal_logistic", include_fault_type=True, ), Candidate("all_extra_trees_per_fault", "all", "extra_trees", per_fault=True), ] CANDIDATE_BY_ID = {candidate.candidate_id: candidate for candidate in CANDIDATES} class OrdinalRidge(BaseEstimator): """Weighted ridge regression rounded to the three ordered severities.""" def __init__(self, alpha: float = 10.0) -> None: self.alpha = alpha def fit(self, X: np.ndarray, y: Iterable[str]) -> "OrdinalRidge": y_array = np.asarray(list(y), dtype=object) y_numeric = np.asarray([SEVERITY_TO_INT[value] for value in y_array], dtype=float) self.scaler_ = StandardScaler().fit(X) counts = pd.Series(y_array).value_counts() weights = np.asarray([len(y_array) / counts[value] for value in y_array]) self.model_ = Ridge(alpha=self.alpha).fit( self.scaler_.transform(X), y_numeric, sample_weight=weights ) return self def predict(self, X: np.ndarray) -> np.ndarray: numeric = self.model_.predict(self.scaler_.transform(X)) levels = np.clip(np.rint(numeric), 0, 2).astype(int) return INT_TO_SEVERITY[levels] class OrdinalLogistic(BaseEstimator): """Two cumulative balanced logistic classifiers for ordered severity.""" def __init__(self, C: float = 1.0, random_state: int = 42) -> None: self.C = C self.random_state = random_state def fit(self, X: np.ndarray, y: Iterable[str]) -> "OrdinalLogistic": y_numeric = np.asarray([SEVERITY_TO_INT[value] for value in y], dtype=int) self.scaler_ = StandardScaler().fit(X) scaled = self.scaler_.transform(X) self.threshold_models_: list[float | LogisticRegression] = [] for threshold in (0, 1): target = (y_numeric > threshold).astype(int) if np.unique(target).size == 1: self.threshold_models_.append(float(target[0])) else: model = LogisticRegression( C=self.C, class_weight="balanced", max_iter=5_000, random_state=self.random_state + threshold, ).fit(scaled, target) self.threshold_models_.append(model) return self def predict(self, X: np.ndarray) -> np.ndarray: scaled = self.scaler_.transform(X) probabilities = [] for model in self.threshold_models_: if isinstance(model, float): probabilities.append(np.full(len(X), model)) else: probabilities.append(model.predict_proba(scaled)[:, 1]) above_male = probabilities[0] above_srednie = np.minimum(probabilities[1], above_male) levels = (above_male >= 0.5).astype(int) + (above_srednie >= 0.5).astype(int) return INT_TO_SEVERITY[levels] class PerFaultEstimator: """Fit a separate severity estimator for each predicted fault type.""" def __init__(self, estimator_name: str, random_state: int, n_jobs: int) -> None: self.estimator_name = estimator_name self.random_state = random_state self.n_jobs = n_jobs def fit( self, X: np.ndarray, y: np.ndarray, fault_labels: np.ndarray, ) -> "PerFaultEstimator": self.global_model_ = make_estimator( self.estimator_name, self.random_state, self.n_jobs ).fit(X, y) self.models_: dict[str, object] = {} for offset, fault_label in enumerate(FAULT_LABELS, start=1): mask = fault_labels == fault_label if not mask.any(): continue if np.unique(y[mask]).size == 1: model = DummyClassifier(strategy="most_frequent").fit(X[mask], y[mask]) else: model = make_estimator( self.estimator_name, self.random_state + offset, self.n_jobs, ).fit(X[mask], y[mask]) self.models_[fault_label] = model return self def predict(self, X: np.ndarray, fault_labels: np.ndarray) -> np.ndarray: result = self.global_model_.predict(X).astype(object) for fault_label, model in self.models_.items(): mask = fault_labels == fault_label if mask.any(): result[mask] = model.predict(X[mask]) return result def make_estimator(estimator_name: str, random_state: int, n_jobs: int): if estimator_name == "extra_trees": return ExtraTreesClassifier( n_estimators=400, class_weight="balanced", max_features="sqrt", min_samples_leaf=1, n_jobs=n_jobs, random_state=random_state, ) if estimator_name == "extra_trees_leaf2": return ExtraTreesClassifier( n_estimators=400, class_weight="balanced", max_features=0.7, min_samples_leaf=2, n_jobs=n_jobs, random_state=random_state, ) if estimator_name == "extra_trees_mf03": return ExtraTreesClassifier( n_estimators=800, class_weight="balanced", max_features=0.3, min_samples_leaf=1, n_jobs=n_jobs, random_state=random_state, ) if estimator_name == "logistic": return sklearn_pipeline( StandardScaler(), LogisticRegression( C=1.0, class_weight="balanced", max_iter=5_000, random_state=random_state, ), ) if estimator_name == "svc": return sklearn_pipeline( StandardScaler(), SVC(C=2.0, kernel="rbf", class_weight="balanced"), ) if estimator_name == "ordinal_ridge": return OrdinalRidge(alpha=10.0) if estimator_name == "ordinal_logistic": return OrdinalLogistic(C=1.0, random_state=random_state) raise ValueError(f"Unknown estimator: {estimator_name}") def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description=__doc__) parser.add_argument("--data", type=Path, default=Path("val.csv")) parser.add_argument("--output-dir", type=Path, default=Path("severity_outputs")) parser.add_argument("--seeds", nargs="+", type=int, default=DEFAULT_SEEDS) parser.add_argument("--n-splits", type=int, default=5) parser.add_argument("--missing-rate", type=float, default=0.05) parser.add_argument("--mask-random-state", type=int, default=2026) parser.add_argument("--model-random-state", type=int, default=42) parser.add_argument("--n-jobs", type=int, default=-1) parser.add_argument( "--candidates", nargs="+", choices=sorted(CANDIDATE_BY_ID), default=[candidate.candidate_id for candidate in CANDIDATES], ) return parser.parse_args() def prepare_labeled_frame(df: pd.DataFrame, labels: np.ndarray) -> pd.DataFrame: result = df.copy() result["diagnostic_label"] = labels return result def fit_candidate( candidate: Candidate, X_train_full: pd.DataFrame, train_fault_mask: np.ndarray, y_severity_fault: np.ndarray, y_fault_label: np.ndarray, random_state: int, n_jobs: int, ): transformer = SeverityFeatures( candidate.feature_set, include_fault_type=candidate.include_fault_type ).fit(X_train_full) train_features = transformer.transform(X_train_full)[train_fault_mask] if candidate.per_fault: estimator = PerFaultEstimator(candidate.estimator, random_state, n_jobs).fit( train_features, y_severity_fault, y_fault_label ) else: estimator = make_estimator(candidate.estimator, random_state, n_jobs).fit( train_features, y_severity_fault ) return transformer, estimator def predict_candidate( candidate: Candidate, transformer: SeverityFeatures, estimator, X_valid_full: pd.DataFrame, predicted_labels: np.ndarray, ) -> np.ndarray: features = transformer.transform(X_valid_full) if candidate.per_fault: return estimator.predict(features, predicted_labels) return estimator.predict(features) def main() -> None: args = parse_args() df_clean = pd.read_csv(args.data).reset_index(drop=True) validate_data(df_clean, args.n_splits) df_masked, mask_manifest = mask_spectral_cells( df_clean, args.missing_rate, args.mask_random_state ) masked_scenario = missing_scenario_name(args.missing_rate) scenarios = {"clean": df_clean, masked_scenario: df_masked} candidates = [CANDIDATE_BY_ID[candidate_id] for candidate_id in args.candidates] args.output_dir.mkdir(parents=True, exist_ok=True) mask_manifest.to_csv(args.output_dir / "mask_manifest.csv", index=False) y_label = df_clean["label"].reset_index(drop=True) y_severity = df_clean["severity"].reset_index(drop=True) true_fault = y_label.isin(FAULT_LABELS).to_numpy() all_run_rows: list[dict[str, object]] = [] all_fold_rows: list[dict[str, object]] = [] for cv_seed in args.seeds: print(f"Running severity CV seed {cv_seed} ...", flush=True) splits = make_splits( df_clean, splitter_name="stratified-group", n_splits=args.n_splits, random_state=cv_seed, ) oof_label = { scenario: np.full(len(df_clean), "", dtype=object) for scenario in scenarios } oof_severity = { (candidate.candidate_id, scenario): np.full( len(df_clean), NOT_APPLICABLE, dtype=object ) for candidate in candidates for scenario in scenarios } for fold, (train_idx, valid_idx) in enumerate(splits, start=1): X_train_clean = df_clean.iloc[train_idx] y_train_label = y_label.iloc[train_idx] label_pipeline = make_pipeline( "logistic", "combined", args.model_random_state + fold, args.n_jobs ).fit(X_train_clean, y_train_label) scenario_predicted_labels: dict[str, np.ndarray] = {} for scenario, scenario_df in scenarios.items(): predicted = label_pipeline.predict(scenario_df.iloc[valid_idx]) scenario_predicted_labels[scenario] = predicted oof_label[scenario][valid_idx] = predicted train_fault_mask = y_train_label.isin(FAULT_LABELS).to_numpy() train_with_labels = prepare_labeled_frame( X_train_clean, y_train_label.to_numpy(dtype=object) ) y_severity_fault = ( y_severity.iloc[train_idx].to_numpy(dtype=object)[train_fault_mask] ) y_fault_label = y_train_label.to_numpy(dtype=object)[train_fault_mask] fold_true_fault = y_label.iloc[valid_idx].isin(FAULT_LABELS).to_numpy() for candidate_offset, candidate in enumerate(candidates): transformer, estimator = fit_candidate( candidate=candidate, X_train_full=train_with_labels, train_fault_mask=train_fault_mask, y_severity_fault=y_severity_fault, y_fault_label=y_fault_label, random_state=args.model_random_state + 1000 + 20 * fold + candidate_offset, n_jobs=args.n_jobs, ) for scenario, scenario_df in scenarios.items(): predicted_labels = scenario_predicted_labels[scenario] valid_with_labels = prepare_labeled_frame( scenario_df.iloc[valid_idx], predicted_labels ) severity_pred = predict_candidate( candidate, transformer, estimator, valid_with_labels, predicted_labels, ) emitted = np.full(len(valid_idx), NOT_APPLICABLE, dtype=object) predicted_fault = np.isin(predicted_labels, FAULT_LABELS) emitted[predicted_fault] = severity_pred[predicted_fault] oof_severity[(candidate.candidate_id, scenario)][valid_idx] = emitted fold_label_f1 = macro_f1( y_label.iloc[valid_idx], predicted_labels ) fold_severity = float( accuracy_score( y_severity.iloc[valid_idx].to_numpy()[fold_true_fault], emitted[fold_true_fault], ) ) fold_score = raw_score(fold_label_f1, fold_severity) all_fold_rows.append( { "cv_seed": cv_seed, "fold": fold, "scenario": scenario, "candidate_id": candidate.candidate_id, "valid_fault_rows": int(fold_true_fault.sum()), "macro_f1": fold_label_f1, "severity_accuracy_submission": fold_severity, "raw_score": fold_score, "ml_points": ml_points(fold_score), } ) for scenario in scenarios: label_f1 = macro_f1(y_label, oof_label[scenario]) for candidate in candidates: emitted = oof_severity[(candidate.candidate_id, scenario)] severity_accuracy = float( accuracy_score( y_severity.to_numpy()[true_fault], emitted[true_fault] ) ) score = raw_score(label_f1, severity_accuracy) all_run_rows.append( { "cv_seed": cv_seed, "scenario": scenario, "candidate_id": candidate.candidate_id, "feature_set": candidate.feature_set, "estimator": candidate.estimator, "include_fault_type": candidate.include_fault_type, "per_fault": candidate.per_fault, "macro_f1_oof": label_f1, "severity_accuracy_submission_oof": severity_accuracy, "raw_score_oof": score, "ml_points_oof": ml_points(score), } ) oof_frame = df_clean[ ["engine_id", "cylinder", "label", "severity"] ].copy() oof_frame["predicted_label"] = oof_label[scenario] oof_frame["predicted_severity_submission"] = emitted oof_frame.to_csv( args.output_dir / f"oof__{candidate.candidate_id}__seed_{cv_seed}__{scenario}.csv", index=False, ) run_frame = pd.DataFrame(all_run_rows) fold_frame = pd.DataFrame(all_fold_rows) summary_frame = ( run_frame.groupby(["scenario", "candidate_id"], sort=False) .agg( seeds=("cv_seed", "nunique"), macro_f1_mean=("macro_f1_oof", "mean"), severity_mean=("severity_accuracy_submission_oof", "mean"), severity_std=("severity_accuracy_submission_oof", "std"), severity_min=("severity_accuracy_submission_oof", "min"), severity_max=("severity_accuracy_submission_oof", "max"), raw_score_mean=("raw_score_oof", "mean"), raw_score_std=("raw_score_oof", "std"), raw_score_min=("raw_score_oof", "min"), ml_points_mean=("ml_points_oof", "mean"), ml_points_min=("ml_points_oof", "min"), ) .reset_index() .sort_values( ["scenario", "raw_score_mean", "severity_mean"], ascending=[True, False, False], ) ) run_frame.to_csv(args.output_dir / "severity_runs.csv", index=False) fold_frame.to_csv(args.output_dir / "severity_folds.csv", index=False) summary_frame.to_csv(args.output_dir / "severity_summary.csv", index=False) metadata = { "data": str(args.data), "cv_seeds": args.seeds, "n_splits": args.n_splits, "splitter": "StratifiedGroupKFold", "missing_rate": args.missing_rate, "mask_random_state": args.mask_random_state, "label_pipeline": "combined features + Logistic Regression (frozen)", "candidates": args.candidates, } (args.output_dir / "metadata.json").write_text( json.dumps(metadata, ensure_ascii=False, indent=2), encoding="utf-8" ) display_columns = [ "scenario", "candidate_id", "severity_mean", "severity_std", "severity_min", "raw_score_mean", "raw_score_std", "ml_points_mean", ] print("\nSeverity ranking across CV seeds:") print( summary_frame[display_columns].to_string( index=False, float_format="{:.4f}".format ) ) print(f"\nDetailed outputs saved to: {args.output_dir.resolve()}") if __name__ == "__main__": main()