A model trained last month is answering questions about this morning. For ads, feeds, search ranking, fraud and pricing, the world moves fast enough that a model's accuracy decays within days or hours: new items appear, user tastes shift, adversaries adapt. Online learning is the architecture that keeps updating the model from live traffic, so the gap between what happened and what the model knows stays at minutes instead of weeks.

The learning algorithm is the small part. The hard parts are getting correct labels for predictions whose outcome arrives late, making sure the features you train on are the ones you served with, and making sure a model that is changing every minute cannot quietly get worse in production. This article builds that system end to end, with runnable code for a streaming learner, a label joiner and a promotion gate.

Advertisement

Three things people mean by online

The word covers three designs, and the first decision is which one you need. Online inference means predictions are served live, but training is a batch job; nothing here is online learning. Frequent retraining means a batch job runs every hour or day, warm-started from the previous model on the newest data; it is simple, uses your existing training stack and is enough for most teams. True online learning updates model parameters continuously from a stream of labelled examples, with a new version in serving every few minutes.

Move to the third only when you can show that freshness pays. Replay a week of logs, train models with different data cut-offs, and measure how quickly quality decays with staleness. If a model six hours old is as good as one five minutes old, hourly retraining gives you the benefit with far less machinery. If quality falls within the hour, as it often does for new-item recommendation and ad click prediction, the streaming architecture is worth building. The drift detection article covers how to measure that decay in production.

The architecture: every stream explained

Serving scores each request and writes a prediction record: a prediction id, the timestamp, the model version and the exact feature values used. Outcomes such as clicks, purchases or chargebacks arrive later on a separate stream, carrying the prediction id. A joiner matches them within a time window and emits labelled training examples. A streaming trainer consumes these, updates the model and writes checkpoints. A validator compares each checkpoint with the live model on recent held-out data and promotes it to a registry if it passes. Serving picks up the new version without a restart.

An online learning loop: log at serve, join labels late, train, gate, swapServingpredict + log featuresPrediction logid, features, scoreLabel eventsclick, buy, timeoutJoinerwindow + watermarkStreaming trainerpredict, then learnValidatorholdout + guardsModel registryversions, rollbackstreamexamplescheckpointpromotehot swapEvery arrow is a stream with lag; the joiner's window sets how stale labels may be.Nothing reaches serving without passing the validator.
Figure 1. Online learning is a loop of streams: prediction logs and label events meet in a joiner, the trainer learns from joined examples, and a validator decides what serving may use.

Two principles hold the design together. First, log at serve: train on the features exactly as serving saw them, not on features recomputed later, which may differ because the feature pipeline has moved on. Recomputing is the most common source of training-serving skew; the feature store article covers the point-in-time lookups needed if you must recompute. Second, every hop is a stream with lag and possible duplicates, so each stage must tolerate reordering and replay.

Advertisement

Delayed labels and the join window

Positive labels arrive when something happens. Negative labels are the absence of an event, so they can only be emitted when you stop waiting. The window length is therefore a trade-off between freshness and label correctness. A short window gives fresh training data but mislabels late positives as negatives; a long window is accurate but delays every example.

Worked example: a click-through model where 90 percent of clicks arrive within 2 minutes, 98 percent within 30 minutes and the rest over hours. With a 30-minute window, 2 percent of true positives are trained as negatives, so if the true click rate is 3 percent the trained base rate is about 2.94 percent, a small, stable bias you can correct in calibration. With a 2-minute window, 10 percent of positives become negatives and the example latency drops from 30 minutes to 2. Measure your own delay distribution before choosing; it varies widely between clicks, purchases and fraud chargebacks, which can take weeks.

# Joiner: emit a positive as soon as a label arrives, a negative only when the window closes.
WINDOW = 30 * 60                                   # seconds; from measured label delays

def on_prediction(ev, state):
    state.put(ev.prediction_id, ev, expire_at=ev.ts + WINDOW)

def on_label(ev, state, out):
    pred = state.pop(ev.prediction_id)
    if pred is None:
        metrics.inc("label_without_prediction")    # late, duplicate or never logged
        return
    out.emit(example(pred.features, label=1, served_by=pred.model_version))

