Tutorial 20: MLOps End-to-End Project
Table of Contents
Introduction
MLOps is the discipline of deploying and maintaining machine learning models in production reliably and efficiently. While building an accurate model is important, the real challenge lies in everything around it: versioning data and code together, tracking experiments reproducibly, automating training and deployment pipelines, monitoring model performance in production, and responding to data drift.
This tutorial walks through a complete MLOps project from raw data to production monitoring. We will build a customer churn prediction system using industry-standard tools: DVC for data versioning, MLflow for experiment tracking and model registry, GitHub Actions for CI/CD, Docker for containerization, and Evidently for production monitoring.
Prerequisites
- Python 3.9+
- Git and GitHub account
- Docker and Docker Compose
- AWS CLI or equivalent cloud CLI (for deployment)
- Basic understanding of ML model training
# Install all required packages
pip install dvc[s3] mlflow scikit-learn pandas evidently
pip install fastapi uvicorn docker boto3
import os
print("MLOps E2E Tutorial - Environment Setup")
Project Overview
Our project structure follows MLOps best practices with clear separation of concerns.
churn-prediction/
├── .github/
│ └── workflows/
│ ├── train.yml
│ ├── test.yml
│ └── deploy.yml
├── data/
│ ├── raw/
│ │ └── customers.csv.dvc
│ └── processed/
│ └── features.csv.dvc
├── src/
│ ├── data/
│ │ ├── init.py
│ │ ├── prepare.py
│ │ └── validate.py
│ ├── features/
│ │ ├── init.py
│ │ └── buildfeatures.py
│ ├── models/
│ │ ├── init.py
│ │ ├── train.py
│ │ └── predict.py
│ └── monitoring/
│ ├── init.py
│ └── driftdetection.py
├── serving/
│ ├── app.py
│ ├── Dockerfile
│ └── requirements.txt
├── tests/
│ ├── testdata.py
│ ├── testmodel.py
│ └── testapi.py
├── configs/
│ └── config.yaml
├── dvc.yaml
├── dvc.lock
├── params.yaml
├── docker-compose.yml
└── requirements.txt
Data Versioning with DVC
DVC (Data Version Control) tracks large datasets and model files alongside your Git repository without storing them in Git itself.
Setting Up DVC
# Initialize DVC in your Git repository
cd churn-prediction
dvc init
Configure remote storage (S3 in this example)
dvc remote add -d myremote s3://my-ml-bucket/dvc-store
dvc remote modify myremote region us-east-1
Track data files
dvc add data/raw/customers.csv
git add data/raw/customers.csv.dvc data/raw/.gitignore
git commit -m "Track raw customer data with DVC"
Push data to remote storage
dvc push
DVC Pipeline Definition
# dvc.yaml - Defines the reproducible ML pipeline
stages:
prepare:
cmd: python src/data/prepare.py
deps:
- src/data/prepare.py
- data/raw/customers.csv
params:
- prepare.testsize
- prepare.randomseed
outs:
- data/processed/train.csv
- data/processed/test.csv
featurize:
cmd: python src/features/buildfeatures.py
deps:
- src/features/buildfeatures.py
- data/processed/train.csv
- data/processed/test.csv
params:
- features
outs:
- data/processed/trainfeatures.csv
- data/processed/testfeatures.csv
train:
cmd: python src/models/train.py
deps:
- src/models/train.py
- data/processed/trainfeatures.csv
params:
- train
outs:
- models/model.pkl
metrics:
- metrics/trainmetrics.json:
cache: false
plots:
- metrics/roccurve.json:
cache: false
x: fpr
y: tpr
evaluate:
cmd: python src/models/evaluate.py
deps:
- src/models/evaluate.py
- models/model.pkl
- data/processed/testfeatures.csv
metrics:
- metrics/evalmetrics.json:
cache: false
# params.yaml - Central parameter configuration
prepare:
testsize: 0.2
randomseed: 42
features:
numerical:
- tenure
- monthlycharges
- totalcharges
categorical:
- contracttype
- paymentmethod
- internetservice
engineered:
- chargespermonth
- tenuregroup
train:
model: xgboost
nestimators: 500
maxdepth: 6
learningrate: 0.01
subsample: 0.8
colsamplebytree: 0.8
earlystoppingrounds: 50
Data Preparation Script
# src/data/prepare.py
import pandas as pd
from sklearn.modelselection import traintestsplit
import yaml
import os
def preparedata():
"""Split raw data into train and test sets."""
with open("params.yaml", "r") as f:
params = yaml.safeload(f)["prepare"]
# Load raw data
df = pd.readcsv("data/raw/customers.csv")
# Basic cleaning
df = df.dropna(subset=["churn"])
df["totalcharges"] = pd.tonumeric(
df["totalcharges"], errors="coerce"
)
df = df.fillna(df.median(numericonly=True))
# Split
traindf, testdf = traintestsplit(
df,
testsize=params["testsize"],
randomstate=params["randomseed"],
stratify=df["churn"]
)
os.makedirs("data/processed", existok=True)
traindf.tocsv("data/processed/train.csv", index=False)
testdf.tocsv("data/processed/test.csv", index=False)
print(f"Train: {len(traindf)} rows, Test: {len(testdf)} rows")
print(f"Churn rate - Train: {traindf['churn'].mean():.3f}, "
f"Test: {testdf['churn'].mean():.3f}")
if name == "main":
preparedata()
Experiment Tracking with MLflow
MLflow provides experiment tracking, model logging, and a model registry — all essential for reproducible ML.
Setting Up MLflow
# Start MLflow tracking server
mlflow server --backend-store-uri sqlite:///mlflow.db \
--default-artifact-root ./mlruns \
--host 0.0.0.0 --port 5000
import mlflow
import mlflow.sklearn
import mlflow.xgboost
Configure MLflow
mlflow.settrackinguri("http://localhost:5000")
mlflow.setexperiment("churn-prediction")
Training Script with MLflow Integration
# src/models/train.py
import pandas as pd
import numpy as np
import xgboost as xgb
from sklearn.metrics import (
accuracyscore, precisionscore, recallscore,
f1score, rocaucscore, roccurve
)
import mlflow
import mlflow.xgboost
import yaml
import json
import joblib
import os
def trainmodel():
"""Train model with full MLflow tracking."""
with open("params.yaml", "r") as f:
params = yaml.safeload(f)
trainparams = params["train"]
# Load features
traindf = pd.readcsv("data/processed/trainfeatures.csv")
Xtrain = traindf.drop(columns=["churn"])
ytrain = traindf["churn"]
# Configure MLflow
mlflow.settrackinguri("http://localhost:5000")
mlflow.setexperiment("churn-prediction")
with mlflow.startrun(runname="xgboosttraining") as run:
# Log parameters
mlflow.logparams(trainparams)
mlflow.logparam("numfeatures", Xtrain.shape[1])
mlflow.logparam("trainsamples", Xtrain.shape[0])
# Log dataset info
mlflow.settag("dataversion", os.popen("dvc version").read().strip())
mlflow.settag("gitcommit",
os.popen("git rev-parse HEAD").read().strip())
# Train model
model = xgb.XGBClassifier(
nestimators=trainparams["nestimators"],
maxdepth=trainparams["maxdepth"],
learningrate=trainparams["learningrate"],
subsample=trainparams["subsample"],
colsamplebytree=trainparams["colsamplebytree"],
evalmetric="logloss",
uselabelencoder=False,
randomstate=42
)
# Train with validation split for early stopping
from sklearn.modelselection import traintestsplit
Xtr, Xval, ytr, yval = traintestsplit(
Xtrain, ytrain, testsize=0.15, randomstate=42
)
model.fit(
Xtr, ytr,
evalset=[(Xval, yval)],
verbose=False
)
# Predictions and metrics
ypred = model.predict(Xval)
yprob = model.predictproba(Xval)[:, 1]
metrics = {
"accuracy": accuracyscore(yval, ypred),
"precision": precisionscore(yval, ypred),
"recall": recallscore(yval, ypred),
"f1": f1score(yval, ypred),
"aucroc": rocaucscore(yval, yprob),
}
# Log metrics
mlflow.logmetrics(metrics)
# Log feature importance
importance = dict(zip(
Xtrain.columns,
model.featureimportances.tolist()
))
mlflow.logdict(importance, "featureimportance.json")
# Log ROC curve data
fpr, tpr, = roccurve(yval, yprob)
rocdata = [{"fpr": f, "tpr": t}
for f, t in zip(fpr.tolist(), tpr.tolist())]
# Log model
mlflow.xgboost.logmodel(
model,
artifactpath="model",
registeredmodelname="churn-predictor",
inputexample=Xval.iloc[:3],
)
# Save locally for DVC tracking
os.makedirs("models", existok=True)
joblib.dump(model, "models/model.pkl")
# Save metrics for DVC
os.makedirs("metrics", existok=True)
with open("metrics/trainmetrics.json", "w") as f:
json.dump(metrics, f, indent=2)
with open("metrics/roccurve.json", "w") as f:
json.dump(rocdata, f)
print(f"Run ID: {run.info.runid}")
for name, value in metrics.items():
print(f" {name}: {value:.4f}")
return run.info.runid
if name == "main":
trainmodel()
Model Registry
The MLflow Model Registry provides a central hub for managing model lifecycle stages.
# src/models/registry.py
import mlflow
from mlflow.tracking import MlflowClient
client = MlflowClient("http://localhost:5000")
def promotemodel(modelname: str, version: int, stage: str):
"""Promote a model version to a new stage."""
validstages = ["Staging", "Production", "Archived"]
if stage not in validstages:
raise ValueError(f"Stage must be one of {validstages}")
client.transitionmodelversionstage(
name=modelname,
version=version,
stage=stage,
archiveexistingversions=(stage == "Production")
)
print(f"Model {modelname} v{version} promoted to {stage}")
def getproductionmodel(modelname: str):
"""Load the current production model."""
modeluri = f"models:/{modelname}/Production"
model = mlflow.pyfunc.loadmodel(modeluri)
return model
def comparemodels(modelname: str, metric: str = "aucroc"):
"""Compare all versions of a model by a given metric."""
versions = client.searchmodelversions(f"name='{modelname}'")
results = []
for v in versions:
run = client.getrun(v.runid)
metricvalue = run.data.metrics.get(metric, 0)
results.append({
"version": v.version,
"stage": v.currentstage,
"metric": metricvalue,
"runid": v.runid,
})
results.sort(key=lambda x: x["metric"], reverse=True)
print(f"\nModel: {modelname} | Metric: {metric}")
print("-" 60)
for r in results:
print(f" v{r['version']} ({r['stage']}): {r['metric']:.4f}")
return results
Usage
if name == "main":
comparemodels("churn-predictor", metric="aucroc")
# promotemodel("churn-predictor", version=3, stage="Production")
CI/CD with GitHub Actions
Training Pipeline Workflow
# .github/workflows/train.yml
name: ML Training Pipeline
on:
push:
paths:
- 'src/'
- 'params.yaml'
- 'dvc.yaml'
workflowdispatch:
inputs:
retrain:
description: 'Force retraining'
required: false
default: 'false'
jobs:
train:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Set up Python
uses: actions/setup-python@v5
with:
python-version: '3.11'
- name: Install dependencies
run: |
pip install -r requirements.txt
pip install dvc[s3]
- name: Configure AWS credentials
uses: aws-actions/configure-aws-credentials@v4
with:
aws-access-key-id: ${{ secrets.AWS
ACCESSKEYID }}
aws-secret-access-key: ${{ secrets.AWSSECRETACCESSKEY }}
aws-region: us-east-1
- name: Pull data from DVC
run: dvc pull
- name: Run DVC pipeline
env:
MLFLOWTRACKINGURI: ${{ secrets.MLFLOWTRACKINGURI }}
run: dvc repro
- name: Run tests
run: pytest tests/ -v --tb=short
- name: Check model quality gate
run: |
python -c "
import json
with open('metrics/evalmetrics.json') as f:
metrics = json.load(f)
assert metrics['aucroc'] > 0.85, \
f'AUC {metrics[\"aucroc\"]:.4f} below threshold 0.85'
assert metrics['f1'] > 0.75, \
f'F1 {metrics[\"f1\"]:.4f} below threshold 0.75'
print('Quality gate PASSED')
"
- name: Push metrics
run: |
git config user.name github-actions
git config user.email github-actions@github.com
git add metrics/
git diff --cached --quiet || git commit -m "Update metrics [skip ci]"
git push
deploy:
needs: train
if: github.ref == 'refs/heads/main'
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Build and push Docker image
run: |
docker build -t churn-api:${{ github.sha }} serving/
docker tag churn-api:${{ github.sha }} \
${{ secrets.ECRREGISTRY }}/churn-api:latest
docker push ${{ secrets.ECRREGISTRY }}/churn-api:latest
- name: Deploy to ECS
run: |
aws ecs update-service \
--cluster ml-production \
--service churn-api \
--force-new-deployment
Testing Workflow
# .github/workflows/test.yml
name: ML Tests
on:
pullrequest:
branches: [main]
jobs:
test:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Set up Python
uses: actions/setup-python@v5
with:
python-version: '3.11'
- name: Install dependencies
run: pip install -r requirements.txt
- name: Run unit tests
run: pytest tests/testdata.py tests/testmodel.py -v
- name: Run API tests
run: pytest tests/testapi.py -v
- name: Lint code
run: |
pip install ruff
ruff check src/
Containerization with Docker
# serving/Dockerfile
FROM python:3.11-slim
WORKDIR /app
Install dependencies
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
Copy application code
COPY app.py .
COPY models/ models/
Health check
HEALTHCHECK --interval=30s --timeout=10s --retries=3 \
CMD curl -f http://localhost:8000/health || exit 1
EXPOSE 8000
CMD ["uvicorn", "app:app", "--host", "0.0.0.0", "--port", "8000", \
"--workers", "4"]
Serving Application
# serving/app.py
from fastapi import FastAPI, HTTPException
from pydantic import BaseModel
import joblib
import numpy as np
import pandas as pd
import mlflow
import logging
import time
from datetime import datetime
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(name)
app = FastAPI(title="Churn Prediction API", version="1.0.0")
Load model at startup
model = None
modelversion = None
@app.onevent("startup")
async def loadmodel():
global model, modelversion
try:
# Try loading from MLflow registry first
mlflowuri = os.environ.get("MLFLOWTRACKINGURI")
if mlflowuri:
mlflow.settrackinguri(mlflowuri)
model = mlflow.pyfunc.loadmodel(
"models:/churn-predictor/Production"
)
modelversion = "mlflow-production"
else:
model = joblib.load("models/model.pkl")
modelversion = "local"
logger.info(f"Model loaded: {modelversion}")
except Exception as e:
logger.error(f"Failed to load model: {e}")
raise
class PredictionRequest(BaseModel):
tenure: float
monthlycharges: float
totalcharges: float
contracttype: str
paymentmethod: str
internetservice: str
class PredictionResponse(BaseModel):
churnprobability: float
churnprediction: bool
modelversion: str
predictiontimems: float
@app.get("/health")
async def health():
return {
"status": "healthy",
"modelloaded": model is not None,
"modelversion": modelversion,
"timestamp": datetime.utcnow().isoformat()
}
@app.post("/predict", responsemodel=PredictionResponse)
async def predict(request: PredictionRequest):
if model is None:
raise HTTPException(statuscode=503, detail="Model not loaded")
starttime = time.time()
try:
# Create feature DataFrame
features = pd.DataFrame([request.dict()])
# Get prediction
if hasattr(model, "predictproba"):
probability = float(model.predictproba(features)[0, 1])
else:
probability = float(model.predict(features)[0])
prediction = probability >= 0.5
elapsedms = (time.time() - starttime) 1000
# Log prediction for monitoring
logger.info(
f"Prediction: prob={probability:.4f}, "
f"churn={prediction}, time={elapsedms:.1f}ms"
)
return PredictionResponse(
churnprobability=round(probability, 4),
churnprediction=prediction,
modelversion=modelversion,
predictiontimems=round(elapsedms, 2)
)
except Exception as e:
logger.error(f"Prediction error: {e}")
raise HTTPException(statuscode=500, detail=str(e))
@app.post("/predict/batch")
async def predictbatch(requests: list[PredictionRequest]):
"""Batch prediction endpoint for multiple customers."""
if model is None:
raise HTTPException(statuscode=503, detail="Model not loaded")
features = pd.DataFrame([r.dict() for r in requests])
if hasattr(model, "predictproba"):
probabilities = model.predictproba(features)[:, 1]
else:
probabilities = model.predict(features)
return {
"predictions": [
{
"churnprobability": round(float(p), 4),
"churnprediction": bool(p >= 0.5)
}
for p in probabilities
],
"count": len(requests),
"modelversion": modelversion
}
Docker Compose for Local Development
# docker-compose.yml
version: '3.8'
services:
mlflow:
image: ghcr.io/mlflow/mlflow:v2.10.0
ports:
- "5000:5000"
command: >
mlflow server
--backend-store-uri sqlite:///mlflow.db
--default-artifact-root /mlruns
--host 0.0.0.0
volumes:
- mlflow-data:/mlruns
api:
build: ./serving
ports:
- "8000:8000"
environment:
- MLFLOWTRACKINGURI=http://mlflow:5000
dependson:
- mlflow
monitoring:
build: ./monitoring
ports:
- "8501:8501"
environment:
- API
URL=http://api:8000
dependson:
- api
volumes:
mlflow-data:
Deployment
Production Deployment Script
# scripts/deploy.py
import boto3
import json
import subprocess
import sys
def deploytoproduction(imagetag: str, cluster: str = "ml-production",
service: str = "churn-api"):
"""Deploy new model version to production."""
# Step 1: Run pre-deployment tests
print("Running pre-deployment tests...")
result = subprocess.run(
["pytest", "tests/", "-v", "--tb=short"],
captureoutput=True, text=True
)
if result.returncode != 0:
print(f"Tests failed:\n{result.stdout}")
sys.exit(1)
# Step 2: Validate model quality
print("Validating model quality...")
with open("metrics/evalmetrics.json") as f:
metrics = json.load(f)
thresholds = {"aucroc": 0.85, "f1": 0.75, "precision": 0.70}
for metric, threshold in thresholds.items():
if metrics.get(metric, 0) < threshold:
print(f"FAILED: {metric}={metrics[metric]:.4f} < {threshold}")
sys.exit(1)
print("All quality gates passed")
# Step 3: Update ECS service
print(f"Deploying {imagetag} to {cluster}/{service}...")
ecs = boto3.client("ecs")
ecs.updateservice(
cluster=cluster,
service=service,
forceNewDeployment=True
)
print("Deployment initiated. Monitor at AWS Console.")
if name == "main":
deploytoproduction(imagetag=sys.argv[1] if len(sys.argv) > 1
else "latest")
Monitoring with Evidently
Evidently detects data drift and model performance degradation in production.
# src/monitoring/driftdetection.py
import pandas as pd
import numpy as np
from evidently.report import Report
from evidently.metric
preset import (
DataDriftPreset,
TargetDriftPreset,
ClassificationPreset
)
from evidently.metrics import (
DataDriftTable,
DatasetDriftMetric,
ColumnDriftMetric
)
from evidently.testsuite import TestSuite
from evidently.tests import (
TestShareOfDriftedColumns,
TestColumnDrift,
TestShareOfMissingValues
)
import json
from datetime import datetime
import logging
logger = logging.getLogger(name)
class ModelMonitor:
"""Monitor model performance and data drift in production."""
def init(self, referencedata: pd.DataFrame,
featurecolumns: list, targetcolumn: str = "churn"):
self.referencedata = referencedata
self.featurecolumns = featurecolumns
self.targetcolumn = targetcolumn
def checkdatadrift(self, currentdata: pd.DataFrame,
driftthreshold: float = 0.3) -> dict:
"""Check for data drift between reference and current data."""
report = Report(metrics=[
DatasetDriftMetric(),
DataDriftTable(),
])
report.run(
referencedata=self.referencedata[self.featurecolumns],
currentdata=currentdata[self.featurecolumns]
)
result = report.asdict()
datasetdrift = result["metrics"][0]["result"]
driftreport = {
"timestamp": datetime.utcnow().isoformat(),
"datasetdriftdetected": datasetdrift["datasetdrift"],
"driftshare": datasetdrift["shareofdriftedcolumns"],
"numberofcolumns": datasetdrift["numberofcolumns"],
"numberofdriftedcolumns":
datasetdrift["numberofdriftedcolumns"],
"driftedcolumns": [],
}
# Identify which columns drifted
columnresults = result["metrics"][1]["result"]["driftbycolumns"]
for colname, coldata in columnresults.items():
if coldata["driftdetected"]:
driftreport["driftedcolumns"].append({
"column": colname,
"driftscore": coldata["driftscore"],
"stattestname": coldata["stattestname"],
})
iscritical = (
driftreport["driftshare"] > driftthreshold
)
driftreport["actionrequired"] = iscritical
return driftreport
def runtestsuite(self, currentdata: pd.DataFrame) -> dict:
"""Run automated quality tests on current data."""
suite = TestSuite(tests=[
TestShareOfDriftedColumns(lt=0.3),
TestShareOfMissingValues(lt=0.05),
])
# Add per-column drift tests for critical features
for col in ["monthlycharges", "tenure", "totalcharges"]:
if col in self.featurecolumns:
suite.tests.append(TestColumnDrift(columnname=col))
suite.run(
referencedata=self.referencedata,
currentdata=currentdata
)
results = suite.asdict()
return {
"timestamp": datetime.utcnow().isoformat(),
"summary": results["summary"],
"allpassed": results["summary"]["allpassed"],
"tests": [
{
"name": t["name"],
"status": t["status"],
"description": t.get("description", ""),
}
for t in results["tests"]
]
}
def generatemonitoringreport(self, currentdata: pd.DataFrame,
outputpath: str = "monitoringreport.html"):
"""Generate a comprehensive HTML monitoring report."""
report = Report(metrics=[
DataDriftPreset(),
TargetDriftPreset(),
])
report.run(
referencedata=self.referencedata,
currentdata=currentdata
)
report.savehtml(outputpath)
logger.info(f"Monitoring report saved to {outputpath}")
Usage
if name == "main":
# Load reference data (training data)
reference = pd.readcsv("data/processed/trainfeatures.csv")
# Simulate production data
current = pd.readcsv("data/production/latestbatch.csv")
featurecols = [
"tenure", "monthlycharges", "totalcharges",
"contracttype", "paymentmethod", "internetservice"
]
monitor = ModelMonitor(reference, featurecols)
# Check drift
drift = monitor.checkdatadrift(current)
print(f"Drift detected: {drift['datasetdriftdetected']}")
print(f"Drifted columns: {len(drift['driftedcolumns'])}")
if drift["actionrequired"]:
print("ACTION REQUIRED: Significant data drift detected!")
# Run quality tests
testresults = monitor.runtestsuite(current)
print(f"Tests passed: {testresults['allpassed']}")
Alerting
# src/monitoring/alerting.py
import requests
import json
import logging
from datetime import datetime
logger = logging.getLogger(name)
class AlertManager:
"""Send alerts when monitoring detects issues."""
def init(self, slackwebhookurl: str = None,
pagerdutykey: str = None):
self.slackwebhook = slackwebhookurl
self.pagerdutykey = pagerdutykey
def sendslackalert(self, title: str, message: str,
severity: str = "warning"):
"""Send alert to Slack channel."""
if not self.slackwebhook:
logger.warning("Slack webhook not configured")
return
colormap = {
"info": "#36a64f",
"warning": "#ff9900",
"critical": "#ff0000"
}
payload = {
"attachments": [{
"color": colormap.get(severity, "#ff9900"),
"title": f"[ML Alert] {title}",
"text": message,
"fields": [
{"title": "Severity", "value": severity, "short": True},
{"title": "Time", "value": datetime.utcnow().isoformat(),
"short": True},
],
"footer": "MLOps Monitoring System"
}]
}
response = requests.post(self.slackwebhook, json=payload)
if response.statuscode == 200:
logger.info(f"Slack alert sent: {title}")
else:
logger.error(f"Slack alert failed: {response.text}")
def checkandalert(self, driftreport: dict,
testresults: dict):
"""Evaluate monitoring results and send alerts."""
# Data drift alert
if driftreport.get("actionrequired"):
drifted = ", ".join(
[d["column"] for d in driftreport["driftedcolumns"]]
)
self.sendslackalert(
title="Data Drift Detected",
message=(
f"Significant data drift detected in production.\n"
f"Drifted columns ({driftreport['driftshare']:.0%}): "
f"{drifted}\n"
f"Consider retraining the model."
),
severity="critical"
)
# Test failure alert
if not testresults.get("allpassed", True):
failed = [t["name"] for t in testresults.get("tests", [])
if t["status"] != "SUCCESS"]
self.sendslackalert(
title="Quality Tests Failed",
message=(
f"Production data quality tests failed.\n"
f"Failed tests: {', '.join(failed)}"
),
severity="warning"
)
Scheduled monitoring job
def runmonitoringjob():
"""Run the full monitoring pipeline (call from cron or scheduler)."""
from driftdetection import ModelMonitor
import pandas as pd
reference = pd.readcsv("data/processed/trainfeatures.csv")
current = pd.readcsv("data/production/latestbatch.csv")
featurecols = [
"tenure", "monthlycharges", "totalcharges",
"contracttype", "paymentmethod", "internetservice"
]
monitor = ModelMonitor(reference, featurecols)
alerter = AlertManager(
slackwebhookurl=os.environ.get("SLACKWEBHOOKURL")
)
drift = monitor.checkdatadrift(current)
tests = monitor.runtestsuite(current)
alerter.checkandalert(drift, tests)
# Save monitoring results
with open("metrics/monitoringlatest.json", "w") as f:
json.dump({"drift": drift, "tests": tests}, f, indent=2)