협업 필터링

협업 필터링 (Collaborative Filtering)

추천 시스템에서 널리 쓰이는 협업 필터링(Collaborative Filtering) 기법을 소개하는 문서예요. spark.ml이 지원하는 모델 기반 협업 필터링과 ALS(교대 최소제곱, Alternating Least Squares) 알고리즘, 그리고 명시적/암시적 피드백 처리, 정규화 파라미터 스케일링, 콜드 스타트 전략을 Python·Scala·Java·R 예제와 함께 알아볼게요.

출처: 문서

본문

협업 필터링 (Collaborative filtering)

협업 필터링(Collaborative filtering)은 추천 시스템에서 흔히 사용돼요. 이 기법들은 사용자-아이템 연관 행렬(user-item association matrix)의 누락된 항목을 채우는 것을 목표로 해요. spark.ml은 현재 모델 기반 협업 필터링을 지원하는데, 여기서 사용자와 상품은 누락된 항목 예측에 사용할 수 있는 작은 잠재 요인(latent factor) 집합으로 표현돼요. spark.ml은 이 잠재 요인들을 학습하기 위해 교대 최소제곱(ALS, alternating least squares) 알고리즘을 사용해요. spark.ml의 구현은 다음과 같은 파라미터를 가져요:

  • numBlocks — 병렬화를 위해 사용자와 아이템을 분할할 블록 수예요 (기본값 10).
  • rank — 모델의 잠재 요인 수예요 (기본값 10).
  • maxIter — 실행할 최대 반복 횟수예요 (기본값 10).
  • regParam — ALS에서 정규화 파라미터를 지정해요 (기본값 1.0).
  • implicitPrefs — 명시적 피드백 ALS 변형을 사용할지, 아니면 암시적 피드백 데이터에 맞게 조정된 변형을 사용할지 지정해요 (기본값 false, 즉 명시적 피드백 사용).
  • alpha — 암시적 피드백 ALS 변형에 적용되는 파라미터로, 선호 관찰에 대한 기준 신뢰도(baseline confidence)를 조절해요 (기본값 1.0).
  • nonnegative — 최소제곱에 비음수(nonnegative) 제약을 사용할지 여부를 지정해요 (기본값 false).

참고: DataFrame 기반 ALS API는 현재 user와 item id에 정수(integer)만 지원해요. 사용자·아이템 id 컬럼에는 다른 숫자 타입도 지원되지만, id는 정수 값 범위 내에 있어야 해요.

명시적 vs. 암시적 피드백 (Explicit vs. implicit feedback)

행렬 분해 기반 협업 필터링의 표준 접근 방식은 사용자-아이템 행렬의 항목을 사용자가 아이템에 준 명시적 선호(예: 사용자가 영화에 주는 평점)로 취급해요.

실세계의 많은 사용 사례에서는 암시적 피드백(예: 조회수, 클릭, 구매, 좋아요, 공유 등)만 접근 가능한 경우가 흔해요. spark.ml이 이런 데이터를 다루는 방식은 Collaborative Filtering for Implicit Feedback Datasets에서 가져왔어요. 기본적으로 평점 행렬을 직접 모델링하는 대신, 이 접근 방식은 데이터를 사용자 행동 관찰의 강도(예: 클릭 수, 또는 누군가 영화를 시청한 누적 시간)를 나타내는 숫자로 취급해요. 그런 다음 그 숫자들을 아이템에 주어진 명시적 평점이 아니라, 관찰된 사용자 선호에 대한 신뢰도 수준과 연결해요. 이후 모델은 사용자가 아이템에 대해 기대하는 선호도를 예측하는 데 사용할 수 있는 잠재 요인을 찾으려고 해요.

정규화 파라미터 스케일링 (Scaling of the regularization parameter)

각 최소제곱 문제를 풀 때 사용자 요인을 갱신할 때는 사용자가 생성한 평점 수로, 상품 요인을 갱신할 때는 상품이 받은 평점 수로 정규화 파라미터 regParam을 스케일링해요. 이 접근 방식은 "ALS-WR"이라 불리며 "Large-Scale Parallel Collaborative Filtering for the Netflix Prize" 논문에서 논의돼요. 이는 regParam이 데이터셋 규모에 덜 의존하게 만들어, 샘플링된 부분집합에서 학습한 최적 파라미터를 전체 데이터셋에 적용해도 비슷한 성능을 기대할 수 있어요.

