Score and Predict Large Datasets
================================

Sometimes you'll train on a smaller dataset that fits in memory, but need to predict or score for a much larger (possibly larger than memory) dataset. Perhaps your [learning curve](http://scikit-learn.org/stable/modules/learning_curve.html) has leveled off, or you only have labels for a subset of the data.

In this situation, you can use [ParallelPostFit](http://ml.dask.org/modules/generated/dask_ml.wrappers.ParallelPostFit.html) to parallelize and distribute the scoring or prediction steps.

In [2]:
from dask.distributed import Client, progress


# Assuming you've already set up your cluster and client
client = Client('localhost:8786')

client

0,1
Connection method: Direct,
Dashboard: http://localhost:8787/status,

0,1
Comm: tcp://10.249.7.222:8786,Workers: 5
Dashboard: http://10.249.7.222:8787/status,Total threads: 1280
Started: Just now,Total memory: 2.33 TiB

0,1
Comm: tcp://10.249.7.222:33777,Total threads: 256
Dashboard: http://10.249.7.222:39129/status,Memory: 476.56 GiB
Nanny: tcp://10.249.7.222:41653,
Local directory: /tmp/dask-scratch-space/worker-fa9ohm5_,Local directory: /tmp/dask-scratch-space/worker-fa9ohm5_
Tasks executing:,Tasks in memory:
Tasks ready:,Tasks in flight:
CPU usage: 2.0%,Last seen: Just now
Memory usage: 107.88 MiB,Spilled bytes: 0 B
Read bytes: 17.03 kiB,Write bytes: 107.41 kiB

0,1
Comm: tcp://10.249.7.222:37135,Total threads: 256
Dashboard: http://10.249.7.222:43763/status,Memory: 476.56 GiB
Nanny: tcp://10.249.7.222:40401,
Local directory: /tmp/dask-scratch-space/worker-g0tepp2s,Local directory: /tmp/dask-scratch-space/worker-g0tepp2s
Tasks executing:,Tasks in memory:
Tasks ready:,Tasks in flight:
CPU usage: 2.0%,Last seen: Just now
Memory usage: 110.01 MiB,Spilled bytes: 0 B
Read bytes: 15.53 kiB,Write bytes: 106.14 kiB

0,1
Comm: tcp://10.249.7.222:39289,Total threads: 256
Dashboard: http://10.249.7.222:40111/status,Memory: 476.56 GiB
Nanny: tcp://10.249.7.222:39509,
Local directory: /tmp/dask-scratch-space/worker-3qlumrt2,Local directory: /tmp/dask-scratch-space/worker-3qlumrt2
Tasks executing:,Tasks in memory:
Tasks ready:,Tasks in flight:
CPU usage: 2.0%,Last seen: Just now
Memory usage: 108.25 MiB,Spilled bytes: 0 B
Read bytes: 16.68 kiB,Write bytes: 107.10 kiB

0,1
Comm: tcp://10.249.7.222:44857,Total threads: 256
Dashboard: http://10.249.7.222:32961/status,Memory: 476.56 GiB
Nanny: tcp://10.249.7.222:43529,
Local directory: /tmp/dask-scratch-space/worker-xop4ozl5,Local directory: /tmp/dask-scratch-space/worker-xop4ozl5
Tasks executing:,Tasks in memory:
Tasks ready:,Tasks in flight:
CPU usage: 2.0%,Last seen: Just now
Memory usage: 108.15 MiB,Spilled bytes: 0 B
Read bytes: 15.94 kiB,Write bytes: 105.36 kiB

0,1
Comm: tcp://10.249.7.222:45439,Total threads: 256
Dashboard: http://10.249.7.222:46117/status,Memory: 476.56 GiB
Nanny: tcp://10.249.7.222:44437,
Local directory: /tmp/dask-scratch-space/worker-pm43fvu_,Local directory: /tmp/dask-scratch-space/worker-pm43fvu_
Tasks executing:,Tasks in memory:
Tasks ready:,Tasks in flight:
CPU usage: 2.0%,Last seen: Just now
Memory usage: 107.89 MiB,Spilled bytes: 0 B
Read bytes: 16.65 kiB,Write bytes: 106.94 kiB


2024-11-19 17:24:05,752 - distributed.client - ERROR - Failed to reconnect to scheduler after 30.00 seconds, closing client


In [3]:
import numpy as np
import dask.array as da
from sklearn.datasets import make_classification

ModuleNotFoundError: No module named 'sklearn'

In [14]:
from dask.distributed import PipInstall
plugin = PipInstall(packages=["dask-ml"], pip_options=["--upgrade"])

We'll generate a small random dataset with scikit-learn.

In [None]:
X_train, y_train = make_classification(
    n_features=2, n_redundant=0, n_informative=2,
    random_state=1, n_clusters_per_class=1, n_samples=1000)
X_train[:5]

And we'll clone that dataset many times with `dask.array`. `X_large` and `y_large` represent our larger than memory dataset.

In [None]:
# Scale up: increase N, the number of times we replicate the data.
N = 100
X_large = da.concatenate([da.from_array(X_train, chunks=X_train.shape)
                          for _ in range(N)])
y_large = da.concatenate([da.from_array(y_train, chunks=y_train.shape)
                          for _ in range(N)])
X_large

Since our training dataset fits in memory, we can use a scikit-learn estimator as the actual estimator fit during traning.
But we know that we'll want to predict for a large dataset, so we'll wrap the scikit-learn estimator with `ParallelPostFit`.

In [17]:
from sklearn.linear_model import LogisticRegressionCV
from dask_ml.wrappers import ParallelPostFit

In [18]:
clf = ParallelPostFit(LogisticRegressionCV(cv=3), scoring="r2")

See the note in the `dask-ml`'s documentation about when and why a `scoring` parameter is needed: https://ml.dask.org/modules/generated/dask_ml.wrappers.ParallelPostFit.html#dask_ml.wrappers.ParallelPostFit.

Now we'll call `clf.fit`. Dask-ML does nothing here, so this step can only use datasets that fit in memory.

In [None]:
clf.fit(X_train, y_train)

Now that training is done, we'll turn to predicting for the full (larger than memory) dataset.

In [None]:
y_pred = clf.predict(X_large)
y_pred

y_pred is Dask arary. Workers can write the predicted values to a shared file system, without ever having to collect the data on a single machine.

Or we can check the models score on the entire large dataset. The computation will be done in parallel, and no single machine will have to hold all the data.

In [21]:
#clf.score(X_large, y_large)