predictive-maintenance-rag-system / src /training /classification_training.py
SyedaArisha's picture
Upload folder using huggingface_hub
baf834b verified
Raw History Blame Contribute Delete
11.2 kB
"""
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.")