这个问题比较高级,但我会提供一些可能有用的建议。
首先,您编写的代码在本地运行大部分内容。要并行执行 ML 训练,您需要:
- 在集群上工作(本地或远程)。
- 将数据存储在 Dask 数组或数据帧中
- 使用
dask.delayed 任务
或
- 使用
client.submit() API
1.创建(本地)集群
从您的代码中不清楚您是否已经实例化了一个客户端,所以也许只需仔细检查您是否在此处关注the dask-ml docs instructions:
from dask.distributed import Client
import joblib
client = Client(processes=False) # create local cluster
# import coiled # or connect to remote cluster
# client = Client(coiled.Cluster())
with joblib.parallel_backend('dask'):
# your scikit-learn code
但是,请注意,scitkit-learn 的 Dask joblib 后端对于扩展受 CPU 限制的工作负载很有用。要扩展到受 RAM 限制的工作负载(大于内存的数据集),您需要考虑使用 dask-ml 并行估计器之一,如下所示。
2。将数据存储在 Dask 数组中
下面的最小代码示例将两个虚拟数据集设置为 Dask 数组并实例化 K-Means 聚类算法。
import dask_ml.datasets
import dask_ml.cluster
import matplotlib.pyplot as plt
# create dummy datasets
X, y = dask_ml.datasets.make_blobs(n_samples=10000000,
chunks=1000000,
random_state=0,
centers=3)
X2, y2 = dask_ml.datasets.make_blobs(n_samples=10000000,
chunks=1000000,
random_state=3,
centers=3)
# persist predictor sets to cluster memory
X = X.persist()
X2 = X2.persist()
# instantiate KM model
km = dask_ml.cluster.KMeans(n_clusters=3, init_max_iter=2, oversampling_factor=10)
3.与 Dask.Delayed 并行训练
下面的代码使用dask.delayed API 并行运行训练。它遵循the best practices outlined in the Dask docs。
from dask import delayed
import dask
X = delayed(X)
X2 = delayed(X2)
@delayed
def train(model, X):
return model.fit(X)
# define task graphs (lazy evaluation, no computation triggered)
km1 = train(km, X)
km2 = train(km, X2)
# trigger computation and yield fitted models in parallel
km1, km2 = dask.compute(km1, km2)
4.与 Futures 和 client.submit
并行训练
或者,您可以使用client.submit() API 并行训练。这会立即返回指向正在进行的计算并最终指向存储结果的未来。阅读更多内容the docs here。
根据您的问题表述,我假设您的主要优先事项是让培训并行进行。这不需要手动将任务分配给特定的工作人员; Dask 会为您处理工作人员之间的调度和最佳分配。如果您真的对手动将特定任务分配给特定工作人员感兴趣,我建议您查看this SO answer。