Code Examples

Production-ready code snippets for fraud detection

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))