id: evidently-drift-champion-challenger-retrain
namespace: company.team
inputs:
- id: source_url
type: URI
displayName: Training data
description: CSV with a numeric target column. The default is a public data
science salaries dataset of about 3,000 rows.
defaults: https://huggingface.co/datasets/kestra/datasets/raw/main/csv/salaries.csv
- id: simulate_drift
type: BOOL
displayName: Simulate drift
description: Raise salaries by 60% and make every job remote in the new batch,
as a market shift would. Watch the flow detect it, retrain, then promote
the challenger.
defaults: false
- id: drift_share_threshold
type: FLOAT
displayName: Drift share threshold
description: Share of monitored columns whose distribution changed (p-value
under 0.05) before the model is considered stale.
defaults: 0.3
- id: min_improvement_pct
type: FLOAT
displayName: Minimum improvement (%)
description: How much lower the challenger's mean absolute error must be, on
held-out data from the new batch, before it replaces the champion.
defaults: 5.0
- id: notify_slack
type: BOOL
displayName: Notify Slack
description: Post the outcome to Slack. Needs the SLACK_WEBHOOK_URL secret.
defaults: false
variables:
model_dir: ml/salary
champion_kv: salary_model_champion
concurrency:
limit: 1
triggers:
- id: weekly
type: io.kestra.plugin.core.trigger.Schedule
description: Check the model against the latest batch every Monday.
cron: "0 6 * * 1"
tasks:
- id: download_batch
type: io.kestra.plugin.core.http.Download
uri: "{{ inputs.source_url }}"
- id: evaluate
type: io.kestra.plugin.scripts.python.Script
description: >
Compare the new batch with the reference window using Evidently, then
decide. No champion yet: train one (BOOTSTRAP). No meaningful drift: keep
it (STABLE). Drift: train a challenger on the new batch and score both
models on held-out rows of that batch. PROMOTE when the challenger is
clearly better, HOLD when it is not.
containerImage: python:3.12-slim
beforeCommands:
- pip install -q --root-user-action=ignore evidently==0.7.23
scikit-learn==1.7.2
namespaceFiles:
enabled: true
include:
- "{{ vars.model_dir }}/champion.joblib"
- "{{ vars.model_dir }}/reference.csv"
inputFiles:
data.csv: "{{ outputs.download_batch.uri }}"
outputFiles:
- challenger.joblib
- previous_champion.joblib
- new_reference.csv
- drift_report.html
env:
SIMULATE_DRIFT: "{{ inputs.simulate_drift }}"
DRIFT_THRESHOLD: "{{ inputs.drift_share_threshold }}"
MIN_IMPROVEMENT: "{{ inputs.min_improvement_pct }}"
CHAMPION_PATH: "{{ vars.model_dir }}/champion.joblib"
script: |
import json, os, warnings
import joblib, numpy as np, pandas as pd
from evidently import DataDefinition, Dataset, Report
from evidently.presets import DataDriftPreset
from sklearn.compose import ColumnTransformer
from sklearn.ensemble import GradientBoostingRegressor
from sklearn.metrics import mean_absolute_error
from sklearn.model_selection import train_test_split
from sklearn.pipeline import Pipeline
from sklearn.preprocessing import OneHotEncoder
warnings.filterwarnings("ignore")
TARGET = "salary_in_usd"
CATS = ["experience_level", "employment_type", "company_size"]
NUMS = ["remote_ratio", "work_year"]
df = pd.read_csv("data.csv")[CATS + NUMS + [TARGET]].dropna()
# The reference is the data the current champion was trained on. Before the first champion
# exists, the first half of a seeded shuffle plays that role and the second half is the new batch.
order = np.random.default_rng(7).permutation(len(df))
batch = df.iloc[order[len(df) // 2 :]].reset_index(drop=True)
reference_path = os.path.join(os.path.dirname(os.environ["CHAMPION_PATH"]), "reference.csv")
if os.path.exists(reference_path):
reference = pd.read_csv(reference_path)
else:
reference = df.iloc[order[: len(df) // 2]].reset_index(drop=True)
if os.environ["SIMULATE_DRIFT"] == "true":
batch[TARGET] = batch[TARGET] * 1.6
batch["remote_ratio"] = 100
# Fixed tests, so the verdict cannot change because the row count crossed a library default.
definition = DataDefinition(numerical_columns=NUMS + [TARGET], categorical_columns=CATS)
snapshot = Report([DataDriftPreset(num_method="ks", cat_method="chisquare")]).run(
Dataset.from_pandas(batch, data_definition=definition),
Dataset.from_pandas(reference, data_definition=definition),
)
snapshot.save_html("drift_report.html")
metrics = snapshot.dict()["metrics"]
share = next(m["value"]["share"] for m in metrics if m["metric_name"].startswith("DriftedColumnsCount"))
drifted = sorted(m["metric_name"].split("column=")[1].split(",")[0]
for m in metrics if m["metric_name"].startswith("ValueDrift") and m["value"] < 0.05)
def new_model():
return Pipeline([
("encode", ColumnTransformer([("cat", OneHotEncoder(handle_unknown="ignore"), CATS)], remainder="passthrough")),
("model", GradientBoostingRegressor(random_state=0)),
])
champion_path = os.environ["CHAMPION_PATH"]
has_champion = os.path.exists(champion_path)
# Hand the current champion back out, so a promotion can archive it.
open("previous_champion.joblib", "wb").write(open(champion_path, "rb").read() if has_champion else b"")
train, holdout = train_test_split(batch, test_size=0.3, random_state=0)
X, y = holdout[CATS + NUMS], holdout[TARGET]
champion_mae = challenger_mae = 0.0
if not has_champion:
verdict = "BOOTSTRAP"
challenger = new_model().fit(reference[CATS + NUMS], reference[TARGET])
challenger_mae = mean_absolute_error(y, challenger.predict(X))
joblib.dump(challenger, "challenger.joblib")
else:
champion = joblib.load(champion_path)
champion_mae = mean_absolute_error(y, champion.predict(X))
if share < float(os.environ["DRIFT_THRESHOLD"]):
verdict = "STABLE"
open("challenger.joblib", "wb").close()
else:
challenger = new_model().fit(train[CATS + NUMS], train[TARGET])
challenger_mae = mean_absolute_error(y, challenger.predict(X))
joblib.dump(challenger, "challenger.joblib")
improvement = (champion_mae - challenger_mae) / champion_mae * 100
verdict = "PROMOTE" if improvement >= float(os.environ["MIN_IMPROVEMENT"]) else "HOLD"
# The data the promoted model learned from becomes the next reference window.
(reference if verdict == "BOOTSTRAP" else train).to_csv("new_reference.csv", index=False)
improvement = round((champion_mae - challenger_mae) / champion_mae * 100, 1) if champion_mae and challenger_mae else 0.0
out = {
"verdict": verdict,
"drift_share": round(share, 3),
"drifted_columns": ", ".join(drifted) or "none",
"champion_mae": round(champion_mae),
"challenger_mae": round(challenger_mae),
"improvement_pct": improvement,
"reference_rows": len(reference),
"batch_rows": len(batch),
}
print(json.dumps(out, indent=2))
print("::" + json.dumps({"outputs": out}) + "::")
- id: publish_drift_report
type: io.kestra.plugin.core.namespace.UploadFiles
description: Keep every Evidently report, so the history of the model's inputs
stays browsable.
namespace: "{{ flow.namespace }}"
filesMap:
"{{ vars.model_dir }}/reports/drift-{{ execution.startDate | date('yyyy-MM-dd-HHmmss') }}.html": "{{ outputs.evaluate.outputFiles['drift_report.html'] }}"
- id: decide
type: io.kestra.plugin.core.flow.Switch
value: "{{ outputs.evaluate.vars.verdict }}"
cases:
STABLE:
- id: keep_champion
type: io.kestra.plugin.core.log.Log
message: "STABLE: drift share {{ outputs.evaluate.vars.drift_share }} is under
{{ inputs.drift_share_threshold }}. The champion stays, with MAE {{
outputs.evaluate.vars.champion_mae }} on the new batch."
BOOTSTRAP:
- id: install_first_champion
type: io.kestra.plugin.core.namespace.UploadFiles
namespace: "{{ flow.namespace }}"
filesMap:
"{{ vars.model_dir }}/champion.joblib": "{{ outputs.evaluate.outputFiles['challenger.joblib'] }}"
"{{ vars.model_dir }}/reference.csv": "{{ outputs.evaluate.outputFiles['new_reference.csv'] }}"
- id: record_first_champion
type: io.kestra.plugin.core.kv.Set
key: "{{ vars.champion_kv }}"
kvType: JSON
value: |
{"version": 1, "mae": {{ outputs.evaluate.vars.challenger_mae }}, "reason": "bootstrap", "execution_id": "{{ execution.id }}"}
PROMOTE:
- id: archive_old_champion
type: io.kestra.plugin.core.namespace.UploadFiles
description: Keep the replaced model, so a promotion can be rolled back by
copying it back.
namespace: "{{ flow.namespace }}"
filesMap:
"{{ vars.model_dir }}/archive/champion-v{{ kv(vars.champion_kv).version }}.joblib": "{{ outputs.evaluate.outputFiles['previous_champion.joblib'] }}"
- id: promote_challenger
type: io.kestra.plugin.core.namespace.UploadFiles
namespace: "{{ flow.namespace }}"
description: The challenger becomes the champion, and the data it learned from
becomes the reference window for the next drift check.
filesMap:
"{{ vars.model_dir }}/champion.joblib": "{{ outputs.evaluate.outputFiles['challenger.joblib'] }}"
"{{ vars.model_dir }}/reference.csv": "{{ outputs.evaluate.outputFiles['new_reference.csv'] }}"
- id: record_promotion
type: io.kestra.plugin.core.kv.Set
key: "{{ vars.champion_kv }}"
kvType: JSON
value: |
{"version": {{ kv(vars.champion_kv).version + 1 }}, "mae": {{ outputs.evaluate.vars.challenger_mae }}, "reason": "drift on {{ outputs.evaluate.vars.drifted_columns }}", "replaced_mae": {{ outputs.evaluate.vars.champion_mae }}, "execution_id": "{{ execution.id }}"}
defaults:
- id: hold_for_review
type: io.kestra.plugin.core.execution.Fail
description: The data moved but retraining did not help enough. A person has to
look before the model is trusted on this data.
errorMessage: "HOLD: drift share {{ outputs.evaluate.vars.drift_share }} on {{
outputs.evaluate.vars.drifted_columns }}, but the challenger MAE {{
outputs.evaluate.vars.challenger_mae }} is only {{
outputs.evaluate.vars.improvement_pct }}% better than the champion {{
outputs.evaluate.vars.champion_mae }} (need {{
inputs.min_improvement_pct }}%). Champion kept."
- id: log_outcome
type: io.kestra.plugin.core.log.Log
message: >-
{{ outputs.evaluate.vars.verdict }}. Drift share {{
outputs.evaluate.vars.drift_share }} (drifted: {{
outputs.evaluate.vars.drifted_columns }}). Champion MAE {{
outputs.evaluate.vars.champion_mae }}, challenger MAE {{
outputs.evaluate.vars.challenger_mae }}. Current champion: {{
kv(vars.champion_kv, errorOnMissing=false) | toJson }}
- id: notify
type: io.kestra.plugin.core.flow.If
condition: "{{ inputs.notify_slack and outputs.evaluate.vars.verdict != 'STABLE' }}"
then:
- id: slack_outcome
type: io.kestra.plugin.slack.notifications.SlackIncomingWebhook
url: "{{ secret('SLACK_WEBHOOK_URL') }}"
payload: |
{
"text": "Salary model {{ outputs.evaluate.vars.verdict }}: drift share {{ outputs.evaluate.vars.drift_share }} on {{ outputs.evaluate.vars.drifted_columns }}, champion MAE {{ outputs.evaluate.vars.champion_mae }}, challenger MAE {{ outputs.evaluate.vars.challenger_mae }}. Execution {{ execution.id }}."
}