AI Engineering
End-to-end fraud detection pipeline
An MLOps pipeline for credit card fraud detection: DVC-versioned data, PySpark training behind an MLflow metric gate, a FastAPI scoring service, and drift-triggered retraining on a MinIO/Prometheus/Grafana stack.
Problem
Credit card fraud detection is a moving target. The transaction distribution shifts over time, and a model that scored well at training slowly decays in production without anyone noticing until fraud slips through or good customers get blocked.
A one-off notebook model does not survive that. It needs versioned data you can reproduce, a training run whose quality is checked before anything ships, a serving path that scores a live transaction exactly the way training saw it, and a monitor that watches for distribution drift and rebuilds the model when the data moves.
This project builds that full lifecycle around a PySpark fraud classifier so a model change is reproducible, gated, served without training/serving skew, and retrained when drift crosses a threshold.
How it works
The pipeline is a four-stage DVC graph: preprocess, feature_engineering, train, evaluate. Each stage declares its inputs, params, and outputs, so a repro reruns only what changed and the whole run is reproducible from the raw dataset.
Training fits a Spark ML Pipeline that bundles the VectorAssembler with a RandomForestClassifier into one saved PipelineModel, and logs parameters, the model, and a signature to MLflow. The evaluate stage scores the held-out test set and writes a tracked metrics file.
Promotion is gated. A separate CI step reads the held-out AUC and fails the run when it falls below the configured floor, so a regression cannot merge even when the repro itself succeeds. The FastAPI service exposes /health, /predict, and /metrics; at startup it loads the persisted PipelineModel and the fitted scaler, and on each request rebuilds the training features through the same code path training used. Observability runs on a Docker Compose stack, with MinIO as the S3-compatible DVC remote.
raw data (DVC) -> preprocess -> feature_engineering -> train (MLflow) -> evaluate
-> metrics.json -> AUC gate -> registry promotion
-> FastAPI (/health /predict /metrics) -> Prometheus -> Grafana
| scheduled drift check -> drift_report.json -> retrain workflowHard parts
- Training/serving skew: the raw request does not carry the fitted scaler or engineered features, so scoring with a bare assembler would feed the model a different vector than it trained on. Serving replays the exact training chain, applying the persisted scaler (loaded, never re-fit) then the shared feature transforms.
- Fail-closed metric gate: NaN < floor is False in Python, so an undefined AUC would silently pass a naive comparison. The gate rejects non-numeric, boolean, and non-finite values and treats a missing metrics file as a failure before the floor comparison.
- Model-shape ambiguity: a single-class model emits a length-1 probability vector; the API detects that and returns a clear error message.
- Real versus placeholder infrastructure: the monitoring and retrain workflows require configured infrastructure and fail loudly when it is missing, so a run never reports success without having retrained.
- Drift-driven retraining: a scheduled workflow scores new data with PSI, the KS test, and Jensen-Shannon distance, writes a drift report with a retrain flag, and dispatches the retrain only when the drifted-feature percentage crosses the threshold.
Results
- The committed evaluation records a held-out AUC of 0.994, above the 0.90 gate floor, with precision, recall, F1, and accuracy around 0.986 to 0.988.
- The AUC gate blocks any merge or publish when the held-out metric drops below the floor, and fails closed on NaN, non-numeric, and missing metrics.
- The serving path reproduces the full training feature chain, removing training/serving skew as a source of silently wrong predictions, guarded by a serving-parity test.
- Drift monitoring runs on a schedule and only triggers a retrain when drift crosses the threshold, so retraining is driven by measured drift.
- The serving image publishes to the registry only on a merge to main after the pipeline gate passes.
Artifacts
- DVC pipeline with preprocess, feature_engineering, train, and evaluate stages and tracked metrics.
- FastAPI serving app exposing /health, /predict, and /metrics, plus a Prometheus metrics module and a non-root serving Dockerfile.
- MLflow training and evaluation code, plus a registry helper that promotes a version only when its AUC clears the floor.
- A CI metrics gate and GitHub Actions workflows for lint, test, coverage, the ML pipeline, monitoring, and retraining.
- A Docker Compose observability stack (MinIO, Prometheus, Grafana) and Kubernetes manifests, plus advanced drift detection using PSI, KS, and Jensen-Shannon distance.