""" Classification Training Module for LLM-Enhanced Predictive Maintenance Author: Antigravity AI Date: August 2026 This module trains failure classification models (XGBoost + Random Forest) on the AI4I 2020 and Pump Sensor datasets. Features: - Handles class imbalance using a hybrid SMOTE + Class Weighting approach. - Cross-validation: StratifiedKFold for tabular AI4I; TimeSeriesSplit for Pump. - Hyperparameter tuning using Optuna (with RandomizedSearchCV fallback). - Model ensembling via Soft Voting Classifier. - Full performance evaluation (Precision, Recall, F1, ROC-AUC, Confusion Matrix). - Saves trained model artifacts securely. """ import os import time import logging import numpy as np import pandas as pd from sklearn.model_selection import StratifiedKFold, TimeSeriesSplit, RandomizedSearchCV from sklearn.metrics import classification_report, confusion_matrix, roc_auc_score from sklearn.ensemble import RandomForestClassifier, VotingClassifier from xgboost import XGBClassifier import joblib try: from utils.project_paths import LOGS_DIR except ImportError: from project_paths import LOGS_DIR # Optional imports with graceful fallbacks try: import optuna optuna.logging.set_verbosity(optuna.logging.WARNING) HAS_OPTUNA = True except ImportError: HAS_OPTUNA = False try: from imblearn.over_sampling import SMOTE HAS_SMOTE = True except ImportError: HAS_SMOTE = False # Setup logging os.makedirs(LOGS_DIR, exist_ok=True) logging.basicConfig( level=logging.INFO, format='%(asctime)s [%(levelname)s] %(message)s', handlers=[ logging.StreamHandler(), logging.FileHandler(os.path.join(LOGS_DIR, "classification_training.log")) ] ) logger = logging.getLogger(__name__) def time_tracker(func): """Decorator to measure and log execution time.""" def wrapper(*args, **kwargs): start_time = time.time() logger.info(f"Starting: {func.__name__}") result = func(*args, **kwargs) elapsed = time.time() - start_time logger.info(f"Finished: {func.__name__}. Duration: {elapsed:.2f} seconds.") return result return wrapper @time_tracker def evaluate_model(y_true, y_pred, y_proba, model_name: str, dataset_name: str): """Print and return evaluation metrics for classification models.""" logger.info(f"--- Evaluation for {model_name} on {dataset_name} ---") # Classification report (Precision, Recall, F1) rep = classification_report(y_true, y_pred) logger.info(f"Classification Report:\n{rep}") # Confusion Matrix cm = confusion_matrix(y_true, y_pred) logger.info(f"Confusion Matrix:\n{cm}") # ROC-AUC Score roc_auc = roc_auc_score(y_true, y_proba) logger.info(f"ROC-AUC Score: {roc_auc:.4f}") return { 'report': rep, 'confusion_matrix': cm, 'roc_auc': roc_auc } @time_tracker def tune_xgboost(X, y, cv, scale_pos_weight: float, n_trials: int = 15) -> dict: """Optimize XGBoost parameters using Optuna if available, else standard fallback.""" if not HAS_OPTUNA: logger.warning("Optuna not installed. Using predefined optimized XGBoost parameters.") return { 'max_depth': 6, 'learning_rate': 0.05, 'n_estimators': 200, 'subsample': 0.8, 'colsample_bytree': 0.8, 'scale_pos_weight': scale_pos_weight } logger.info("Starting Optuna hyperparameter tuning for XGBoost...") def objective(trial): params = { 'n_estimators': trial.suggest_int('n_estimators', 100, 300), 'max_depth': trial.suggest_int('max_depth', 3, 9), 'learning_rate': trial.suggest_float('learning_rate', 0.01, 0.2, log=True), 'subsample': trial.suggest_float('subsample', 0.6, 1.0), 'colsample_bytree': trial.suggest_float('colsample_bytree', 0.6, 1.0), 'scale_pos_weight': scale_pos_weight, 'random_state': 42, 'n_jobs': -1, 'eval_metric': 'logloss' } # Cross validation scores scores = [] for train_idx, val_idx in cv.split(X, y): X_tr, X_va = X.iloc[train_idx], X.iloc[val_idx] y_tr, y_va = y.iloc[train_idx], y.iloc[val_idx] # Apply SMOTE strictly to training fold to avoid leakage if HAS_SMOTE: sm = SMOTE(random_state=42) X_tr, y_tr = sm.fit_resample(X_tr, y_tr) model = XGBClassifier(**params) model.fit(X_tr, y_tr) preds = model.predict(X_va) # Use F1-score as optimization metric due to severe class imbalance from sklearn.metrics import f1_score scores.append(f1_score(y_va, preds)) return np.mean(scores) study = optuna.create_study(direction='maximize') study.optimize(objective, n_trials=n_trials) logger.info(f"Best XGBoost Trial: Value {study.best_value:.4f}") logger.info(f"Best XGBoost Parameters: {study.best_params}") best_params = study.best_params best_params['scale_pos_weight'] = scale_pos_weight return best_params @time_tracker def tune_random_forest(X, y, cv, n_trials: int = 15) -> dict: """Optimize Random Forest parameters using Optuna if available, else standard fallback.""" if not HAS_OPTUNA: logger.warning("Optuna not installed. Using predefined optimized Random Forest parameters.") return { 'n_estimators': 150, 'max_depth': 8, 'min_samples_split': 5, 'class_weight': 'balanced' } logger.info("Starting Optuna hyperparameter tuning for Random Forest...") def objective(trial): params = { 'n_estimators': trial.suggest_int('n_estimators', 50, 200), 'max_depth': trial.suggest_int('max_depth', 4, 12), 'min_samples_split': trial.suggest_int('min_samples_split', 2, 10), 'class_weight': 'balanced', 'random_state': 42, 'n_jobs': -1 } scores = [] for train_idx, val_idx in cv.split(X, y): X_tr, X_va = X.iloc[train_idx], X.iloc[val_idx] y_tr, y_va = y.iloc[train_idx], y.iloc[val_idx] if HAS_SMOTE: sm = SMOTE(random_state=42) X_tr, y_tr = sm.fit_resample(X_tr, y_tr) model = RandomForestClassifier(**params) model.fit(X_tr, y_tr) preds = model.predict(X_va) from sklearn.metrics import f1_score scores.append(f1_score(y_va, preds)) return np.mean(scores) study = optuna.create_study(direction='maximize') study.optimize(objective, n_trials=n_trials) logger.info(f"Best Random Forest Trial: Value {study.best_value:.4f}") logger.info(f"Best Random Forest Parameters: {study.best_params}") best_params = study.best_params best_params['class_weight'] = 'balanced' return best_params @time_tracker def train_pipeline(train_path: str, test_path: str, dataset_name: str, cv_type: str, output_dir: str): """ Complete model training pipeline (tune -> cross validate -> ensemble -> evaluate -> save). Args: train_path: Path to preprocessed train CSV. test_path: Path to preprocessed test CSV. dataset_name: Name of dataset (e.g. 'AI4I' or 'Pump'). cv_type: Cross-validation type ('stratified' or 'timeseries'). output_dir: Directory where model files will be saved. """ os.makedirs(output_dir, exist_ok=True) logger.info(f"Starting classification pipeline for {dataset_name} dataset...") # 1. Load data train_df = pd.read_csv(train_path) test_df = pd.read_csv(test_path) # Drops timestamp identifier if present drop_cols = ['timestamp', 'target', 'TWF', 'HDF', 'PWF', 'OSF', 'RNF'] features = [col for col in train_df.columns if col not in drop_cols] X_train = train_df[features] y_train = train_df['target'] X_test = test_df[features] y_test = test_df['target'] logger.info(f"Dataset shape - Train: {X_train.shape}, Test: {X_test.shape}") logger.info(f"Class imbalance - Normal: {sum(y_train == 0)}, Failures: {sum(y_train == 1)} (Ratio: {sum(y_train == 0)/sum(y_train == 1):.2f}:1)") # 2. Setup Cross-Validation if cv_type == 'stratified': cv = StratifiedKFold(n_splits=5, shuffle=True, random_state=42) elif cv_type == 'timeseries': cv = TimeSeriesSplit(n_splits=5) else: raise ValueError(f"Unknown cv_type: {cv_type}") # Scale positive weight calculation for XGBoost scale_pos_weight = float(sum(y_train == 0) / sum(y_train == 1)) # 3. Hyperparameter Tuning xgb_params = tune_xgboost(X_train, y_train, cv, scale_pos_weight) rf_params = tune_random_forest(X_train, y_train, cv) # 4. Fit individual best models on the FULL training set logger.info("Fitting best models on full training dataset (applying SMOTE + weights)...") X_tr_final, y_tr_final = X_train, y_train if HAS_SMOTE: logger.info("Applying SMOTE oversampling to training set...") sm = SMOTE(random_state=42) X_tr_final, y_tr_final = sm.fit_resample(X_train, y_train) best_xgb = XGBClassifier(**xgb_params, random_state=42, n_jobs=-1, eval_metric='logloss') best_xgb.fit(X_tr_final, y_tr_final) best_rf = RandomForestClassifier(**rf_params, random_state=42, n_jobs=-1) best_rf.fit(X_tr_final, y_tr_final) # Evaluate individual models on Test set xgb_test_preds = best_xgb.predict(X_test) xgb_test_proba = best_xgb.predict_proba(X_test)[:, 1] evaluate_model(y_test, xgb_test_preds, xgb_test_proba, "Tuned XGBoost", dataset_name) rf_test_preds = best_rf.predict(X_test) rf_test_proba = best_rf.predict_proba(X_test)[:, 1] evaluate_model(y_test, rf_test_preds, rf_test_proba, "Tuned Random Forest", dataset_name) # 5. Build Soft Voting Ensemble logger.info("Building Soft Voting Ensemble (XGBoost + Random Forest)...") ensemble = VotingClassifier( estimators=[ ('xgb', best_xgb), ('rf', best_rf) ], voting='soft', n_jobs=-1 ) # Fit voting classifier on resampled train set ensemble.fit(X_tr_final, y_tr_final) # Evaluate Ensemble on Test set ens_test_preds = ensemble.predict(X_test) ens_test_proba = ensemble.predict_proba(X_test)[:, 1] evals = evaluate_model(y_test, ens_test_preds, ens_test_proba, "Soft Voting Ensemble", dataset_name) # 6. Save the ensemble model and parameters model_save_path = os.path.join(output_dir, f'{dataset_name.lower()}_ensemble_model.pkl') joblib.dump(ensemble, model_save_path) logger.info(f"Saved {dataset_name} ensemble model to {model_save_path}") return evals if __name__ == '__main__': logger.info("Classification Training Module template loaded.")