from __future__ import annotations import datetime as dt import inspect import itertools import time import typing from river import base, metrics, stream, utils from river.compose.pipeline import Pipeline __all__ = ["progressive_val_score"] def _progressive_validation( dataset: base.typing.Dataset, model, metric: metrics.base.Metric, checkpoints: typing.Iterator[int], moment: str | typing.Callable[[dict], dt.datetime] | None = None, delay: str | int | dt.timedelta | typing.Callable | None = None, measure_time=False, measure_memory=False, yield_predictions=False, ): # Check that the model and the metric are in accordance if not metric.works_with(model): raise ValueError(f"{metric.__class__.__name__} metric is not compatible with {model}") # Build predict/learn/update closures once — shared by both fast and general paths. # Using closures avoids per-sample isinstance checks and branching. # Predict closure: score_one + classify for anomaly filters, predict_proba_one or predict_one if isinstance(model, base.AnomalyFilter): _score = model.score_one _classify = model.classify def predict(x, **kwargs): return _classify(_score(x, **kwargs)) elif isinstance(model, base.AnomalyDetector): predict = model.score_one # type: ignore[assignment] elif isinstance(model, base.Classifier) and not metric.requires_labels: # type: ignore predict = model.predict_proba_one # type: ignore[assignment] else: predict = model.predict_one # type: ignore[assignment] # Learn closure: supervised vs unsupervised dispatch is_supervised = model._supervised if is_supervised: learn = model.learn_one else: def learn(x, _y=None, **kwargs): model.learn_one(x, **kwargs) # Metric update closure: uniform (x, y, y_pred) signature for all metric types. metric_update: typing.Callable if isinstance(metric, metrics.base.ClusteringMetric): predict, metric_update = _build_clustering_closures(metric, model) else: _metric_update = metric.update def metric_update(x, y, y_pred): # type: ignore[no-redef] _metric_update(y_true=y, y_pred=y_pred) # If we are dealing with an active learner, we need to check whether or not a label should be # used for training or not. We'll also record how many times labels were used. from river.active.base import ActiveLearningClassifier active_learning = isinstance(model, ActiveLearningClassifier) n_samples_learned = 0 # Check once whether the model's learn_one accepts a sample weight parameter. _learn_params = inspect.signature(model.learn_one).parameters _model_accepts_w = "w" in _learn_params or any( p.kind == inspect.Parameter.VAR_KEYWORD for p in _learn_params.values() ) prev_checkpoint = None next_checkpoint = next(checkpoints, None) n_total_answers = 0 if measure_time: start = time.perf_counter() def report(y_pred): if isinstance(metric, metrics.base.Metrics): state = {m.__class__.__name__: m for m in metric} else: state = {metric.__class__.__name__: metric} state["Step"] = n_total_answers if active_learning: state["Samples used"] = n_samples_learned if measure_time: now = time.perf_counter() state["Time"] = dt.timedelta(seconds=now - start) if measure_memory: state["Memory"] = model._raw_memory_usage if yield_predictions: state["Prediction"] = y_pred return state # Fast path: no delay and no moment — the common case. # Iterates the dataset directly, skipping simulate_qa and the preds dict. if moment is None and delay is None and not active_learning: extra: list for x, y, *extra in dataset: kwargs: dict = extra[0] if extra else {} w = kwargs.pop("w", None) y_pred = predict(x, **kwargs) if y_pred is not None and y_pred != {}: metric_update(x, y, y_pred) if w is not None and _model_accepts_w: learn(x, y, w=w, **kwargs) else: learn(x, y, **kwargs) n_total_answers += 1 if n_total_answers == next_checkpoint: yield report(y_pred=y_pred) prev_checkpoint = next_checkpoint next_checkpoint = next(checkpoints, None) else: if prev_checkpoint and n_total_answers != prev_checkpoint: yield report(y_pred=None) return # General path: delayed labels or active learning — uses simulate_qa with preds dict. preds = {} for i, x, y, *extra in stream.simulate_qa(dataset, moment, delay, copy=True): kwargs = extra[0] if extra else {} w = kwargs.pop("w", None) # Case 1: no ground truth, just make a prediction if y is None: y_pred = predict(x, **kwargs) y_pred, ask_for_label = y_pred if active_learning else (y_pred, True) # type: ignore[misc] preds[i] = y_pred, ask_for_label, w continue # Case 2: there's a ground truth, model and metric can be updated y_pred, use_label, sample_w = preds.pop(i) if y_pred is not None and y_pred != {}: metric_update(x, y, y_pred) if use_label: n_samples_learned += 1 if sample_w is not None and _model_accepts_w: learn(x, y, w=sample_w, **kwargs) else: learn(x, y, **kwargs) # Yield current results n_total_answers += 1 if n_total_answers == next_checkpoint: yield report(y_pred=y_pred) prev_checkpoint = next_checkpoint next_checkpoint = next(checkpoints, None) else: # If the dataset was exhausted, we need to make sure that we yield the final results if prev_checkpoint and n_total_answers != prev_checkpoint: yield report(y_pred=None) def iter_progressive_val_score( dataset: base.typing.Dataset, model, metric: metrics.base.Metric, moment: str | typing.Callable | None = None, delay: str | int | dt.timedelta | typing.Callable | None = None, step=1, measure_time=False, measure_memory=False, yield_predictions=False, ) -> typing.Generator: """Evaluates the performance of a model on a streaming dataset and yields results. This does exactly the same as `evaluate.progressive_val_score`. The only difference is that this function returns an iterator, yielding results at every step. This can be useful if you want to have control over what you do with the results. For instance, you might want to plot the results. Parameters ---------- dataset The stream of observations against which the model will be evaluated. Each element is an `(x, y)` pair or an `(x, y, kwargs)` triple where `kwargs` is a `dict` of extra parameters passed to `learn_one`. To supply per-sample weights, include a `"w"` key in `kwargs`, e.g. `(x, y, {"w": 2.0})`. The weight is forwarded to `learn_one` for models that accept a `w` parameter (e.g. `linear_model.LogisticRegression`). model The model to evaluate. metric The metric used to evaluate the model's predictions. moment The attribute used for measuring time. If a callable is passed, then it is expected to take as input a `dict` of features. If `None`, then the observations are implicitly timestamped in the order in which they arrive. delay The amount to wait before revealing the target associated with each observation to the model. This value is expected to be able to sum with the `moment` value. For instance, if `moment` is a `datetime.date`, then `delay` is expected to be a `datetime.timedelta`. If a callable is passed, then it is expected to take as input a `dict` of features and the target. If a `str` is passed, then it will be used to access the relevant field from the features. If `None` is passed, then no delay will be used, which leads to doing standard online validation. step Iteration number at which to yield results. This only takes into account the predictions, and not the training steps. measure_time Whether or not to measure the elapsed time. measure_memory Whether or not to measure the memory usage of the model. yield_predictions Whether or not to include predictions. If step is 1, then this is equivalent to yielding the predictions at every iterations. Otherwise, not all predictions will be yielded. Examples -------- Take the following model: >>> from river import linear_model >>> from river import preprocessing >>> model = ( ... preprocessing.StandardScaler() | ... linear_model.LogisticRegression() ... ) We can evaluate it on the `Phishing` dataset as so: >>> from river import datasets >>> from river import evaluate >>> from river import metrics >>> steps = evaluate.iter_progressive_val_score( ... model=model, ... dataset=datasets.Phishing(), ... metric=metrics.ROCAUC(), ... step=200 ... ) >>> for step in steps: ... print(step) {'ROCAUC': ROCAUC: 90.20%, 'Step': 200} {'ROCAUC': ROCAUC: 92.25%, 'Step': 400} {'ROCAUC': ROCAUC: 93.23%, 'Step': 600} {'ROCAUC': ROCAUC: 94.05%, 'Step': 800} {'ROCAUC': ROCAUC: 94.79%, 'Step': 1000} {'ROCAUC': ROCAUC: 95.07%, 'Step': 1200} {'ROCAUC': ROCAUC: 95.07%, 'Step': 1250} The `yield_predictions` parameter can be used to include the predictions in the results: >>> import itertools >>> steps = evaluate.iter_progressive_val_score( ... model=model, ... dataset=datasets.Phishing(), ... metric=metrics.ROCAUC(), ... step=1, ... yield_predictions=True ... ) >>> for step in itertools.islice(steps, 100, 105): ... print(step) {'ROCAUC': ROCAUC: 94.68%, 'Step': 101, 'Prediction': {False: 0.966..., True: 0.033...}} {'ROCAUC': ROCAUC: 94.75%, 'Step': 102, 'Prediction': {False: 0.035..., True: 0.964...}} {'ROCAUC': ROCAUC: 94.82%, 'Step': 103, 'Prediction': {False: 0.043..., True: 0.956...}} {'ROCAUC': ROCAUC: 94.89%, 'Step': 104, 'Prediction': {False: 0.816..., True: 0.183...}} {'ROCAUC': ROCAUC: 94.96%, 'Step': 105, 'Prediction': {False: 0.041..., True: 0.958...}} References ---------- [^1]: [Beating the Hold-Out: Bounds for K-fold and Progressive Cross-Validation](http://hunch.net/~jl/projects/prediction_bounds/progressive_validation/coltfinal.pdf) [^2]: [Grzenda, M., Gomes, H.M. and Bifet, A., 2019. Delayed labelling evaluation for data streams. Data Mining and Knowledge Discovery, pp.1-30](https://link.springer.com/content/pdf/10.1007%2Fs10618-019-00654-y.pdf) """ yield from _progressive_validation( dataset, model, metric, checkpoints=itertools.count(step, step) if step else iter([]), moment=moment, delay=delay, measure_time=measure_time, measure_memory=measure_memory, yield_predictions=yield_predictions, ) def progressive_val_score( dataset: base.typing.Dataset, model, metric: metrics.base.Metric, moment: str | typing.Callable | None = None, delay: str | int | dt.timedelta | typing.Callable | None = None, print_every=0, show_time=False, show_memory=False, **print_kwargs, ) -> metrics.base.Metric: """Evaluates the performance of a model on a streaming dataset. This method is the canonical way to evaluate a model's performance. When used correctly, it allows you to exactly assess how a model would have performed in a production scenario. `dataset` is converted into a stream of questions and answers. At each step the model is either asked to predict an observation, or is either updated. The target is only revealed to the model after a certain amount of time, which is determined by the `delay` parameter. Note that under the hood this uses the `stream.simulate_qa` function to go through the data in arrival order. By default, there is no delay, which means that the samples are processed one after the other. When there is no delay, this function essentially performs progressive validation. When there is a delay, then we refer to it as delayed progressive validation. It is recommended to use this method when you want to determine a model's performance on a dataset. In particular, it is advised to use the `delay` parameter in order to get a reliable assessment. Indeed, in a production scenario, it is often the case that ground truths are made available after a certain amount of time. By using this method, you can reproduce this scenario and therefore truthfully assess what would have been the performance of a model on a given dataset. Parameters ---------- dataset The stream of observations against which the model will be evaluated. Each element is an `(x, y)` pair or an `(x, y, kwargs)` triple where `kwargs` is a `dict` of extra parameters passed to `learn_one`. To supply per-sample weights, include a `"w"` key in `kwargs`, e.g. `(x, y, {"w": 2.0})`. The weight is forwarded to `learn_one` for models that accept a `w` parameter (e.g. `linear_model.LogisticRegression`). model The model to evaluate. metric The metric used to evaluate the model's predictions. moment The attribute used for measuring time. If a callable is passed, then it is expected to take as input a `dict` of features. If `None`, then the observations are implicitly timestamped in the order in which they arrive. delay The amount to wait before revealing the target associated with each observation to the model. This value is expected to be able to sum with the `moment` value. For instance, if `moment` is a `datetime.date`, then `delay` is expected to be a `datetime.timedelta`. If a callable is passed, then it is expected to take as input a `dict` of features and the target. If a `str` is passed, then it will be used to access the relevant field from the features. If `None` is passed, then no delay will be used, which leads to doing standard online validation. print_every Iteration number at which to print the current metric. This only takes into account the predictions, and not the training steps. show_time Whether or not to display the elapsed time. show_memory Whether or not to display the memory usage of the model. print_kwargs Extra keyword arguments are passed to the `print` function. For instance, this allows providing a `file` argument, which indicates where to output progress. Examples -------- Take the following model: >>> from river import linear_model >>> from river import preprocessing >>> model = ( ... preprocessing.StandardScaler() | ... linear_model.LogisticRegression() ... ) We can evaluate it on the `Phishing` dataset as so: >>> from river import datasets >>> from river import evaluate >>> from river import metrics >>> evaluate.progressive_val_score( ... model=model, ... dataset=datasets.Phishing(), ... metric=metrics.ROCAUC(), ... print_every=200 ... ) [200] ROCAUC: 90.20% [400] ROCAUC: 92.25% [600] ROCAUC: 93.23% [800] ROCAUC: 94.05% [1,000] ROCAUC: 94.79% [1,200] ROCAUC: 95.07% [1,250] ROCAUC: 95.07% ROCAUC: 95.07% We haven't specified a delay, therefore this is strictly equivalent to the following piece of code: >>> model = ( ... preprocessing.StandardScaler() | ... linear_model.LogisticRegression() ... ) >>> metric = metrics.ROCAUC() >>> for x, y in datasets.Phishing(): ... y_pred = model.predict_proba_one(x) ... metric.update(y, y_pred) ... model.learn_one(x, y) >>> metric ROCAUC: 95.07% When `print_every` is specified, the current state is printed at regular intervals. Under the hood, Python's `print` method is being used. You can pass extra keyword arguments to modify its behavior. For instance, you may use the `file` argument if you want to log the progress to a file of your choice. >>> with open('progress.log', 'w') as f: ... metric = evaluate.progressive_val_score( ... model=model, ... dataset=datasets.Phishing(), ... metric=metrics.ROCAUC(), ... print_every=200, ... file=f ... ) >>> with open('progress.log') as f: ... for line in f.read().splitlines(): ... print(line) [200] ROCAUC: 94.00% [400] ROCAUC: 94.70% [600] ROCAUC: 95.17% [800] ROCAUC: 95.42% [1,000] ROCAUC: 95.82% [1,200] ROCAUC: 96.00% [1,250] ROCAUC: 96.04% Note that the performance is slightly better than above because we haven't used a fresh copy of the model. Instead, we've reused the existing model which has already done a full pass on the data. >>> import os; os.remove('progress.log') References ---------- [^1]: [Beating the Hold-Out: Bounds for K-fold and Progressive Cross-Validation](http://hunch.net/~jl/projects/prediction_bounds/progressive_validation/coltfinal.pdf) [^2]: [Grzenda, M., Gomes, H.M. and Bifet, A., 2019. Delayed labelling evaluation for data streams. Data Mining and Knowledge Discovery, pp.1-30](https://link.springer.com/content/pdf/10.1007%2Fs10618-019-00654-y.pdf) """ checkpoints = iter_progressive_val_score( dataset=dataset, model=model, metric=metric, moment=moment, delay=delay, step=print_every, measure_time=show_time, measure_memory=show_memory, ) from river.active.base import ActiveLearningClassifier active_learning = isinstance(model, ActiveLearningClassifier) for checkpoint in checkpoints: msg = f"[{checkpoint['Step']:,d}] {metric}" if active_learning: msg += f" – {checkpoint['Samples used']:,d} samples used" if show_time: H, rem = divmod(checkpoint["Time"].seconds, 3600) M, S = divmod(rem, 60) msg += f" – {H:02d}:{M:02d}:{S:02d}" if show_memory: msg += f" – {utils.pretty.humanize_bytes(checkpoint['Memory'])}" print(msg, **print_kwargs) return metric def _build_clustering_closures(metric, model): """Build predict and metric-update closures for clustering evaluation. Clustering metrics expect ``(x, y_pred, centers)`` rather than ``(y_true, y_pred)``. When the model is a pipeline the raw features must be transformed through every step except the final clusterer so that ``x`` and ``centers`` live in the same feature space. Returns ``(predict, metric_update)`` where both closures share a single transform pass for pipelines, avoiding redundant work. """ if isinstance(model, Pipeline): clusterer = model._last_step _predict = clusterer.predict_one # Stash the last transformed x so metric_update can reuse it. _state = {} def predict(x, **kwargs): x_transformed, _ = model._transform_one(x) _state["x"] = x_transformed return _predict(x_transformed, **kwargs) def metric_update(x, y, y_pred): centers = getattr(clusterer, "centers", None) if centers is not None: metric.update(_state["x"], y_pred, centers) else: predict = model.predict_one # type: ignore[assignment] def metric_update(x, y, y_pred): centers = getattr(model, "centers", None) if centers is not None: metric.update(x, y_pred, centers) return predict, metric_update