Quick Start
Basic Fraud Scoring
from fraud_detector import FraudScorer # Initialize the scorer scorer = FraudScorer(model_path="models/xgboost_v2.4.pkl") # Score a transaction transaction = { "amount": 250.00, "merchant_category": "electronics", "card_present": False, "distance_from_home": 150.5, "time_since_last_txn": 0.5 } result = scorer.score(transaction) print(f"Score: {result.score:.3f}, Decision: {result.decision}") # Output: Score: 0.234, Decision: APPROVE
XGBoost Training
Train XGBoost with Class Imbalance Handling
import xgboost as xgb from sklearn.model_selection import train_test_split from sklearn.metrics import roc_auc_score, precision_recall_curve, auc from imblearn.over_sampling import SMOTE # Prepare data X_train, X_test, y_train, y_test = train_test_split( features, labels, test_size=0.2, stratify=labels ) # Handle class imbalance with SMOTE smote = SMOTE(sampling_strategy=0.3, random_state=42) X_train_balanced, y_train_balanced = smote.fit_resample(X_train, y_train) # Calculate class weight for remaining imbalance scale_pos_weight = (y_train_balanced == 0).sum() / (y_train_balanced == 1).sum() # Train XGBoost model = xgb.XGBClassifier( n_estimators=200, max_depth=6, learning_rate=0.1, scale_pos_weight=scale_pos_weight, subsample=0.8, colsample_bytree=0.8, eval_metric="auc", early_stopping_rounds=20, random_state=42 ) model.fit( X_train_balanced, y_train_balanced, eval_set=[(X_test, y_test)], verbose=10 ) # Evaluate y_pred_proba = model.predict_proba(X_test)[:, 1] roc_auc = roc_auc_score(y_test, y_pred_proba) precision, recall, _ = precision_recall_curve(y_test, y_pred_proba) pr_auc = auc(recall, precision) print(f"ROC-AUC: {roc_auc:.4f}, PR-AUC: {pr_auc:.4f}")
Feature Engineering
Real-Time Feature Computation
import redis from datetime import datetime, timedelta import numpy as np class FeatureStore: def __init__(self): self.redis = redis.Redis(host='localhost', port=6379) def compute_velocity_features(self, card_id: str) -> dict: """Compute transaction velocity features for a card.""" now = datetime.utcnow() # Get transaction history from Redis txn_key = f"txn:{card_id}:history" transactions = self.redis.zrangebyscore( txn_key, (now - timedelta(hours=24)).timestamp(), now.timestamp(), withscores=True ) # Parse transactions txn_list = [ {"amount": float(t[0].decode().split(":")[0]), "ts": t[1]} for t in transactions ] # Compute features amounts = [t["amount"] for t in txn_list] return { "txn_count_24h": len(txn_list), "txn_count_1h": sum(1 for t in txn_list if t["ts"] > (now - timedelta(hours=1)).timestamp()), "amount_sum_24h": sum(amounts), "amount_avg_24h": np.mean(amounts) if amounts else 0, "amount_std_24h": np.std(amounts) if len(amounts) > 1 else 0, }
Isolation Forest
Anomaly Detection Pipeline
from sklearn.ensemble import IsolationForest from sklearn.preprocessing import StandardScaler import numpy as np class AnomalyDetector: def __init__(self, contamination=0.01): self.scaler = StandardScaler() self.model = IsolationForest( n_estimators=100, contamination=contamination, max_samples="auto", random_state=42, n_jobs=-1 ) def fit(self, X): """Fit on historical transaction data.""" X_scaled = self.scaler.fit_transform(X) self.model.fit(X_scaled) return self def score_transaction(self, transaction: np.ndarray) -> float: """ Score a single transaction. Returns: anomaly score between 0 and 1 (higher = more anomalous) """ X_scaled = self.scaler.transform(transaction.reshape(1, -1)) # Get raw anomaly score (negative = anomaly) raw_score = self.model.decision_function(X_scaled)[0] # Convert to 0-1 range where 1 = most anomalous anomaly_score = 1 - (1 / (1 + np.exp(-raw_score))) return anomaly_score # Usage detector = AnomalyDetector(contamination=0.01) detector.fit(historical_transactions) score = detector.score_transaction(new_transaction) print(f"Anomaly score: {score:.3f}")
Kafka Streaming
Real-Time Transaction Consumer
from kafka import KafkaConsumer, KafkaProducer import json import time class FraudStreamProcessor: def __init__(self, scorer, feature_store): self.scorer = scorer self.feature_store = feature_store self.consumer = KafkaConsumer( 'transactions', bootstrap_servers=['localhost:9092'], value_deserializer=lambda m: json.loads(m.decode('utf-8')), auto_offset_reset='latest' ) self.producer = KafkaProducer( bootstrap_servers=['localhost:9092'], value_serializer=lambda v: json.dumps(v).encode('utf-8') ) def process_transaction(self, txn: dict) -> dict: """Score a transaction and return decision.""" start_time = time.time() # Fetch real-time features features = self.feature_store.get_features(txn['card_id']) features.update(txn) # Score transaction result = self.scorer.score(features) latency_ms = (time.time() - start_time) * 1000 return { 'transaction_id': txn['transaction_id'], 'score': result.score, 'decision': result.decision, 'latency_ms': latency_ms } def run(self): """Start processing transactions.""" print("Starting fraud detection stream...") for message in self.consumer: txn = message.value result = self.process_transaction(txn) # Publish decision self.producer.send('fraud_decisions', result) if result['decision'] == 'BLOCK': self.producer.send('fraud_alerts', result)
API Integration
FastAPI Endpoint
from fastapi import FastAPI, HTTPException from pydantic import BaseModel import time app = FastAPI(title="Fraud Detection API") class Transaction(BaseModel): transaction_id: str card_id: str amount: float merchant_category: str card_present: bool class ScoreResponse(BaseModel): transaction_id: str score: float decision: str latency_ms: float @app.post("/api/v1/score", response_model=ScoreResponse) async def score_transaction(txn: Transaction): start = time.time() try: # Fetch features and score features = feature_store.get_features(txn.card_id) result = scorer.score({**txn.dict(), **features}) latency = (time.time() - start) * 1000 return ScoreResponse( transaction_id=txn.transaction_id, score=result.score, decision=result.decision, latency_ms=latency ) except Exception as e: raise HTTPException(status_code=500, detail=str(e))