File size: 8,279 Bytes
e146383
 
 
 
c8db108
 
 
 
 
 
 
 
 
 
 
 
 
e146383
 
bc38a2f
c8db108
 
 
 
e146383
 
 
 
 
c8db108
 
e146383
c8db108
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
e146383
 
 
c8db108
e146383
 
 
c8db108
 
 
 
 
 
e146383
 
c8db108
 
 
 
e146383
 
 
 
 
 
c8db108
e146383
c8db108
e146383
 
 
 
 
 
 
c8db108
e146383
c8db108
 
e146383
 
c8db108
e146383
 
 
c8db108
 
 
 
e146383
c8db108
e146383
 
c8db108
e146383
 
 
 
 
 
c8db108
 
e146383
 
 
 
 
 
 
 
 
 
 
c8db108
e146383
 
c8db108
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
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")