Collaborative filtering 入門

使用 PySpark 的 Big Data 基礎

Upendra Devisetty

Science Analyst, CyVerse

什麼是 Collaborative filtering?

  • Collaborative filtering 是找出有共同興趣的使用者

  • Collaborative filtering 常用於推薦系統

  • Collaborative filtering 方法:

    • User-User Collaborative filtering:找出與目標使用者相似的使用者

    • Item-Item Collaborative filtering:找出並推薦與目標使用者已互動項目相似的項目

使用 PySpark 的 Big Data 基礎

pyspark.mllib.recommendation 子模組中的 Rating 類別

  • Rating 類別是包住 tuple(user、product、rating)的外層封裝

  • 有助於解析 RDD,建立由 user、product、rating 組成的 tuple

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 的 Big Data 基礎

使用 randomSplit() 切分資料

  • 將資料切成訓練集與測試集,有助於評估預測模型

  • 通常訓練集比例會大於測試集

  • 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]
使用 PySpark 的 Big Data 基礎

Alternating Least Squares(ALS)

  • spark.mllib 的 Alternating Least Squares(ALS)演算法可用於 collaborative filtering

  • 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)
使用 PySpark 的 Big Data 基礎

predictAll()

  • 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)]
使用 PySpark 的 Big Data 基礎

模型評估

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 為 (actual rating - predicted rating) 的平方之平均值。

MSE = rates_preds.map(lambda r: (r[1][0] - r[1][1])**2).mean()
使用 PySpark 的 Big Data 基礎

一起來練習吧!

使用 PySpark 的 Big Data 基礎

Preparing Video For Download...