협업 필터링
협업 필터링 (Collaborative Filtering)
협업 필터링은 추천 시스템에서 흔히 쓰이는 기법이에요. 이 기법은 사용자-아이템 연관 행렬에서 누락된 항목을 채우는 것을 목표로 해요. spark.mllib는 현재 모델 기반 협업 필터링을 지원하는데, 여기서 사용자와 상품은 적은 수의 잠재 요인(latent factors)으로 표현되고, 그 요인들로 누락된 항목을 예측할 수 있어요. spark.mllib는 이 잠재 요인을 학습하기 위해 교대 최소 제곱(ALS, alternating least squares) 알고리즘을 사용해요.
spark.mllib의 구현에는 다음과 같은 파라미터가 있어요.
numBlocks— 계산을 병렬화하는 데 쓰는 블록 수 (자동 구성하려면 -1로 설정)rank— 사용할 피처 수 (잠재 요인의 수라고도 불러요)iterations— 실행할 ALS 반복 횟수. ALS는 보통 20회 이하의 반복으로 합리적인 해에 수렴해요.lambda— ALS에서 정규화 파라미터를 지정해요.implicitPrefs— 명시적 피드백(explicit feedback) ALS 변형을 쓸지, 아니면 암묵적 피드백(implicit feedback) 데이터에 맞게 조정된 변형을 쓸지 지정해요.alpha— 암묵적 피드백 변형에 적용되는 파라미터로, 선호도 관측에 대한 기준(baseline) 신뢰도를 결정해요.
본문
명시적 vs 암묵적 피드백
행렬 분해 기반 협업 필터링의 표준 접근법은 사용자-아이템 행렬의 항목을 사용자가 아이템에 준 명시적 선호도로 다뤄요. 예를 들어 사용자가 영화에 평점을 매기는 경우죠.
많은 실제 사례에서는 암묵적 피드백(예: 조회, 클릭, 구매, 좋아요, 공유 등)만 접근 가능한 경우가 흔해요. spark.mllib에서 이런 데이터를 다루는 방식은 Collaborative Filtering for Implicit Feedback Datasets 논문에서 가져왔어요. 기본적으로 평점 행렬을 직접 모델링하는 대신, 사용자 행동 관측의 강도(예: 클릭 수, 누군가 영화를 시청한 누적 시간)를 나타내는 숫자로 데이터를 다뤄요. 그 숫자들은 아이템에 주어진 명시적 평점이 아니라, 관측된 사용자 선호도에 대한 신뢰도 수준과 연결돼요. 그러면 모델은 사용자가 아이템에 기대하는 선호도를 예측하는 데 쓸 수 있는 잠재 요인을 찾으려 해요.
정규화 파라미터의 스케일링
v1.1부터 각 최소 제곱 문제를 풀 때 정규화 파라미터 lambda를, 사용자 요인을 갱신할 때는 사용자가 발생시킨 평점 수로, 상품 요인을 갱신할 때는 상품이 받은 평점 수로 스케일링해요. 이 방식을 "ALS-WR"이라고 부르며 Large-Scale Parallel Collaborative Filtering for the Netflix Prize 논문에서 다뤄요. 이 방식은 lambda가 데이터셋의 규모에 덜 의존하게 만들어서, 샘플링된 부분 집합에서 학습한 최적 파라미터를 전체 데이터셋에 적용해도 비슷한 성능을 기대할 수 있게 해줘요.
예제
다음 예제에서는 평점 데이터를 불러와요. 각 행은 사용자, 상품, 평점으로 구성돼요. 평점이 명시적이라고 가정하는 기본 ALS.train() 메서드를 사용하고, 평점 예측의 평균 제곱 오차(MSE)를 측정해 추천을 평가해요.
API에 대한 자세한 내용은 ALS Python 문서를 참고하세요.
from pyspark.mllib.recommendation import ALS, MatrixFactorizationModel, Rating
# Load and parse the data
data = sc.textFile("data/mllib/als/test.data")
ratings = data.map(lambda l: l.split(','))\
.map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))
# Build the recommendation model using Alternating Least Squares
rank = 10
numIterations = 10
model = ALS.train(ratings, rank, numIterations)
# Evaluate the model on training data
testdata = ratings.map(lambda p: (p[0], p[1]))
predictions = model.predictAll(testdata).map(lambda r: ((r[0], r[1]), r[2]))
ratesAndPreds = ratings.map(lambda r: ((r[0], r[1]), r[2])).join(predictions)
MSE = ratesAndPreds.map(lambda r: (r[1][0] - r[1][1])**2).mean()
print("Mean Squared Error = " + str(MSE))
# Save and load model
model.save(sc, "target/tmp/myCollaborativeFilter")
sameModel = MatrixFactorizationModel.load(sc, "target/tmp/myCollaborativeFilter")
전체 예제 코드는 Spark 저장소의 "examples/src/main/python/mllib/recommendation_example.py"에서 찾을 수 있어요.
평점 행렬이 다른 정보 원천에서 파생됐다면(즉, 다른 신호에서 추론됐다면) trainImplicit 메서드를 쓰면 더 좋은 결과를 얻을 수 있어요.
# Build the recommendation model using Alternating Least Squares based on implicit ratings
model = ALS.trainImplicit(ratings, rank, numIterations, alpha=0.01)
다음 예제에서도 평점 데이터를 불러와요. 각 행은 사용자, 상품, 평점으로 구성돼요. 평점이 명시적이라고 가정하는 기본 ALS.train() 메서드를 사용하고, 평점 예측의 평균 제곱 오차를 측정해 추천 모델을 평가해요.
API에 대한 자세한 내용은 ALS Scala 문서를 참고하세요.
import org.apache.spark.mllib.recommendation.ALS
import org.apache.spark.mllib.recommendation.MatrixFactorizationModel
import org.apache.spark.mllib.recommendation.Rating
// Load and parse the data
val data = sc.textFile("data/mllib/als/test.data")
val ratings = data.map(_.split(',') match { case Array(user, item, rate) =>
Rating(user.toInt, item.toInt, rate.toDouble)
})
// Build the recommendation model using ALS
val rank = 10
val numIterations = 10
val model = ALS.train(ratings, rank, numIterations, 0.01)
// Evaluate the model on rating data
val usersProducts = ratings.map { case Rating(user, product, rate) =>
(user, product)
}
val predictions =
model.predict(usersProducts).map { case Rating(user, product, rate) =>
((user, product), rate)
}
val ratesAndPreds = ratings.map { case Rating(user, product, rate) =>
((user, product), rate)
}.join(predictions)
val MSE = ratesAndPreds.map { case ((user, product), (r1, r2)) =>
val err = (r1 - r2)
err * err
}.mean()
println(s"Mean Squared Error = $MSE")
// Save and load model
model.save(sc, "target/tmp/myCollaborativeFilter")
val sameModel = MatrixFactorizationModel.load(sc, "target/tmp/myCollaborativeFilter")
전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/mllib/RecommendationExample.scala"에서 찾을 수 있어요.
평점 행렬이 다른 정보 원천에서 파생됐다면(즉, 다른 신호에서 추론됐다면) trainImplicit 메서드를 쓰면 더 좋은 결과를 얻을 수 있어요.
val alpha = 0.01
val lambda = 0.01
val model = ALS.trainImplicit(ratings, rank, numIterations, lambda, alpha)
MLlib의 모든 메서드는 Java 친화적 타입을 사용하므로, Scala에서 쓰는 것과 같은 방식으로 Java에서도 import해서 호출할 수 있어요. 유일한 주의점은 메서드가 Scala RDD 객체를 받는데, Spark Java API는 별도의 JavaRDD 클래스를 사용한다는 점이에요. JavaRDD 객체에서 .rdd()를 호출하면 Java RDD를 Scala RDD로 변환할 수 있어요. Scala의 예제와 동등한 자체 포함 애플리케이션 예제는 아래와 같아요.
API에 대한 자세한 내용은 ALS Java 문서를 참고하세요.
import scala.Tuple2;
import org.apache.spark.api.java.*;
import org.apache.spark.mllib.recommendation.ALS;
import org.apache.spark.mllib.recommendation.MatrixFactorizationModel;
import org.apache.spark.mllib.recommendation.Rating;
import org.apache.spark.SparkConf;
SparkConf conf = new SparkConf().setAppName("Java Collaborative Filtering Example");
JavaSparkContext jsc = new JavaSparkContext(conf);
// Load and parse the data
String path = "data/mllib/als/test.data";
JavaRDD<String> data = jsc.textFile(path);
JavaRDD<Rating> ratings = data.map(s -> {
String[] sarray = s.split(",");
return new Rating(Integer.parseInt(sarray[0]),
Integer.parseInt(sarray[1]),
Double.parseDouble(sarray[2]));
});
// Build the recommendation model using ALS
int rank = 10;
int numIterations = 10;
MatrixFactorizationModel model = ALS.train(JavaRDD.toRDD(ratings), rank, numIterations, 0.01);
// Evaluate the model on rating data
JavaRDD<Tuple2<Object, Object>> userProducts =
ratings.map(r -> new Tuple2<>(r.user(), r.product()));
JavaPairRDD<Tuple2<Integer, Integer>, Double> predictions = JavaPairRDD.fromJavaRDD(
model.predict(JavaRDD.toRDD(userProducts)).toJavaRDD()
.map(r -> new Tuple2<>(new Tuple2<>(r.user(), r.product()), r.rating()))
);
JavaRDD<Tuple2<Double, Double>> ratesAndPreds = JavaPairRDD.fromJavaRDD(
ratings.map(r -> new Tuple2<>(new Tuple2<>(r.user(), r.product()), r.rating())))
.join(predictions).values();
double MSE = ratesAndPreds.mapToDouble(pair -> {
double err = pair._1() - pair._2();
return err * err;
}).mean();
System.out.println("Mean Squared Error = " + MSE);
// Save and load model
model.save(jsc.sc(), "target/tmp/myCollaborativeFilter");
MatrixFactorizationModel sameModel = MatrixFactorizationModel.load(jsc.sc(),
"target/tmp/myCollaborativeFilter");
전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/mllib/JavaRecommendationExample.java"에서 찾을 수 있어요.
위 애플리케이션을 실행하려면 Spark Quick Start 가이드의 자체 포함 애플리케이션(Self-Contained Applications) 섹션에 있는 안내를 따라요. 빌드 파일에 spark-mllib를 의존성으로 반드시 포함해야 해요.
튜토리얼
Spark Summit 2014의 training exercises에는 spark.mllib로 개인화된 영화 추천을 다루는 실습 튜토리얼이 포함돼 있어요.
더 알아보기 (Learn more)
- MLlib Main Guide — DataFrame 기반 ML API의 전반적인 내용.
- MLlib: RDD 기반 API 가이드 —
spark.mllib패키지의 각 가이드 모음. - ML 파이프라인 — 협업 필터링 모델을 파이프라인에 통합하는 방법.