Fundamentele Big Data cu PySpark
Upendra Devisetty
Science Analyst, CyVerse
Filtrarea colaborativă identifică utilizatori cu interese comune
Este utilizată frecvent în sistemele de recomandare
Abordări ale filtrării colaborative:
Filtrare colaborativă Utilizator-Utilizator: Identifică utilizatori similari cu utilizatorul țintă
Filtrare colaborativă Articol-Articol: Identifică și recomandă articole similare cu cele ale utilizatorului țintă
Clasa Rating este un wrapper pentru tuplu (utilizator, produs și rating)
Utilă pentru parsarea RDD și crearea unui tuplu de utilizator, produs și 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)
Împărțirea datelor în seturi de antrenament și testare este esențială pentru evaluarea modelelor predictive
De obicei, o proporție mai mare din date este alocată antrenamentului față de testare
Metoda randomSplit() din PySpark împarte aleatoriu datele conform ponderilor furnizate și returnează mai multe RDD-uri
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]
Algoritmul Alternating Least Squares (ALS) din spark.mllib oferă filtrare colaborativă
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)
Metoda predictAll() returnează o listă de ratinguri prezise pentru perechile de utilizator și produs furnizate
Metoda primește un RDD fără ratinguri și generează ratingurile corespunzătoare
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 este media pătratelor diferenței (rating actual - rating prezis)
MSE = rates_preds.map(lambda r: (r[1][0] - r[1][1])**2).mean()
Fundamentele Big Data cu PySpark