콜드 스타트 전략 (Cold-start strategy)

ALSModel로 예측할 때, 테스트 데이터셋에서 모델 학습 시점에 존재하지 않았던 사용자 및/또는 아이템을 만나는 것이 흔해요. 이는 주로 두 가지 시나리오에서 발생해요:

  • 프로덕션 환경에서, 평점 이력이 없고 모델이 학습하지 않은 새 사용자나 새 아이템인 경우 (이것이 바로 "콜드 스타트 문제(cold start problem)"예요).
  • 교차 검증 중에, 데이터가 학습 세트와 평가 세트로 분할되는 경우예요. Spark의 CrossValidatorTrainValidationSplit에서처럼 단순한 무작위 분할을 사용할 때, 평가 세트에 학습 세트에 없는 사용자 및/또는 아이템이 나타나는 것은 실제로 매우 흔해요.

기본적으로 Spark는 ALSModel.transform 중에 사용자 및/또는 아이템 요인이 모델에 없으면 NaN 예측값을 할당해요. 이는 프로덕션 시스템에서 유용할 수 있는데, 새 사용자나 새 아이템임을 나타내므로 시스템이 예측으로 사용할 폴백(fallback)을 결정할 수 있게 해줘요.

하지만 교차 검증 중에는 바람직하지 않아요. NaN 예측값이 있으면 평가 지표 결과도 NaN이 되기 때문이에요 (예: RegressionEvaluator 사용 시). 이러면 모델 선택이 불가능해져요.

Spark는 사용자가 coldStartStrategy 파라미터를 "drop"으로 설정해, 예측 DataFrame에서 NaN 값을 포함하는 행을 제거할 수 있게 해줘요. 평가 지표는 그러면 NaN이 아닌 데이터에 대해 계산되어 유효해져요. 이 파라미터의 사용법은 아래 예제에 나와 있어요.

참고: 현재 지원되는 콜드 스타트 전략은 "nan"(위에서 언급한 기본 동작)과 "drop"이에요. 향후 더 많은 전략이 지원될 수 있어요.

예제 (Examples)

다음 예제에서는 MovieLens 데이터셋에서 평점 데이터를 로드해요. 각 행은 사용자, 영화, 평점, 타임스탬프로 구성돼요. 그런 다음 기본적으로 평점이 명시적이라고 가정하는(implicitPrefsFalse) ALS 모델을 학습시켜요. 평점 예측의 제곱근 평균 제곱 오차(RMSE)를 측정해 추천 모델을 평가해요.

API에 대한 자세한 내용은 ALS Python 문서를 참고하세요.

from pyspark.ml.evaluation import RegressionEvaluator
from pyspark.ml.recommendation import ALS
from pyspark.sql import Row

lines = spark.read.text("data/mllib/als/sample_movielens_ratings.txt").rdd
parts = lines.map(lambda row: row.value.split("::"))
ratingsRDD = parts.map(lambda p: Row(userId=int(p[0]), movieId=int(p[1]),
                                     rating=float(p[2]), timestamp=int(p[3])))
ratings = spark.createDataFrame(ratingsRDD)
(training, test) = ratings.randomSplit([0.8, 0.2])

# Build the recommendation model using ALS on the training data
# Note we set cold start strategy to 'drop' to ensure we don't get NaN evaluation metrics
als = ALS(maxIter=5, regParam=0.01, userCol="userId", itemCol="movieId", ratingCol="rating",
          coldStartStrategy="drop")
model = als.fit(training)

# Evaluate the model by computing the RMSE on the test data
predictions = model.transform(test)
evaluator = RegressionEvaluator(metricName="rmse", labelCol="rating",
                                predictionCol="prediction")
rmse = evaluator.evaluate(predictions)
print("Root-mean-square error = " + str(rmse))

# Generate top 10 movie recommendations for each user
userRecs = model.recommendForAllUsers(10)
# Generate top 10 user recommendations for each movie
movieRecs = model.recommendForAllItems(10)

# Generate top 10 movie recommendations for a specified set of users
users = ratings.select(als.getUserCol()).distinct().limit(3)
userSubsetRecs = model.recommendForUserSubset(users, 10)
# Generate top 10 user recommendations for a specified set of movies
movies = ratings.select(als.getItemCol()).distinct().limit(3)
movieSubSetRecs = model.recommendForItemSubset(movies, 10)

