Grunderna i Big Data med PySpark
Upendra Devisetty
Science Analyst, CyVerse
Kollaborativ filtrering handlar om att hitta användare med gemensamma intressen
Kollaborativ filtrering används ofta i rekommendationssystem
Metoder för kollaborativ filtrering:
Användare–användare-filtrering: Hittar användare som liknar målanvändaren
Objekt–objekt-filtrering: Hittar och rekommenderar objekt som liknar dem målanvändaren interagerat med
Klassen Rating är ett omslag runt en tupel (användare, produkt och betyg)
Användbar för att tolka RDD och skapa en tupel med användare, produkt och betyg
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)
Att dela upp data i tränings- och testmängder är viktigt för att utvärdera prediktiva modeller
Vanligtvis tilldelas en större del av data till träning än till testning
PySparks metod randomSplit() delar upp data slumpmässigt utifrån angivna vikter och returnerar flera RDD:er
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]
Algoritmen Alternating Least Squares (ALS) i spark.mllib möjliggör kollaborativ filtrering
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)
Metoden predictAll() returnerar en lista med förutsagda betyg för angivna användar- och produktpar
Metoden tar emot en RDD utan betyg och genererar betygen
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 är medelvärdet av kvadraten av (actual rating - predicted rating)
MSE = rates_preds.map(lambda r: (r[1][0] - r[1][1])**2).mean()
Grunderna i Big Data med PySpark