我有一个相当大的训练矩阵(超过10亿行,每行两个特征).有两个类(0和1).这对于一台机器来说太大了,但幸运的是我有大约200台MPI主机供我使用.每个都是一个适度的双核工作站.
功能生成已成功分发.
Multiprocessing scikit-learn中的答案表明可以分发SGDClassifier的工作:
您可以跨核心分发数据集,执行partial_fit,获取权重向量,对它们求平均值,将它们分配给估算器,再次进行部分拟合.
当我在每个估算器上第二次运行partial_fit时,我从哪里开始获得最终的聚合估算器?
我最好的猜测是再次对coefs和截距进行平均,并使用这些值进行估算.结果估计器给出的结果与使用fit()在整个数据上构造的估计量不同.
每个主机生成局部矩阵和局部矢量.这是测试集的n行和相应的n个目标值.
每个主机使用局部矩阵和局部向量来制作SGDC分类器并进行部分拟合.然后每个都将coef向量和截距发送到root.Root对这些进行平均并将它们发送回主机.主机执行另一个partial_fit并将coef向量和截距发送到root.
Root构造具有这些值的新估计器.
local_matrix = get_local_matrix()
local_vector = get_local_vector()
estimator = linear_model.SGDClassifier()
estimator.partial_fit(local_matrix, local_vector, [0,1])
comm.send((estimator.coef_,estimator.intersept_),dest=0,tag=rank)
average_coefs = None
avg_intercept = None
comm.bcast(0,root=0)
if rank > 0:
comm.send( (estimator.coef_, estimator.intercept_ ), dest=0, tag=rank)
else:
pairs = [comm.recv(source=r, tag=r) for r in range(1,size)]
pairs.append( (estimator.coef_, estimator.intercept_) )
average_coefs = np.average([ a[0] for a in pairs ],axis=0)
avg_intercept = np.average( [ a[1][0] for a in pairs ] )
estimator.coef_ = comm.bcast(average_coefs,root=0)
estimator.intercept_ = …Run Code Online (Sandbox Code Playgroud) python parallel-processing machine-learning mpi scikit-learn