전체 예제 코드는 Spark 저장소의 examples/src/main/python/ml/als_example.py 에서 찾을 수 있어요.

평점 행렬이 다른 정보 소스에서 파생된 경우(즉, 다른 신호에서 추론된 경우) 더 나은 결과를 얻기 위해 implicitPrefsTrue로 설정할 수 있어요:

als = ALS(maxIter=5, regParam=0.01, implicitPrefs=True,
          userCol="userId", itemCol="movieId", ratingCol="rating")

다음 예제에서는 MovieLens 데이터셋에서 평점 데이터를 로드해요. 각 행은 사용자, 영화, 평점, 타임스탬프로 구성돼요. 그런 다음 기본적으로 평점이 명시적이라고 가정하는(implicitPrefsfalse) ALS 모델을 학습시켜요. 평점 예측의 제곱근 평균 제곱 오차(RMSE)를 측정해 추천 모델을 평가해요.

API에 대한 자세한 내용은 ALS Scala 문서를 참고하세요.

import org.apache.spark.ml.evaluation.RegressionEvaluator
import org.apache.spark.ml.recommendation.ALS

case class Rating(userId: Int, movieId: Int, rating: Float, timestamp: Long)
def parseRating(str: String): Rating = {
  val fields = str.split("::")
  assert(fields.size == 4)
  Rating(fields(0).toInt, fields(1).toInt, fields(2).toFloat, fields(3).toLong)
}

val ratings = spark.read.textFile("data/mllib/als/sample_movielens_ratings.txt")
  .map(parseRating)
  .toDF()
val Array(training, test) = ratings.randomSplit(Array(0.8, 0.2))

// Build the recommendation model using ALS on the training data
val als = new ALS()
  .setMaxIter(5)
  .setRegParam(0.01)
  .setUserCol("userId")
  .setItemCol("movieId")
  .setRatingCol("rating")
val model = als.fit(training)

// Evaluate the model by computing the RMSE on the test data
// Note we set cold start strategy to 'drop' to ensure we don't get NaN evaluation metrics
model.setColdStartStrategy("drop")
val predictions = model.transform(test)

val evaluator = new RegressionEvaluator()
  .setMetricName("rmse")
  .setLabelCol("rating")
  .setPredictionCol("prediction")
val rmse = evaluator.evaluate(predictions)
println(s"Root-mean-square error = $rmse")

// Generate top 10 movie recommendations for each user
val userRecs = model.recommendForAllUsers(10)
// Generate top 10 user recommendations for each movie
val movieRecs = model.recommendForAllItems(10)

// Generate top 10 movie recommendations for a specified set of users
val users = ratings.select(als.getUserCol).distinct().limit(3)
val userSubsetRecs = model.recommendForUserSubset(users, 10)
// Generate top 10 user recommendations for a specified set of movies
val movies = ratings.select(als.getItemCol).distinct().limit(3)
val movieSubSetRecs = model.recommendForItemSubset(movies, 10)

전체 예제 코드는 Spark 저장소의 examples/src/main/scala/org/apache/spark/examples/ml/ALSExample.scala 에서 찾을 수 있어요.

평점 행렬이 다른 정보 소스에서 파생된 경우(즉, 다른 신호에서 추론된 경우) 더 나은 결과를 얻기 위해 implicitPrefstrue로 설정할 수 있어요:

val als = new ALS()
  .setMaxIter(5)
  .setRegParam(0.01)
  .setImplicitPrefs(true)
  .setUserCol("userId")
  .setItemCol("movieId")
  .setRatingCol("rating")

다음 예제에서는 MovieLens 데이터셋에서 평점 데이터를 로드해요. 각 행은 사용자, 영화, 평점, 타임스탬프로 구성돼요. 그런 다음 기본적으로 평점이 명시적이라고 가정하는(implicitPrefsfalse) ALS 모델을 학습시켜요. 평점 예측의 제곱근 평균 제곱 오차(RMSE)를 측정해 추천 모델을 평가해요.

API에 대한 자세한 내용은 ALS Java 문서를 참고하세요.

import java.io.Serializable;

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.ml.evaluation.RegressionEvaluator;
import org.apache.spark.ml.recommendation.ALS;
import org.apache.spark.ml.recommendation.ALSModel;

