id: kpi-changepoint-release-attribution
namespace: company.team
inputs:
- id: scenario
type: SELECT
displayName: Demo scenario
description: >
TRACKING_BREAK: a release on day 41 stops logging most Android orders, a
lasting drop. PROMO_SPIKE: a two-day campaign that returns to normal,
which must not be flagged. STABLE: weekly seasonality and noise only.
values:
- TRACKING_BREAK
- PROMO_SPIKE
- STABLE
defaults: TRACKING_BREAK
- id: min_shift_pct
type: FLOAT
displayName: Minimum lasting shift (%)
description: Ignore changepoints whose before/after levels differ by less than this.
defaults: 15.0
- id: min_segment_days
type: INT
displayName: Minimum segment length (days)
description: A level must hold at least this long to count, so one-off spikes
are not changepoints.
defaults: 5
- id: recent_days
type: INT
displayName: Alert on shifts in the last N days
defaults: 21
tasks:
- id: detect
type: io.kestra.plugin.scripts.python.Script
description: >
Remove the weekly pattern from 90 days of a daily KPI per segment, find
level shifts with ruptures (PELT), keep shifts that are large and lasting,
then name the release deployed closest before each shift.
containerImage: python:3.12-slim
beforeCommands:
- pip install -q --root-user-action=ignore ruptures==1.1.10 numpy==2.3.3
env:
SCENARIO: "{{ inputs.scenario }}"
MIN_SHIFT: "{{ inputs.min_shift_pct }}"
MIN_SEG: "{{ inputs.min_segment_days }}"
RECENT: "{{ inputs.recent_days }}"
outputFiles:
- changepoints.md
script: |
import json, os
from datetime import date, timedelta
import numpy as np
import ruptures as rpt
rng = np.random.default_rng(3)
scenario = os.environ["SCENARIO"]
min_shift, min_seg, recent = float(os.environ["MIN_SHIFT"]), int(os.environ["MIN_SEG"]), int(os.environ["RECENT"])
start, days = date(2026, 7, 8), 90
dates = [start + timedelta(d) for d in range(days)]
# Demo KPI: daily orders per channel with a weekly pattern. Replace with a warehouse query.
weekly = np.array([1.0, 0.97, 0.98, 1.02, 1.12, 1.25, 1.18])
base = {"web": 4200, "ios": 2600, "android": 2300}
series = {}
for ch, level in base.items():
y = level * weekly[[d.weekday() for d in dates]] * rng.normal(1, 0.04, days)
if scenario == "TRACKING_BREAK" and ch == "android":
y[71:] *= 0.58
if scenario == "PROMO_SPIKE" and ch == "web":
y[75:77] *= 1.9
series[ch] = y
releases = {date(2026, 8, 3): "v2.14 checkout copy", date(2026, 9, 1): "v2.15 payments SDK", date(2026, 9, 17): "v2.16 Android analytics SDK upgrade", date(2026, 9, 24): "v2.17 search ranking"}
findings, rows = [], []
for ch, y in series.items():
# Divide out the weekday pattern so Saturdays are not a "shift".
factors = np.array([np.median(y[[i for i in range(days) if dates[i].weekday() == w]]) for w in range(7)])
deseason = y / (factors[[d.weekday() for d in dates]] / factors.mean())
signal = deseason / deseason.mean()
bkps = rpt.Pelt(model="l2", min_size=min_seg).fit(signal.reshape(-1, 1)).predict(pen=0.05 * days * signal.var() + 0.5)
edges = [0] + bkps
for i in range(1, len(edges) - 1):
cp = edges[i]
before, after = deseason[edges[i - 1]:cp].mean(), deseason[cp:edges[i + 1]].mean()
shift = (after - before) / before * 100
if abs(shift) < min_shift:
continue
when = dates[cp]
# A changepoint estimate can land a day early, so look 3 days before to 1 day after.
near = sorted((abs((when - d).days), d) for d in releases if -1 <= (when - d).days <= 3)
cause = f"{releases[near[0][1]]} on {near[0][1]}" if near else "no release from 3 days before to 1 day after"
f = {"segment": ch, "date": str(when), "shift_pct": round(shift, 1), "before": round(before), "after": round(after),
"lasting_days": edges[i + 1] - cp, "release": cause, "recent": (dates[-1] - when).days < recent}
findings.append(f)
rows.append(f"| {ch} | {round(deseason[:30].mean())} | {round(deseason[-14:].mean())} |")
alerts = [f for f in findings if f["recent"]]
lines = [f"# KPI changepoints ({scenario})", "", "| segment | level first 30 days | level last 14 days |", "|---|---|---|", *rows, "",
"## Lasting shifts", *([f"- {f['segment']}: {f['shift_pct']:+}% from {f['date']} ({f['before']} to {f['after']} orders/day, held {f['lasting_days']} days). Release: {f['release']}" for f in findings] or ["- none"])]
open("changepoints.md", "w").write("\n".join(lines) + "\n")
print("\n".join(lines))
print("::" + json.dumps({"outputs": {"shifts": len(findings), "alerts": len(alerts),
"detail": "; ".join(f"{f['segment']} {f['shift_pct']:+}% from {f['date']} ({f['release']})" for f in alerts)}}) + "::")
- id: publish_report
type: io.kestra.plugin.core.namespace.UploadFiles
namespace: "{{ flow.namespace }}"
filesMap:
"kpi/changepoints-{{ execution.startDate | date('yyyy-MM-dd') }}.md": "{{ outputs.detect.outputFiles['changepoints.md'] }}"
- id: alert_gate
type: io.kestra.plugin.core.flow.If
condition: "{{ outputs.detect.vars.alerts > 0 }}"
then:
- id: lasting_shift
type: io.kestra.plugin.core.execution.Fail
description: A lasting shift right after a release is usually broken tracking or
a real regression. Either way someone has to look.
errorMessage: "Lasting KPI shift in the last {{ inputs.recent_days }} days: {{
outputs.detect.vars.detail }}"
else:
- id: no_shift
type: io.kestra.plugin.core.log.Log
message: "No lasting shift of {{ inputs.min_shift_pct }}% or more in the last {{
inputs.recent_days }} days ({{ outputs.detect.vars.shifts }} older
shift(s) in the report)."
triggers:
- id: daily
type: io.kestra.plugin.core.trigger.Schedule
cron: "30 7 * * *"