You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.04 06:25:29