public static class Rating implements Serializable {
  private int userId;
  private int movieId;
  private float rating;
  private long timestamp;

  public Rating() {}

  public Rating(int userId, int movieId, float rating, long timestamp) {
    this.userId = userId;
    this.movieId = movieId;
    this.rating = rating;
    this.timestamp = timestamp;
  }

  public int getUserId() {
    return userId;
  }

  public int getMovieId() {
    return movieId;
  }

  public float getRating() {
    return rating;
  }

  public long getTimestamp() {
    return timestamp;
  }

  public static Rating parseRating(String str) {
    String[] fields = str.split("::");
    if (fields.length != 4) {
      throw new IllegalArgumentException("Each line must contain 4 fields");
    }
    int userId = Integer.parseInt(fields[0]);
    int movieId = Integer.parseInt(fields[1]);
    float rating = Float.parseFloat(fields[2]);
    long timestamp = Long.parseLong(fields[3]);
    return new Rating(userId, movieId, rating, timestamp);
  }
}

JavaRDD<Rating> ratingsRDD = spark
  .read().textFile("data/mllib/als/sample_movielens_ratings.txt").javaRDD()
  .map(Rating::parseRating);
Dataset<Row> ratings = spark.createDataFrame(ratingsRDD, Rating.class);
Dataset<Row>[] splits = ratings.randomSplit(new double[]{0.8, 0.2});
Dataset<Row> training = splits[0];
Dataset<Row> test = splits[1];

// Build the recommendation model using ALS on the training data
ALS als = new ALS()
  .setMaxIter(5)
  .setRegParam(0.01)
  .setUserCol("userId")
  .setItemCol("movieId")
  .setRatingCol("rating");
ALSModel model = als.fit(training);

// Evaluate the model by computing the RMSE on the test data
// Note we set cold start strategy to 'drop' to ensure we don't get NaN evaluation metrics
model.setColdStartStrategy("drop");
Dataset<Row> predictions = model.transform(test);

RegressionEvaluator evaluator = new RegressionEvaluator()
  .setMetricName("rmse")
  .setLabelCol("rating")
  .setPredictionCol("prediction");
double rmse = evaluator.evaluate(predictions);
System.out.println("Root-mean-square error = " + rmse);

// Generate top 10 movie recommendations for each user
Dataset<Row> userRecs = model.recommendForAllUsers(10);
// Generate top 10 user recommendations for each movie
Dataset<Row> movieRecs = model.recommendForAllItems(10);

// Generate top 10 movie recommendations for a specified set of users
Dataset<Row> users = ratings.select(als.getUserCol()).distinct().limit(3);
Dataset<Row> userSubsetRecs = model.recommendForUserSubset(users, 10);
// Generate top 10 user recommendations for a specified set of movies
Dataset<Row> movies = ratings.select(als.getItemCol()).distinct().limit(3);
Dataset<Row> movieSubSetRecs = model.recommendForItemSubset(movies, 10);

전체 예제 코드는 Spark 저장소의 examples/src/main/java/org/apache/spark/examples/ml/JavaALSExample.java 에서 찾을 수 있어요.

평점 행렬이 다른 정보 소스에서 파생된 경우(즉, 다른 신호에서 추론된 경우) 더 나은 결과를 얻기 위해 implicitPrefstrue로 설정할 수 있어요:

ALS als = new ALS()
  .setMaxIter(5)
  .setRegParam(0.01)
  .setImplicitPrefs(true)
  .setUserCol("userId")
  .setItemCol("movieId")
  .setRatingCol("rating");

자세한 내용은 R API 문서를 참고하세요.

# Load training data
data <- list(list(0, 0, 4.0), list(0, 1, 2.0), list(1, 1, 3.0),
             list(1, 2, 4.0), list(2, 1, 1.0), list(2, 2, 5.0))
df <- createDataFrame(data, c("userId", "movieId", "rating"))
training <- df
test <- df

# Fit a recommendation model using ALS with spark.als
model <- spark.als(training, maxIter = 5, regParam = 0.01, userCol = "userId",
                   itemCol = "movieId", ratingCol = "rating")

# Model summary
summary(model)

# Prediction
predictions <- predict(model, test)
head(predictions)

전체 예제 코드는 Spark 저장소의 examples/src/main/r/ml/als.R 에서 찾을 수 있어요.

더 알아보기 (Learn more)