def on_watermark(wm, state, out):
    # Watermark = min(prediction stream time, label stream time). A lagging label stream
    # holds it back, so windows cannot close early and flood training with false negatives.
    for pred in state.expire_before(wm):
        out.emit(example(pred.features, label=0, served_by=pred.model_version))

The watermark is what keeps this correct under lag. If the label stream falls behind, say its consumer restarts, a joiner that closes windows by wall-clock time emits a flood of false negatives and the model learns that nobody clicks anything. Closing windows only when both streams have progressed past the window end prevents that. The stateful stream join article covers join state and timers in depth.

The learner: algorithms that update one example at a time

Streaming learners must update cheaply per example, cope with features that appear and disappear, and never need a pass over all history. For sparse, high-dimensional problems such as ad click prediction, the standard choice is logistic regression with FTRL-Proximal, the algorithm McMahan and colleagues described in 2013 for Google's ad click prediction. It keeps two numbers per feature, gives each feature its own adaptive learning rate, and uses L1 regularisation to keep rarely seen features at exactly zero, which keeps the model small as hashed feature spaces grow.

import math
from collections import defaultdict

class FTRLProximal:
    """Per-coordinate FTRL-Proximal logistic regression (McMahan et al., 2013)."""

    def __init__(self, alpha=0.1, beta=1.0, l1=1.0, l2=1.0):
        self.alpha, self.beta, self.l1, self.l2 = alpha, beta, l1, l2
        self.z = defaultdict(float)   # accumulated adjusted gradients
        self.n = defaultdict(float)   # accumulated squared gradients

    def _weight(self, i):
        z = self.z[i]
        if abs(z) <= self.l1:
            return 0.0                # L1 keeps rare features exactly zero
        sign = 1.0 if z > 0 else -1.0
        return -(z - sign * self.l1) / ((self.beta + math.sqrt(self.n[i])) / self.alpha + self.l2)

    def predict(self, x):
        """x: dict of hashed feature index -> value."""
        s = sum(self._weight(i) * v for i, v in x.items())
        return 1.0 / (1.0 + math.exp(-max(min(s, 35.0), -35.0)))

    def learn(self, x, y):
        p = self.predict(x)           # progressive validation: score before learning
        for i, v in x.items():
            g = (p - y) * v
            w = self._weight(i)
            sigma = (math.sqrt(self.n[i] + g * g) - math.sqrt(self.n[i])) / self.alpha
            self.z[i] += g - sigma * w
            self.n[i] += g * g
        return p

Note the order inside learn: the example is scored before the model learns from it. The running average of those scores' losses is called progressive validation. Every example is out-of-sample at the moment it is scored, so it is an honest estimate of live quality that costs nothing extra. On a synthetic stream of 20,000 examples, the class above gives progressive log loss of about 0.515 in the first 5,000 events, dropping to about 0.436 as it converges.

For deep models the pattern is to keep the network architecture fixed and continue training from the latest checkpoint on small batches of the newest examples, often updating embeddings and the final layers more often than the whole network. For tabular problems in Python, the River library provides incremental versions of many models with the same score-then-learn loop.

from river import compose, linear_model, metrics, preprocessing

model = compose.Pipeline(preprocessing.StandardScaler(), linear_model.LogisticRegression())
metric = metrics.Accuracy()

for x, y in joined_examples():          # dicts of features, bool labels, in event-time order
    y_pred = model.predict_one(x)       # score first: an honest out-of-sample estimate
    metric.update(y, y_pred)
    model.learn_one(x, y)

Checkpoint, validate, promote

A model that updates continuously must still be released deliberately. The trainer writes a checkpoint every few minutes. A validator scores the candidate and the live model on the same recent held-out slice, a sample of joined examples the trainer never sees, and promotes only if the candidate is no worse and passes sanity guards. Serving loads the new version in the background and swaps an atomic pointer, so requests never see a half-loaded model. The registry keeps the last several versions so rollback is a pointer change.

