Spaces:
Sleeping
Sleeping
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") |