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

如何在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稀疏矩阵文件。

Spark分布式实现方案

整体逻辑和单机版完全对齐,所有重计算步骤分布式执行,仅在最后阶段将非零结果收集到驱动节点生成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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 09:06:22