Course Content
AI Monitoring and Observability
3 sections · 7 lessons
Mini Project: Observability Dashboard for LLM API
You are going to build the thing that would have caught the failure nobody catches. A mock LLM API serves 6,000 requests. Somewhere around request 3,000, without warning you in the code that consumes it, the input distribution shifts: prompts get longer, the retrieval step starts returning weaker matches, and answer quality degrades. Latency barely moves. The error rate stays at zero. Every request returns 200.
By the end you will have a system that detects that shift, attributes it to a specific feature, decides whether it warrants a page or a ticket, routes it accordingly, and refuses to retrain on the corrupted data. All of it runs locally with no external services — you will mock Prometheus scraping, Grafana panels, Slack, PagerDuty, and a retraining pipeline, because the point is the logic, not the vendors.
+---------------------------+ request ------------>| mock LLM API service | | retrieve -> generate | | -> postprocess | +---------------------------+ | | | structured | | metrics | OTel spans JSON logs v v v +-------+ +-------+ +--------+ | logs | |registry| |exporter| +---+---+ +---+---+ +----+---+ | | | v v v +------------------------------------+ | detectors: PSI drift | robust-z | | anomaly | SLO burn | +------------------+-----------------+ | +------------------v-----------------+ | alert engine (persistence rules) | +---+------------------------+-------+ | | dashboard notifier -> retrain gateSetup
1mkdir llm-observability && cd llm-observability2python -m venv .venv && source .venv/bin/activate3pip install numpy prometheus-client opentelemetry-api opentelemetry-sdk45# obs/service.py mock API with logging + metrics6# obs/tracing.py tracer setup and instrumented pipeline7# obs/detectors.py PSI drift + robust z-score anomaly8# obs/alerts.py rules, persistence, notifiers9# obs/dashboard.py terminal panels10# obs/run.py the simulation that ties it togetherPart 1 — Structured logs and metrics on every request
Build the service so that instrumentation cannot be skipped. The whole thing lives in a finally block, because a service that stops emitting metrics when it starts failing is blind at exactly the wrong moment.
1# obs/service.py2import json, time, uuid, random, math3from prometheus_client import Counter, Histogram, CollectorRegistry45REGISTRY = CollectorRegistry()6REQUESTS = Counter("llm_requests_total", "requests",7 ["model", "status"], registry=REGISTRY)8LATENCY = Histogram("llm_latency_seconds", "end-to-end latency", ["model"],9 buckets=(.1,.25,.5,.75,1,1.5,2,3,4,6,8,12), registry=REGISTRY)10TOKENS = Histogram("llm_tokens_total", "tokens", ["model", "direction"],11 buckets=(64,128,256,512,1024,2048,4096), registry=REGISTRY)12QUALITY = Histogram("llm_answer_quality", "0-1 quality score", ["model"],13 buckets=(.1,.3,.5,.6,.7,.8,.9,.95), registry=REGISTRY)1415MODEL = "assistant-v4"16EVENTS = [] # stands in for a log store1718def log_event(**fields):19 fields["ts"] = time.time()20 EVENTS.append(fields)21 print(json.dumps(fields))2223def handle_request(prompt_tokens: int, drifted: bool = False) -> dict:24 """One request. Latency is deliberately dominated by generation."""25 request_id = str(uuid.uuid4())26 started = time.perf_counter()27 status = "ok"28 try:29 out_tokens = max(40, int(random.gauss(300, 60)))30 # weaker retrieval after the shift -> lower quality, longer answers31 top_score = random.gauss(0.55 if drifted else 0.78, 0.08)32 quality = min(1.0, max(0.0, random.gauss(33 0.62 if drifted else 0.86, 0.07)))34 # generation time is roughly linear in output tokens35 latency = 0.04 + 0.18 + out_tokens * 0.0094 + random.gauss(0, 0.05)36 latency = max(0.05, latency)37 time.sleep(0) # no real waiting in the simulation38 return dict(request_id=request_id, input_tokens=prompt_tokens,39 output_tokens=out_tokens, top_score=top_score,40 quality=quality, latency=latency)41 except Exception:42 status = "error"43 raise44 finally:45 elapsed = time.perf_counter() - started46 REQUESTS.labels(MODEL, status).inc()47 LATENCY.labels(MODEL).observe(locals().get("latency", elapsed))48 if status == "ok":49 TOKENS.labels(MODEL, "in").observe(prompt_tokens)50 TOKENS.labels(MODEL, "out").observe(out_tokens)51 QUALITY.labels(MODEL).observe(quality)52 log_event(event="inference_complete", request_id=request_id,53 model=MODEL, model_version="4.2.1",54 input_tokens=prompt_tokens, output_tokens=out_tokens,55 top_score=round(top_score, 3), quality=round(quality, 3),56 latency_ms=round(latency * 1000, 1))Two decisions to notice. Labels are bounded — model and status only. Adding a per-user label here would multiply the stored series by the user count, which is how metrics backends run out of memory. Anything high-cardinality goes on the log event, where storage grows with the number of events rather than with the product of label values.
The histogram buckets are chosen, not defaulted. Quantiles are interpolated inside whichever bucket contains the answer, so your worst-case error is the bucket's width. With default boundaries jumping 1 → 2.5 → 5 seconds, a true p95 of 2.6 s could be reported anywhere in that 2.5-second-wide band. The boundaries above are dense where this service's traffic actually sits.
Verify by hand before going further — instrumentation bugs are invisible later:
1from prometheus_client import generate_latest2for _ in range(50):3 handle_request(prompt_tokens=600)4print(generate_latest(REGISTRY).decode())5# expect llm_requests_total{model="assistant-v4",status="ok"} 50.06# and a bucket ladder whose +Inf count is also 50Part 2 — Tracing the pipeline and finding the bottleneck
Metrics tell you a request took 3.9 seconds. Only a trace tells you where the 3.9 seconds went. Use an in-memory exporter so the spans are inspectable in the same process.
1# obs/tracing.py2from opentelemetry import trace3from opentelemetry.sdk.resources import Resource4from opentelemetry.sdk.trace import TracerProvider5from opentelemetry.sdk.trace.export import SimpleSpanProcessor6from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter78exporter = InMemorySpanExporter()9provider = TracerProvider(resource=Resource.create({10 "service.name": "llm-api", "service.version": "4.2.1"}))11provider.add_span_processor(SimpleSpanProcessor(exporter))12trace.set_tracer_provider(provider)13tracer = trace.get_tracer("llm-api")1415def traced_request(prompt_tokens: int, drifted: bool = False):16 with tracer.start_as_current_span("chat.answer") as root:17 root.set_attribute("gen_ai.usage.input_tokens", prompt_tokens)1819 with tracer.start_as_current_span("retrieval.search") as s:20 s.set_attribute("retrieval.k", 5)21 r = handle_request(prompt_tokens, drifted)22 s.set_attribute("retrieval.top_score", round(r["top_score"], 3))2324 with tracer.start_as_current_span("llm.generate") as s:25 s.set_attribute("gen_ai.request.model", MODEL)26 s.set_attribute("gen_ai.usage.output_tokens", r["output_tokens"])2728 with tracer.start_as_current_span("postprocess"):29 pass3031 root.set_attribute("gen_ai.usage.output_tokens", r["output_tokens"])32 root.set_attribute("answer.quality", round(r["quality"], 3))33 return rSimpleSpanProcessor is correct here and wrong in production — it exports synchronously on every span end, adding a network round trip to every operation. Production uses BatchSpanProcessor, which exports in the background from a queue of 2,048 spans by default. When the Collector stalls, that queue fills in seconds at peak traffic and new spans are dropped with only a log warning. Size the queue for your peak before you ship anything.
Now the analysis that justifies the whole exercise. A representative trace decomposes like this:
chat.answer 4,200 ms retrieval.search 40 ms 1.0% llm.generate (402 output tokens) 3,780 ms 90.0% postprocess 25 ms 0.6% sum of children = 3,845 ms unaccounted = 355 ms 8.5% <- missing instrumentationTwo readings. The unaccounted 355 ms is the gap between the root duration and the sum of its children; a gap that large means real work is happening in an uninstrumented region and you should go find it before optimising anything. And the attribution decides what to work on, usually against intuition:
| Optimisation | New total | Improvement |
|---|---|---|
| Halve retrieval: 40 → 20 ms | 4,180 ms | 0.5% |
| Cut output tokens 402 → 250 with a "be concise" instruction | 2,771 ms | 34.0% |
Generation time is roughly linear in output tokens, so 3,780 × (250/402) = 2,351 ms, saving 1,429 ms of a 4,200 ms request. Weeks of retrieval work buy half a percent; one sentence in a prompt buys a third of the latency. You cannot make that comparison without span-level attribution.
Part 3 — Drift and anomaly detectors
Two detectors doing genuinely different jobs: PSI compares a distribution against a fixed baseline, and a robust z-score flags a single window that does not belong.
1# obs/detectors.py2import numpy as np34def fit_reference(reference, n_bins=10):5 """Compute edges ONCE on known-good data. Store with the model."""6 edges = np.unique(np.percentile(reference, np.linspace(0, 100, n_bins + 1)))7 edges[0], edges[-1] = -np.inf, np.inf8 expected = np.histogram(reference, bins=edges)[0] / len(reference)9 return edges, expected1011def psi(current, edges, expected, floor=0.005):12 if len(current) < 200:13 return None # too few samples to trust14 actual = np.histogram(current, bins=edges)[0] / len(current)15 a = np.clip(actual, floor, None); e = np.clip(expected, floor, None)16 a, e = a / a.sum(), e / e.sum()17 return float(np.sum((a - e) * np.log(a / e)))1819def robust_z(x, baseline):20 med = np.median(baseline)21 mad = np.median(np.abs(baseline - med)) or 1e-922 return float((x - med) / (1.4826 * mad))Two unit tests you should assert on, because both detectors have a silent-failure mode:
1# PSI, computed by hand from five bins:2# expected 0.10 0.25 0.30 0.25 0.103# actual 0.05 0.15 0.28 0.32 0.204# (0.05-0.10)ln(0.5) = 0.034665# (0.15-0.25)ln(0.6) = 0.051086# (0.28-0.30)ln(0.933) = 0.001387# (0.32-0.25)ln(1.28) = 0.017288# (0.20-0.10)ln(2.0) = 0.069319# PSI = 0.17371 -> "moderate shift, investigate"1011# robust z: baseline median 340 tokens, MAD 1712# sigma_hat = 1.4826 x 17 = 25.2013# z(448) = (448 - 340) / 25.20 = 4.29 -> firesThe PSI bands are < 0.10 no meaningful shift, 0.10–0.25 investigate, > 0.25 act. The bin edges must come from the reference and never be recomputed on current data — if you re-quantile today's data, every bin holds 1/B of the mass by construction and PSI reads near zero no matter what happened. That single bug produces a detector that never fires, and it is the most common implementation error in this whole area.
A drift detector compares today against a stored reference. The moment the reference is derived from the data you are testing, it stops being a comparison and starts being a tautology.
The robust z-score exists because the ordinary one destroys itself. Take a 24-hour baseline of mean output length where 23 hours scatter around 340 tokens with a standard deviation of 25, and one hour — a past incident — hit 1,200:
mean = (23 x 340 + 1200) / 24 = 375.83sum of squared deviations = 13,750 + 29,533 + 679,251 = 722,534sample variance = 722,534 / 23 = 31,415 sd = 177.2classical z(448) = (448 - 375.83) / 177.2 = 0.41 -> silentrobust z(448) = (448 - 340) / 25.20 = 4.29 -> firesOne old incident sitting in the baseline inflated the standard deviation sevenfold and blinded the detector. The median has a 50% breakdown point, so one outlier in 24 moves it essentially not at all. Use the median and MAD for any baseline computed from production data, because production data contains incidents.
1# obs/run.py -- the request stream with an injected shift2import random, numpy as np3from obs.detectors import fit_reference, psi, robust_z45random.seed(7); np.random.seed(7)6BASELINE = [max(50, int(random.gauss(600, 120))) for _ in range(3000)]7EDGES, EXPECTED = fit_reference(np.array(BASELINE))89def stream(n=6000, drift_at=3000):10 for i in range(n):11 drifted = i >= drift_at12 mu, sigma = (950, 200) if drifted else (600, 120)13 yield i, max(50, int(random.gauss(mu, sigma))), driftedPart 4 — The dashboard
Collect four panels' worth of time series in fixed windows and render them. Four panels, not forty: a dashboard exists so that one person can decide in five seconds whether anything is wrong.
1# obs/dashboard.py2def sparkline(values, lo=None, hi=None, width=48):3 if not values: return ""4 blocks = "_.-~^"5 lo = min(values) if lo is None else lo6 hi = max(values) if hi is None else hi7 span = (hi - lo) or 1.08 step = max(1, len(values) // width)9 return "".join(10 blocks[min(len(blocks) - 1,11 int((v - lo) / span * (len(blocks) - 1)))]12 for v in values[::step])1314def render(panels):15 print("=" * 74)16 for name, vals, unit, lo, hi in panels:17 cur = vals[-1] if vals else 018 print(f"{name:<22} {cur:8.3f} {unit:<6} {sparkline(vals, lo, hi)}")19 print("=" * 74)2021# panels = [("p95 latency", p95_series, "s", 0, 8),22# ("error ratio", err_series, "", 0, 0.05),23# ("input tokens (mean)", tok_series, "tok", 400, 1200),24# ("answer quality (mean)", qual_series, "", 0.4, 1.0)]When the run completes, the quality panel drops from roughly 0.86 to 0.62 and the token panel climbs from about 600 to 950 at the halfway mark. The latency and error panels barely move — which is precisely the failure this whole system exists to catch.
Part 5 — Alert rules, notifiers, and a retraining gate
Before writing a single rule, compute what it will cost in interruptions. Evaluate a metric once a minute and you get 10,080 checks a week. For a stable, roughly normal metric, a one-sided threshold at k standard deviations breaches with probability 1−Φ(k) per check:
| Rule | Per-check probability | False alarms per week |
|---|---|---|
| Single breach at 2σ | 0.022750 | 229.3 |
| Single breach at 3σ | 0.0013499 | 13.6 |
| Two consecutive at 2σ | p2 = 5.176e−4 | 5.22 |
| Two consecutive at 3σ | p2 = 1.822e−6 | 0.018 |
A loose 2σ threshold requiring two consecutive breaches (5.22 a week) is 2.6 times quieter than a tight 3σ threshold firing on a single check (13.6 a week) — because squaring a small probability is a far stronger lever than pushing further into the tail. And it is not slower. Hold the false-alarm budget fixed at one page per week and solve for each design's threshold:
single breach: 10,080 x p = 1 -> p = 9.92e-5 -> k = 3.72 sigmatwo consecutive: 10,079 x p^2 = 1 -> p = 9.96e-3 -> k = 2.33 sigmaFacing a real regression that shifts the mean by 2σ, each check now breaches with probability 1−Φ(k−2). The single-breach wait is geometric at 1/p; the two-consecutive wait is (1+p)/p2:
single at 3.72 sigma: p = 0.0427 -> 1/0.0427 = 23.4 checkstwo-consec at 2.33: p = 0.3707 -> 1.3707 / 0.3707^2 = 9.97 checksSame false-alarm rate, 2.3 times faster detection. Persistence is the lever; the threshold is not. Do not push it further, though — requiring five consecutive breaches at a 3σ threshold takes 62 checks to detect a genuine 3σ shift, turning a six-minute detection into an hour for a false-alarm rate that was already negligible at two.
Persistence is the lever, not the threshold: requiring a breach to repeat squares its probability, which buys far more quiet than pushing further into the tail ever will — and costs almost nothing in detection speed.
1# obs/alerts.py2from dataclasses import dataclass, field34@dataclass5class Rule:6 name: str7 severity: str # "page" | "ticket"8 predicate: callable # (state) -> bool9 consecutive: int = 2 # the persistence requirement10 summary: str = ""11 runbook: str = ""12 _streak: int = field(default=0, repr=False)13 _firing: bool = field(default=False, repr=False)1415 def evaluate(self, state):16 self._streak = self._streak + 1 if self.predicate(state) else 017 should_fire = self._streak >= self.consecutive18 transition = should_fire and not self._firing # edge, not level19 self._firing = should_fire20 return transition2122RULES = [23 Rule("ErrorBudgetBurnFast", "page",24 lambda s: s["error_ratio_5m"] > 14.4 * 0.001, consecutive=2,25 summary="Burning error budget 14.4x: 2% of the month in one hour",26 runbook="runbooks/error-budget"),27 Rule("LatencySLOBreach", "page",28 lambda s: s["p95_latency"] > 5.0, consecutive=5,29 summary="p95 latency above the 5s SLO", runbook="runbooks/latency"),30 Rule("InputDriftHigh", "ticket",31 lambda s: (s["psi"] or 0) > 0.25, consecutive=2,32 summary="Input distribution has moved from the training baseline",33 runbook="runbooks/drift"),34 Rule("QualityAnomaly", "ticket",35 lambda s: abs(s["quality_z"]) > 3.0, consecutive=2,36 summary="Answer quality outside the robust baseline",37 runbook="runbooks/quality"),38]The transition line matters more than it looks. Returning an edge rather than a level means the rule notifies once when it starts firing, not once per evaluation for the whole incident. Level-triggered notification is how a single two-hour outage produces 120 identical Slack messages.
1def notify(rule, state):2 payload = {"alert": rule.name, "severity": rule.severity,3 "summary": rule.summary, "runbook": rule.runbook,4 "context": {k: round(v, 4) for k, v in state.items()5 if isinstance(v, (int, float))}}6 if rule.severity == "page":7 print(f"[PAGERDUTY] {payload}") # would create an incident8 else:9 print(f"[SLACK #ml-alerts] {payload}")Every notification carries the runbook link and the numbers that triggered it. An alert whose entire content is QualityAnomaly hands the responder a puzzle; this one hands them a task.
The retraining gate
Now the part that is easiest to get dangerously wrong. The tempting wiring is "drift alert fires → retraining pipeline runs → new model deploys". Do not build that. The most common cause of a large PSI is an upstream data bug, and retraining on corrupted inputs moves the corruption from a reversible pipeline problem into the model weights, where fixing the pipeline no longer fixes anything.
1def retraining_decision(state):2 """Drift alone never triggers a retrain. It opens a ticket."""3 if state["data_validation_failed"]:4 return "BLOCK: upstream data failed validation - fix the pipeline first"5 if state["accuracy_on_labels"] is None:6 return "TICKET: drift detected, no labels yet - human review"7 if state["labelled_examples_since_last_train"] < 20_000:8 return "WAIT: not enough new ground truth to learn from"9 if state["accuracy_on_labels"] < 0.88:10 return "RETRAIN: measured performance below floor"11 return "NO ACTION: inputs moved, measured performance is fine"1213def promotion_gate(champion, challenger, slices):14 # holdout n = 2,000, champion accuracy 0.9115 # SE = sqrt(0.91 x 0.09 / 2000) = sqrt(4.095e-5) = 0.0064 = 0.64 pp16 # so require ~2 SE of improvement: +1.5 pp, not "any improvement"17 return all([18 challenger.accuracy - champion.accuracy > 0.015,19 all(challenger.acc_on(s) >= champion.acc_on(s) - 0.01 for s in slices),20 challenger.score_on(REGRESSION_SET) >= 0.95,21 challenger.p95_latency_ms <= champion.p95_latency_ms * 1.10,22 ])The standard-error line is the one people skip. With a 2,000-example holdout, a challenger scoring 91.5% against a champion's 91.0% is +0.5 pp — well inside the 0.64 pp noise of a single measurement, and no evidence at all. And acc_on(slice) earns its place repeatedly: a model can gain 2 points overall while losing 9 on a minority segment, because the aggregate is dominated by the majority.
Part 6 — Wiring it together
1import numpy as np2from collections import deque34WINDOW = 2005tokens, quality = deque(maxlen=WINDOW), deque(maxlen=WINDOW)6latency = deque(maxlen=WINDOW)7quality_baseline = deque(maxlen=30)8panels = {"p95": [], "err": [], "tok": [], "qual": []}910for i, prompt_tokens, drifted in stream():11 r = traced_request(prompt_tokens, drifted)12 tokens.append(prompt_tokens); quality.append(r["quality"])13 latency.append(r["latency"])1415 if i % WINDOW == 0 and i > 0:16 mean_q = float(np.mean(quality))17 state = {18 "psi": psi(np.array(tokens), EDGES, EXPECTED),19 "quality_z": robust_z(mean_q, np.array(quality_baseline))20 if len(quality_baseline) >= 8 else 0.0,21 "p95_latency": float(np.percentile(latency, 95)),22 "error_ratio_5m": 0.0,23 "data_validation_failed": False,24 "labelled_examples_since_last_train": i,25 "accuracy_on_labels": None,26 }27 quality_baseline.append(mean_q)28 for rule in RULES:29 if rule.evaluate(state):30 notify(rule, state)31 if rule.name == "InputDriftHigh":32 print(" ->", retraining_decision(state))33 panels["tok"].append(float(np.mean(tokens)))34 panels["qual"].append(mean_q)3536render([("input tokens (mean)", panels["tok"], "tok", 400, 1200),37 ("answer quality (mean)", panels["qual"], "", 0.4, 1.0)])A correct run produces roughly this shape. Nothing fires for the first sixteen windows; p95 latency sits around 4 seconds the whole time, inside the 5-second SLO this mock service is held to. At window 17 — two windows after the injected shift, because both rules require two consecutive breaches — InputDriftHigh and QualityAnomaly notify to the ticket channel, and the retraining decision returns TICKET: drift detected, no labels yet rather than starting a retrain. No page is sent, correctly: nothing is down, and nobody can fix answer quality at 3 a.m.
The two-window delay is the persistence rule working, and it is what stops a single noisy window from generating a ticket.
Where this goes wrong
| Mistake | Symptom | Fix |
|---|---|---|
| Recomputing PSI bin edges on current data | PSI stays near 0 through an obvious shift | Fit edges once on the reference; store them with the model |
| Mean and standard deviation for the anomaly baseline | One past incident blinds the detector (z = 0.41 instead of 4.29) | Median and MAD, scaled by 1.4826 |
| Rolling reference window only | Slow degradation is chased and never flagged | Run both: rolling for incidents, fixed training baseline for validity |
| Alerting on a KS or chi-squared p-value | Everything is significant once traffic grows | Threshold on effect size (PSI, D), never on significance |
| Level-triggered notification | 120 identical messages during one incident | Notify on the transition into firing, not on every evaluation |
| Drift wired directly to retraining | A pipeline bug gets baked into the weights | Drift opens a ticket; a human decides |
| Promoting on any accuracy improvement | Models promoted on noise; performance random-walks | Require roughly 2 standard errors, plus per-slice checks |
| Instrumentation on the happy path | Dashboards go blank during the incident | Emit metrics from finally, always |
| Default histogram buckets | p95 off by seconds; SLO compliance unmeasurable | Boundaries dense where traffic and thresholds are |
| Unbounded metric labels | Metrics backend runs out of memory weeks after launch | Bounded labels only; identifiers go on log events |
What to carry into a real system
Everything above is mocked, and that is deliberate — swapping the print statements for real Slack webhooks and the in-memory exporter for an OTLP endpoint is an afternoon. The parts that took judgement are the ones that transfer.
The layered signal is the design. Latency and error rate caught nothing in this run; token distribution and quality score caught everything. A production system needs both layers, and teams reliably build only the first.
Every threshold got its false-alarm rate computed before it was written. Two minutes of arithmetic is what separates a channel people read from a channel people mute. The specific finding worth keeping: at a matched false-alarm budget, a loose threshold requiring two consecutive breaches detects a small regression more than twice as fast as a tight single-breach threshold.
Detection and response are separated by a human. The system detects drift, diagnoses which feature moved, decides on a severity, and routes to a queue. It does not retrain. That boundary is not timidity — it is the difference between a reversible incident and a corrupted model.
Baselines are stored artefacts, not code. Bin edges, expected proportions, medians, MADs, quality reference windows — all versioned alongside the model they describe. A drift number computed against a baseline that lives in someone's notebook is not reproducible, and nobody will trust it during an incident, which is the only time it matters.