【发布时间】:2014-01-29 21:11:03
【问题描述】:
我有一个相当大的训练矩阵(超过 10 亿行,每行两个特征)。有两个类(0 和 1)。 这对于单台机器来说太大了,但幸运的是我有大约 200 个 MPI 主机可供我使用。每个都是一个普通的双核工作站。
特征生成已经成功分发。
Multiprocessing scikit-learn 中的答案表明可以分发 SGDClassifier 的工作:
您可以跨核分布数据集,进行部分拟合,获取权重向量,对其进行平均,将它们分配给估计器,再次进行部分拟合。
当我对每个估算器第二次运行 partial_fit 时,我应该从哪里获得最终的聚合估算器?
我最好的猜测是再次平均系数和截距,并使用这些值进行估计。结果估计器给出的结果与使用 fit() 对整个数据构造的估计器不同。
详情
每个主机生成一个局部矩阵和一个局部向量。这是测试集的n行和对应的n个目标值。
每个主机都使用局部矩阵和局部向量来制作 SGDClassifier 并进行部分拟合。然后每个都将 coef 向量和截距发送到根。 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_ = np.array( [comm.bcast(avg_intercept,root=0)] )
estimator.partial_fit(metric_matrix, edges_exist,[0,1])
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_ = average_coefs
estimator.intercept_ = np.array( [avg_intercept] )
print("The estimator at rank 0 should now be working")
谢谢!
【问题讨论】:
标签: python parallel-processing machine-learning mpi scikit-learn