如何用Python multiprocessing starmap并行化多输入矩阵生成函数?
并行化生成U_new和V_new矩阵的实现方案
问题背景
我是数学专业学生,刚接触并行化技术。目前通过两段循环调用sample_features函数生成矩阵U_new和V_new:该函数接收迭代索引i/j,以及train_data、train_rating等多个需重复传入的数据集与参数。尝试用multiprocessing库的starmap并行化但未成功,需要正确的实现方案。
现有代码
def sample_features(train_data, train_rating, Item_vector, mu_U, Lambda_U, i, alpha, name='User'): if name=='User': idx=(train_data[:,0]==i) V_j = Item_vector[:,train_data[idx,1]] else: idx=(train_data[:,1]==i) V_j = Item_vector[:,train_data[idx,0]] Lambda_i_star=Lambda_U + alpha*np.dot(V_j, V_j.T) Lambda_i_star_inv=np.linalg.inv(Lambda_i_star) mu_i_star=np.dot(Lambda_i_star_inv,(alpha*np.dot(train_rating[idx],V_j.T)+np.dot(Lambda_U,mu_U))) return multivariate_normal(mu_i_star, Lambda_i_star_inv) for i in range(num_User): U_new[:,i]=sample_features(train_data, train_rating, Item_vector, mu_U, Lambda_U, i, alpha, name='User') for j in range(num_Item): V_new[:,j]=sample_features(train_data, train_rating, U_new, mu_V, Lambda_V, j, alpha, name='Item')
变量维度信息
- U_new (N x D)
- V_new & Item_vector (M x D)
- train_data (Rx2)
- train_rating (Rx1)
- mu_U & mu_V (D x 1)
- Lambda_U & Lambda_V (D x D)
- i & j & alpha (标量)
正确并行化实现方案
1. 核心问题修正
原代码存在一个关键错误:sample_features返回的是分布对象,而非采样后的数值数组,直接赋值给U_new/V_new会报错。首先修改函数,让它返回采样结果:
from scipy.stats import multivariate_normal # 确保导入该模块 def sample_features(train_data, train_rating, Item_vector, mu_U, Lambda_U, i, alpha, name='User'): if name=='User': idx=(train_data[:,0]==i) V_j = Item_vector[:,train_data[idx,1]] else: idx=(train_data[:,1]==i) V_j = Item_vector[:,train_data[idx,0]] Lambda_i_star=Lambda_U + alpha*np.dot(V_j, V_j.T) Lambda_i_star_inv=np.linalg.inv(Lambda_i_star) mu_i_star=np.dot(Lambda_i_star_inv,(alpha*np.dot(train_rating[idx],V_j.T)+np.dot(Lambda_U,mu_U))) # 从分布中采样1个D维样本,返回可直接赋值的数组 return multivariate_normal.rvs(mu_i_star.flatten(), Lambda_i_star_inv)
2. 并行化实现步骤
步骤1:用functools.partial绑定固定参数
减少每个任务的参数传递冗余,让并行调用更简洁:
from functools import partial from multiprocessing import Pool import numpy as np # 绑定生成User特征的固定参数 sample_user_task = partial( sample_features, train_data=train_data, train_rating=train_rating, Item_vector=Item_vector, mu_U=mu_U, Lambda_U=Lambda_U, alpha=alpha, name='User' ) # 绑定生成Item特征的固定参数(注意依赖U_new,需等U生成后再初始化) sample_item_task = partial( sample_features, train_data=train_data, train_rating=train_rating, mu_U=mu_V, # 对应原函数的mu_U参数,实际传入Item的mu_V Lambda_U=Lambda_V, # 同理,传入Item的Lambda_V alpha=alpha, name='Item' )
步骤2:并行生成U_new
# Windows系统必须把主逻辑放在if __name__ == '__main__'块中,Linux/Mac建议添加 if __name__ == '__main__': # 初始化进程池,默认使用全部CPU核心,也可指定数量如Pool(4) with Pool() as pool: # 给每个User索引分配任务,并行执行 user_samples = pool.map(sample_user_task, range(num_User)) # 将采样结果整理成NxD的U_new矩阵 U_new = np.column_stack(user_samples)
步骤3:并行生成V_new
由于V_new依赖已生成的U_new,需在U_new完成后执行:
if __name__ == '__main__': # 更新Item任务的Item_vector为刚生成的U_new sample_item_task.keywords['Item_vector'] = U_new with Pool() as pool: item_samples = pool.map(sample_item_task, range(num_Item)) V_new = np.column_stack(item_samples)
3. 常见问题排查
- 进程启动失败:Windows系统下必须将主逻辑包裹在
if __name__ == '__main__':中,避免子进程重复执行初始化代码。 - 内存占用过高:如果数据集极大,可改用
Pool.imap分批处理,或使用multiprocessing.shared_memory共享大数组,减少进程间数据拷贝。 - 参数混淆:注意原函数的
mu_U/Lambda_U是通用参数名,生成Item特征时要传入mu_V/Lambda_V,别搞混对应关系。
内容的提问来源于stack exchange,提问作者Antoniouss
相关产品推荐
相关产品推荐

