使用 PySpark 的大数据基础
Upendra Devisetty
Science Analyst, CyVerse
协同过滤用于发现具有共同兴趣的用户
协同过滤常用于推荐系统
协同过滤方法:
用户-用户 协同过滤:查找与目标用户相似的用户
物品-物品 协同过滤:查找并推荐与目标用户已互动物品相似的物品
Rating 类是对三元组(user、product、rating)的封装
用于解析 RDD 并创建 user、product、rating 三元组
from pyspark.mllib.recommendation import Rating
r = Rating(user = 1, product = 2, rating = 5.0)
(r[0], r[1], r[2])
(1, 2, 5.0)
将数据拆分为训练集和测试集对评估预测模型很重要
训练集通常占更大比例
PySpark 的 randomSplit() 根据权重随机拆分并返回多个 RDD
data = sc.parallelize([1, 2, 3, 4, 5, 6, 7, 8, 9, 10])
training, test=data.randomSplit([0.6, 0.4])
training.collect()
test.collect()
[1, 2, 5, 6, 9, 10]
[3, 4, 7, 8]
spark.mllib 中的交替最小二乘(ALS)算法提供协同过滤
ALS.train(ratings, rank, iterations)
r1 = Rating(1, 1, 1.0)
r2 = Rating(1, 2, 2.0)
r3 = Rating(2, 1, 2.0)
ratings = sc.parallelize([r1, r2, r3])
ratings.collect()
[Rating(user=1, product=1, rating=1.0),
Rating(user=1, product=2, rating=2.0),
Rating(user=2, product=1, rating=2.0)]
model = ALS.train(ratings, rank=10, iterations=10)
predictAll() 返回给定用户-物品对的预测评分列表
输入为不含评分的 RDD,用于生成预测评分
unrated_RDD = sc.parallelize([(1, 2), (1, 1)])
predictions = model.predictAll(unrated_RDD)
predictions.collect()
[Rating(user=1, product=1, rating=1.0000278574351853),
Rating(user=1, product=2, rating=1.9890355703778122)]
rates = ratings.map(lambda x: ((x[0], x[1]), x[2]))
rates.collect()
[((1, 1), 1.0), ((1, 2), 2.0), ((2, 1), 2.0)]
preds = predictions.map(lambda x: ((x[0], x[1]), x[2])) preds.collect()[((1, 1), 1.000027857), ((1, 2), 1.9890355703)]
rates_preds = rates.join(preds)
rates_preds.collect()
[((1, 2), (2.0, 1.9890355703)), ((1, 1), (1.0, 1.000027857))]
MSE 是 (实际评分 - 预测评分) 的平方的平均值
MSE = rates_preds.map(lambda r: (r[1][0] - r[1][1])**2).mean()
使用 PySpark 的大数据基础