def maybe_promote(candidate, live, holdout, guards):
    """Run every few minutes on the newest checkpoint. Returns the model serving should use."""
    c, l = evaluate(candidate, holdout), evaluate(live, holdout)   # same recent, held-out slice
    checks = {
        "logloss": c.logloss <= l.logloss * 1.01,         # no worse than 1 percent
        "calibration": abs(c.mean_pred - c.base_rate) < guards.max_calib_gap,
        "weights": c.weight_norm < guards.max_weight_norm,  # catches divergence
        "volume": candidate.examples_since(live) > guards.min_examples,
    }
    if all(checks.values()):
        registry.promote(candidate)                      # serving swaps an atomic pointer
        return candidate
    alert_if_repeated(checks)                            # a stuck trainer is also a failure
    return live

Every prediction record carries the model version that produced it, so if quality drops you can attribute it to a specific promotion. A trainer that stops promoting is also a failure: the model is silently going stale, so alert on model age as well as on quality.

Stability: feedback loops, forgetting and poisoning

An online model trains on the consequences of its own decisions. Items it ranks low are rarely shown, so they rarely collect positive labels, so they stay low. This feedback loop narrows what the model knows. The standard counter is a small amount of exploration, randomising a fraction of traffic or using bandit-style selection, and logging the propensity so training can reweight. The multi-armed bandit article covers exploration strategies.

Recency is also a liability. A model that adapts quickly to a weekend or a holiday forgets the weekday pattern and must relearn it on Monday. Adaptive per-feature learning rates, as in FTRL, damp this for well-established features. Keep a batch-trained baseline and compare against it, so slow drift away from a good solution becomes visible. Finally, anything that can generate labels can steer the model: bots, click fraud and coordinated abuse. Filter known invalid traffic before the joiner, cap how much any single user or source can contribute per window, and alert on sudden changes in label volume.

Failure modes

  • Label stream lag. Without watermarks, windows close early and the model learns false negatives. Watch the join rate and the base rate of emitted examples.
  • Duplicate events. Retries and replays double-count positives. Deduplicate by event id in the joiner and make trainer offsets and checkpoints consistent; the exactly-once processing article covers the options.
  • Training-serving skew. Features recomputed at training time differ from served ones. Log at serve and compare feature distributions between the prediction log and serving.
  • Divergence. A bad batch or a learning-rate bug blows up weights. The validator's weight-norm and calibration guards stop it from reaching serving.
  • Clock skew. Event timestamps from different hosts disagree, and joins silently miss. Use server-assigned timestamps for windowing.
  • Silent staleness. The trainer crashes and serving keeps the last model forever. Alert on model age.

Operating it: what to watch

Watch the pipeline and the model separately. For the pipeline: lag of each consumer, join rate (the fraction of predictions that receive a label or time out cleanly), labels without predictions, the delay distribution of labels, and examples per minute. For the model: progressive log loss against a batch-trained baseline, calibration (mean prediction versus observed rate) per segment, weight norm, number of non-zero features, time since last promotion and promotions rejected per hour. Replay capability matters more than any single metric: keep the prediction and label logs long enough to rebuild a model from a known-good point, because some failures are only noticed days later.

Trade-offs at a glance

ChoiceGainCost
Hourly retraining instead of streamingReuses batch stack, easy to debugUp to an hour stale
Short join windowFresher examplesLate positives become negatives
Linear model with FTRLCheap, sparse, stableLess expressive than deep models
Continued training of deep modelRicher model stays freshHarder to validate and to roll back
Exploration trafficBreaks feedback loopsSome requests get worse answers

What to do next

  1. Measure how fast your model decays: train with data cut-offs of 5 minutes, 1 hour and 1 day on replayed logs and compare.
  2. If hourly retraining is enough, stop there; otherwise continue with this list.
  3. Log features, model version and a prediction id at serve time for every scored request.
  4. Measure your label delay distribution and choose a window, then build the joiner with watermarks on both streams.
  5. Start with a progressively validated streaming learner, such as the FTRL class above, beside your batch model.
  6. Add the promotion gate with a held-out slice, calibration and weight-norm guards, and alert on model age.
  7. Add a small exploration slice and cap per-source label contributions before the first launch.
Key takeaway: Online learning keeps a model current by training on a stream of labelled live traffic, but most of the work is in the data path. Log features at serve time, join delayed labels with windows sized from measured delays and closed by watermarks, and score each example before learning from it so progressive validation measures live quality. Release every checkpoint through a gate with rollback, and guard against feedback loops, poisoning and silent staleness. If hourly retraining is fresh enough, use that instead.