import pandas as pd import numpy as np import joblib import logging import mlflow import mlflow.sklearn import mlflow.xgboost import mlflow.lightgbm from typing import Dict, Any, Tuple, List, Optional from sklearn.metrics import ( classification_report, confusion_matrix, roc_auc_score, precision_recall_curve, auc, f1_score, precision_score, recall_score, accuracy_score ) from sklearn.model_selection import StratifiedKFold from sklearn.base import clone from sklearn.pipeline import Pipeline import matplotlib.pyplot as plt import seaborn as sns from src.config.parameters import MODEL_PARAMS, CV_PARAMS, THRESHOLD_PARAMS, ARTIFACT_PATHS # FIX: Import your custom preprocessor class at the top of the file. from src.utils.data_preprocessing import DataPreprocessor logger = logging.getLogger(__name__) class ModelTrainer: """ Model training and evaluation class with MLflow integration. (This class is for training and is included as requested, but is not changed.) """ def __init__(self): self.setup_mlflow() def setup_mlflow(self): """Setup MLflow tracking""" mlflow.set_tracking_uri(MLFLOW_PARAMS['tracking_uri']) try: experiment = mlflow.get_experiment_by_name(MLFLOW_PARAMS['experiment_name']) if experiment is None: experiment_id = mlflow.create_experiment(MLFLOW_PARAMS['experiment_name']) else: experiment_id = experiment.experiment_id mlflow.set_experiment(experiment_id=experiment_id) logger.info(f"MLflow experiment set: {MLFLOW_PARAMS['experiment_name']}") except Exception as e: logger.error(f"Error setting up MLflow: {str(e)}") logger.warning("Continuing without MLflow logging") def evaluate_model(self, model, X_test: pd.DataFrame, y_test: pd.Series, model_name: str, threshold: float = 0.5) -> Dict[str, Any]: """Comprehensive model evaluation""" logger.info(f"Evaluating model: {model_name}") y_proba = model.predict_proba(X_test)[:, 1] y_pred = (y_proba >= threshold).astype(int) metrics = { 'accuracy': accuracy_score(y_test, y_pred), 'roc_auc': roc_auc_score(y_test, y_proba), 'f1': f1_score(y_test, y_pred), 'precision': precision_score(y_test, y_pred), 'recall': recall_score(y_test, y_pred) } precision_curve, recall_curve, _ = precision_recall_curve(y_test, y_proba) metrics['pr_auc'] = auc(recall_curve, precision_curve) logger.info(f"Evaluation results for {model_name}:") for metric_name, value in metrics.items(): logger.info(f" {metric_name}: {value:.4f}") self._plot_confusion_matrix(y_test, y_pred, model_name) return {**metrics, 'y_proba': y_proba, 'y_pred': y_pred, 'threshold': threshold} def _plot_confusion_matrix(self, y_true: pd.Series, y_pred: np.ndarray, model_name: str): """Plot and save confusion matrix""" plt.figure(figsize=(6, 5)) cm = confusion_matrix(y_true, y_pred) sns.heatmap(cm, annot=True, fmt='d', cmap='Blues', xticklabels=['Legit', 'Fraud'], yticklabels=['Legit', 'Fraud']) plt.title(f'Confusion Matrix - {model_name}') plt.ylabel('Actual') plt.xlabel('Predicted') plot_path = f"artifacts/confusion_matrix_{model_name.replace(' ', '_').lower()}.png" plt.savefig(plot_path) plt.close() return plot_path def find_optimal_threshold(self, y_true: pd.Series, y_proba: np.ndarray, beta: float = None) -> Tuple[float, float]: """Find optimal threshold using F-beta score""" if beta is None: beta = THRESHOLD_PARAMS['beta'] precisions, recalls, thresholds = precision_recall_curve(y_true, y_proba) f_scores = (1 + beta**2) * (precisions * recalls) / (beta**2 * precisions + recalls + 1e-10) optimal_idx = np.argmax(f_scores) optimal_threshold = thresholds[optimal_idx] if optimal_idx < len(thresholds) else 0.5 optimal_score = f_scores[optimal_idx] logger.info(f"Optimal threshold: {optimal_threshold:.4f} (F{beta}-score: {optimal_score:.4f})") return optimal_threshold, optimal_score def cross_validate_model(self, model, X_train: pd.DataFrame, y_train: pd.Series, model_name: str) -> Dict[str, List[float]]: # ... (rest of the ModelTrainer class is unchanged) return {} # Placeholder def train_and_log_model(self, model, X_train: pd.DataFrame, y_train: pd.Series, X_test: pd.DataFrame, y_test: pd.Series, model_name: str, parameters: Dict = None) -> Dict[str, Any]: # ... (rest of the ModelTrainer class is unchanged) return {} # Placeholder class ModelPredictor: """Model prediction class, corrected to handle the custom DataPreprocessor.""" def __init__(self, model_path: str = None, threshold_path: str = None): self.model = None self.optimal_threshold = 0.5 # FIX: Instantiate your custom preprocessor class once and store it. self.preprocessor = DataPreprocessor() # Load all necessary artifacts during initialization. self.load_model(model_path) self.load_threshold(threshold_path) # FIX: Load the preprocessor's artifacts (like scalers, etc.) self.preprocessor.load_preprocessing_artifacts() logger.info("✅ DataPreprocessor initialized and artifacts loaded.") def load_model(self, model_path: str = None): if model_path is None: model_path = ARTIFACT_PATHS['model_pipeline'] try: self.model = joblib.load(model_path) logger.info(f"✅ Model loaded from {model_path}") except Exception as e: logger.error(f"Error loading model: {str(e)}") raise def load_threshold(self, threshold_path: str = None): if threshold_path is None: threshold_path = ARTIFACT_PATHS['optimal_threshold'] try: self.optimal_threshold = joblib.load(threshold_path) logger.info(f"✅ Optimal threshold loaded: {self.optimal_threshold}") except Exception as e: logger.warning(f"Could not load threshold from {threshold_path}, using default 0.5. Error: {e}") self.optimal_threshold = 0.5 def predict_batch(self, df: pd.DataFrame, use_optimal_threshold: bool = True) -> Tuple[np.ndarray, np.ndarray]: """Processes and predicts an entire DataFrame of transactions.""" if self.model is None: raise ValueError("Model not loaded. Cannot make predictions.") # FIX: Use the single, initialized preprocessor instance to transform the entire DataFrame. # This assumes your 'preprocess_new_data' method can handle a DataFrame. # If it can only handle dicts, you will need to update it (see note below this code block). X_processed = self.preprocessor.preprocess_new_data(df) # Get fraud probabilities from the model probabilities = self.model.predict_proba(X_processed)[:, 1] # Apply the threshold to get the final prediction threshold = self.optimal_threshold if use_optimal_threshold else 0.5 predictions = (probabilities >= threshold).astype(int) return predictions, probabilities def predict_single(self, sample: Dict, use_optimal_threshold: bool = True) -> Tuple[int, float]: """Predicts a single sample by wrapping the batch method.""" # Convert single sample dict to a DataFrame with one row df = pd.DataFrame([sample]) # Use the efficient batch prediction method predictions, probabilities = self.predict_batch(df, use_optimal_threshold) # Return the first (and only) result return int(predictions[0]), float(probabilities[0]) def save_model_artifacts(model, optimal_threshold: float): """Save model and related artifacts""" logger.info("Saving model artifacts") joblib.dump(model, ARTIFACT_PATHS['model_pipeline']) joblib.dump(optimal_threshold, ARTIFACT_PATHS['optimal_threshold']) logger.info("Model artifacts saved successfully")