如何在Spark中实现大规模矩阵乘法计算物品相似度并导出scipy稀疏矩阵
现有存储(user, item, rating)三元组的Spark DataFrame,样例数据如下:
user item rating ust001 ipx001 5 ust002 ipx04 2 ust001 itx001 4 ust002 iox04 5
设共有n个用户、m个物品,可构建规模为n×m的用户-物品评分矩阵A。计算目标为基于该矩阵得到物品-物品相似度矩阵B = A^T * A,最终保存为scipy稀疏矩阵格式的B.npz文件。
原有单机Python实现逻辑在小数据集下可正常运行:先完成用户、物品字符串ID到数值索引的映射,存储映射字典后将评分数据转换为数值索引格式;再基于scipy构建csr_matrix格式的稀疏矩阵,通过矩阵乘法计算得到物品相似度矩阵,最终保存为npz文件。
ID映射阶段代码如下:
import numpy as np import pandas as pd import pickle df = pd.read('user_item.paruet') # mapping string to index user2num = {} item2num = {} UID = 0 IID = 0 # remaping index to string num2user = {} num2ite ={} # loop over all emelemt and map string to index for i in range(len(df['user'])): if df['user'][i] not in user2num: user2num[df['user'][i]] = UID num2ser[UID] = df['user'][i] UID += 1 if df['item'][i] not in item2num: item2num[df['item'][i]] = IID num2item[IID] = df['item'][i] IID += 1 # save the pair of string-index with open('num2item.pickle', 'wb') as handle: pickle.dump(num2item, handle, protocol=pickle.HIGHEST_PROTOCOL) with open('item2num.pickle', 'wb') as handle: pickle.dump(item2num, handle, protocol=pickle.HIGHEST_PROTOCOL) with open('num2user.pickle', 'wb') as handle: pickle.dump(num2user, handle, protocol=pickle.HIGHEST_PROTOCOL) with open('user2num.pickle', 'wb') as handle: pickle.dump(user2num, handle, protocol=pickle.HIGHEST_PROTOCOL) df["user"] = df["user"].map(pan2num) df["item"] = df["item"].map(mrch2num) df.to_parquet('ID_user-item.parquet')
矩阵计算阶段代码如下:
# another file to compte item-item similarity import numpy as np import pandas as pd import pickle from scipy.sparse import csr_matrix from scipy import sparse import scipy df = pd.read_parquet('ID_user-item.parquet') with open('num2item.pickle', 'rb') as handle: item_id = pickle.load(handle) with open('num2user.pickle', 'rb') as handle: user_id = pickle.load(handle) row = df['user'].values col = df['item'].values data = df['rating'].values A = csr_matrix((data,(row, col)), shape=(len(user_id), len(item_id))) B = csr_matrix((data,(col, row)), shape=(len(item_id), len(user_id))) C = sparse.csr_matrix.dot(B, A) scipy.sparse.save_npz('item-item.npz', C) #based on num2item, I can remap the index to string to retrival the item-item similarity.
该方案存在明确瓶颈:当用户-物品评分数据量达到500G时,单机Python程序会持续出现内存溢出问题,需要沿用原计算逻辑,基于Spark实现大规模矩阵乘法运算,最终得到item-item.npz格式的scipy稀疏矩阵文件。
整体逻辑和单机版完全对齐,所有重计算步骤分布式执行,仅在最后阶段将非零结果收集到驱动节点生成scipy稀疏矩阵,避免全量数据加载到单机内存。
步骤1:分布式生成ID映射
放弃单机循环遍历的映射方式,用Spark分布式算子生成用户、物品的数值索引,同时保留和单机版一致的映射字典输出:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, udf from pyspark.sql.types import IntegerType import pickle spark = SparkSession.builder.appName("item_sim_calc").getOrCreate() # 读取原始全量评分数据 raw_df = spark.read.parquet("hdfs:///path/to/your/user_item.parquet") # 分布式生成用户索引映射 user2num = raw_df.select("user").distinct().rdd.map(lambda x: x[0]).zipWithIndex().collectAsMap() n_user = len(user2num) user_bc = spark.sparkContext.broadcast(user2num) # 分布式生成物品索引映射 item2num = raw_df.select("item").distinct().rdd.map(lambda x: x[0]).zipWithIndex().collectAsMap() n_item = len(item2num) item_bc = spark.sparkContext.broadcast(item2num) # 生成反向映射字典 num2user = {v:k for k,v in user2num.items()} num2item = {v:k for k,v in item2num.items()} # 按原有格式保存映射字典,和单机逻辑完全兼容 dict_list = [ ("num2item.pickle", num2item), ("item2num.pickle", item2num), ("num2user.pickle", num2user), ("user2num.pickle", user2num) ] for path, obj in dict_list: with open(path, "wb") as f: pickle.dump(obj, f, protocol=pickle.HIGHEST_PROTOCOL) # 将原始字符串ID替换为数值索引 @udf(returnType=IntegerType()) def map_user(uid): return user_bc.value[uid] @udf(returnType=IntegerType()) def map_item(iid): return item_bc.value[iid] indexed_df = raw_df.withColumn("user_idx", map_user(col("user")))\ .withColumn("item_idx", map_item(col("item")))\ .select("user_idx", "item_idx", "rating")
注意:如果用户/物品量级超过1亿,不要直接collect映射字典到驱动节点,可将映射关系存储为Parquet表,通过join的方式完成ID转换,避免驱动节点OOM。
步骤2:分布式计算物品相似度矩阵
不要将全量索引数据拉取到单机构建scipy矩阵,直接用Spark分布式计算A^T * A。推荐直接使用Spark MLlib内置的分布式矩阵相似度计算接口,底层做了分块优化,性能比手写RDD逻辑高3~10倍:
from pyspark.mllib.linalg.distributed import CoordinateMatrix, MatrixEntry # 转换为分布式坐标矩阵 entries = indexed_df.rdd.map(lambda x: MatrixEntry(x[0], x[1], float(x[2]))) coord_mat = CoordinateMatrix(entries, numRows=n_user, numCols=n_item) # 计算列相似度,本质就是A^T * A,返回结果为物品x物品的相似度矩阵 item_sim_mat = coord_mat.toRowMatrix().columnSimilarities() # 提取非零元素格式为((item_i, item_j), sim_value) sim_entries = item_sim_mat.entries.map(lambda x: ((x.i, x.j), x.value)).cache()
如果需要手动实现矩阵乘法逻辑,可按“按用户分组生成物品对、聚合求和”的逻辑编写RDD算子,注意对交互物品数过长的用户做截断,避免笛卡尔积数据爆炸。
步骤3:生成scipy稀疏npz文件
所有非零相似度结果计算完成后,分批收集到驱动节点,拼接为csr_matrix后保存,输出格式和单机版完全一致:
import numpy as np from scipy.sparse import csr_matrix, save_npz rows = [] cols = [] data = [] # 收集所有非零元素 for (i,j), val in sim_entries.collect(): rows.append(i) cols.append(j) data.append(val) # 构建csr稀疏矩阵 B = csr_matrix((data, (rows, cols)), shape=(n_item, n_item)) # 保存为npz文件 save_npz("item-item.npz", B) spark.stop()
如果最终物品相似度矩阵非零元素量级过大(例如超过10亿),不要直接collect全量结果到驱动节点,可将sim_entries按物品ID分块存储为Parquet文件,后续按需加载构建局部稀疏矩阵使用。
内容的提问来源于stack exchange,提问作者jason

