클러스터링 - RDD 기반 API

클러스터링 - RDD 기반 API (Clustering)

비지도 학습(unsupervised learning) 문제인 클러스터링을 다루는 문서예요. K-means, 가우시안 혼합(Gaussian mixture), 파워 반복 클러스터링(PIC), 잠재 디리클레 할당(LDA), 이분 K-means(Bisecting k-means), 스트리밍 K-means까지 MLlib이 지원하는 클러스터링 모델들을 Python·Scala·Java 예제와 함께 알아볼게요.

출처: 문서

본문

클러스터링(Clustering)은 어떤 유사성(similarity) 개념에 기반해 엔터티 부분집합들을 서로 묶는 것을 목표로 하는 비지도 학습 문제예요. 클러스터링은 탐색적 분석(exploratory analysis)이나 계층적 지도 학습 파이프라인(각 클러스터마다 별도의 분류기·회귀 모델을 학습시키는)의 구성 요소로 자주 사용돼요.

spark.mllib 패키지는 다음 모델들을 지원해요:

  • K-means
  • 가우시안 혼합 (Gaussian mixture)
  • 파워 반복 클러스터링 (Power iteration clustering, PIC)
  • 잠재 디리클레 할당 (Latent Dirichlet allocation, LDA)
  • 이분 K-means (Bisecting k-means)
  • 스트리밍 K-means (Streaming k-means)

K-means

K-means는 데이터 포인트를 미리 정해진 수의 클러스터로 묶는 가장 흔히 사용되는 클러스터링 알고리즘 중 하나예요. spark.mllib 구현에는 k-means++ 방법의 병렬화 변형인 kmeans||이 포함돼요. spark.mllib의 구현은 다음과 같은 파라미터를 가져요:

  • k는 원하는 클러스터 수예요. 참고로 k보다 적은 수의 클러스터가 반환될 수 있어요. 예를 들어 클러스터링할 서로 다른 포인트가 k개 미만인 경우 그럴 수 있어요.
  • maxIterations는 실행할 최대 반복 횟수예요.
  • initializationMode는 무작위 초기화 또는 k-means||를 통한 초기화를 지정해요.
  • runs 이 파라미터는 Spark 2.0.0 이후로 효과가 없어요.
  • initializationSteps는 k-means|| 알고리즘의 단계 수를 결정해요.
  • epsilon은 k-means가 수렴했다고 간주할 거리 임계값을 결정해요.
  • initialModel은 초기화에 사용되는 선택적인 클러스터 중심 집합이에요. 이 파라미터가 제공되면 한 번의 실행만 수행돼요.

예제 (Examples)

다음 예제들은 PySpark 셸에서 테스트할 수 있어요.

다음 예제에서는 데이터를 로드·파싱한 후 KMeans 객체를 사용해 데이터를 두 클러스터로 묶어요. 원하는 클러스터 수를 알고리즘에 전달해요. 그런 다음 클러스터 내 제곱 오차 합(Within Set Sum of Squared Error, WSSSE)을 계산해요. 이 오차 측정치는 k를 늘리면 줄일 수 있어요. 실제로 최적의 k는 보통 WSSSE 그래프에 "elbow(팔꿈치)"가 생기는 지점이에요.

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

from numpy import array
from math import sqrt

from pyspark.mllib.clustering import KMeans, KMeansModel

# Load and parse the data
data = sc.textFile("data/mllib/kmeans_data.txt")
parsedData = data.map(lambda line: array([float(x) for x in line.split(' ')]))

# Build the model (cluster the data)
clusters = KMeans.train(parsedData, 2, maxIterations=10, initializationMode="random")

# Evaluate clustering by computing Within Set Sum of Squared Errors
def error(point):
    center = clusters.centers[clusters.predict(point)]
    return sqrt(sum([x**2 for x in (point - center)]))

WSSSE = parsedData.map(lambda point: error(point)).reduce(lambda x, y: x + y)
print("Within Set Sum of Squared Error = " + str(WSSSE))

# Save and load model
clusters.save(sc, "target/org/apache/spark/PythonKMeansExample/KMeansModel")
sameModel = KMeansModel.load(sc, "target/org/apache/spark/PythonKMeansExample/KMeansModel")

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

다음 코드 조각들은 spark-shell에서 실행할 수 있어요.

다음 예제에서는 데이터를 로드·파싱한 후 KMeans 객체를 사용해 데이터를 두 클러스터로 묶어요. 원하는 클러스터 수를 알고리즘에 전달해요. 그런 다음 클러스터 내 제곱 오차 합(WSSSE)을 계산해요. 이 오차 측정치는 k를 늘리면 줄일 수 있어요. 실제로 최적의 k는 보통 WSSSE 그래프에 "elbow(팔꿈치)"가 생기는 지점이에요.

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

import org.apache.spark.mllib.clustering.{KMeans, KMeansModel}
import org.apache.spark.mllib.linalg.Vectors

// Load and parse the data
val data = sc.textFile("data/mllib/kmeans_data.txt")
val parsedData = data.map(s => Vectors.dense(s.split(' ').map(_.toDouble))).cache()

// Cluster the data into two classes using KMeans
val numClusters = 2
val numIterations = 20
val clusters = KMeans.train(parsedData, numClusters, numIterations)

// Evaluate clustering by computing Within Set Sum of Squared Errors
val WSSSE = clusters.computeCost(parsedData)
println(s"Within Set Sum of Squared Errors = $WSSSE")

// Save and load model
clusters.save(sc, "target/org/apache/spark/KMeansExample/KMeansModel")
val sameModel = KMeansModel.load(sc, "target/org/apache/spark/KMeansExample/KMeansModel")

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

MLlib의 모든 메서드는 Java 친화적인 타입을 사용하므로 Scala에서와 같은 방식으로 import하고 호출할 수 있어요. 유일한 주의점은 메서드가 Scala RDD 객체를 받는데, Spark Java API는 별도의 JavaRDD 클래스를 사용한다는 점이에요. JavaRDD 객체에서 .rdd()를 호출하면 Java RDD를 Scala RDD로 변환할 수 있어요. Scala에서 제공된 예제와 동등한 독립 실행형 애플리케이션 예제는 아래에 있어요.

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

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.mllib.clustering.KMeans;
import org.apache.spark.mllib.clustering.KMeansModel;
import org.apache.spark.mllib.linalg.Vector;
import org.apache.spark.mllib.linalg.Vectors;

// Load and parse data
String path = "data/mllib/kmeans_data.txt";
JavaRDD<String> data = jsc.textFile(path);
JavaRDD<Vector> parsedData = data.map(s -> {
  String[] sarray = s.split(" ");
  double[] values = new double[sarray.length];
  for (int i = 0; i < sarray.length; i++) {
    values[i] = Double.parseDouble(sarray[i]);
  }
  return Vectors.dense(values);
});
parsedData.cache();

// Cluster the data into two classes using KMeans
int numClusters = 2;
int numIterations = 20;
KMeansModel clusters = KMeans.train(parsedData.rdd(), numClusters, numIterations);

System.out.println("Cluster centers:");
for (Vector center: clusters.clusterCenters()) {
  System.out.println(" " + center);
}
double cost = clusters.computeCost(parsedData.rdd());
System.out.println("Cost: " + cost);

// Evaluate clustering by computing Within Set Sum of Squared Errors
double WSSSE = clusters.computeCost(parsedData.rdd());
System.out.println("Within Set Sum of Squared Errors = " + WSSSE);

// Save and load model
clusters.save(jsc.sc(), "target/org/apache/spark/JavaKMeansExample/KMeansModel");
KMeansModel sameModel = KMeansModel.load(jsc.sc(),
  "target/org/apache/spark/JavaKMeansExample/KMeansModel");

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

가우시안 혼합 (Gaussian mixture)

가우시안 혼합 모델(Gaussian Mixture Model)은 포인트가 각각 고유한 확률을 가진 k개의 가우시안 하위 분포 중 하나에서 추출되는 복합 분포를 나타내요. spark.mllib 구현은 샘플 집합이 주어졌을 때 최대우도 모델(maximum-likelihood model)을 유도하기 위해 기대값 최대화(expectation-maximization) 알고리즘을 사용해요. 이 구현은 다음과 같은 파라미터를 가져요:

  • k는 원하는 클러스터 수예요.
  • convergenceTol은 수렴이 달성되었다고 간주하는 로그-우도의 최대 변화값이에요.
  • maxIterations는 수렴에 도달하지 않고 수행할 최대 반복 횟수예요.
  • initialModel은 EM 알고리즘을 시작할 선택적인 시작점이에요. 이 파라미터를 생략하면 데이터에서 무작위 시작점이 만들어져요.

예제 (Examples)

다음 예제에서는 데이터를 로드·파싱한 후 GaussianMixture 객체를 사용해 데이터를 두 클러스터로 묶어요. 원하는 클러스터 수를 알고리즘에 전달해요. 그런 다음 혼합 모델의 파라미터를 출력해요.

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

from numpy import array

from pyspark.mllib.clustering import GaussianMixture, GaussianMixtureModel

# Load and parse the data
data = sc.textFile("data/mllib/gmm_data.txt")
parsedData = data.map(lambda line: array([float(x) for x in line.strip().split(' ')]))

# Build the model (cluster the data)
gmm = GaussianMixture.train(parsedData, 2)

# Save and load model
gmm.save(sc, "target/org/apache/spark/PythonGaussianMixtureExample/GaussianMixtureModel")
sameModel = GaussianMixtureModel\
    .load(sc, "target/org/apache/spark/PythonGaussianMixtureExample/GaussianMixtureModel")

# output parameters of model
for i in range(2):
    print("weight = ", gmm.weights[i], "mu = ", gmm.gaussians[i].mu,
          "sigma = ", gmm.gaussians[i].sigma.toArray())

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

다음 예제에서는 데이터를 로드·파싱한 후 GaussianMixture 객체를 사용해 데이터를 두 클러스터로 묶어요. 원하는 클러스터 수를 알고리즘에 전달해요. 그런 다음 혼합 모델의 파라미터를 출력해요.

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

import org.apache.spark.mllib.clustering.{GaussianMixture, GaussianMixtureModel}
import org.apache.spark.mllib.linalg.Vectors

// Load and parse the data
val data = sc.textFile("data/mllib/gmm_data.txt")
val parsedData = data.map(s => Vectors.dense(s.trim.split(' ').map(_.toDouble))).cache()

// Cluster the data into two classes using GaussianMixture
val gmm = new GaussianMixture().setK(2).run(parsedData)

// Save and load model
gmm.save(sc, "target/org/apache/spark/GaussianMixtureExample/GaussianMixtureModel")
val sameModel = GaussianMixtureModel.load(sc,
  "target/org/apache/spark/GaussianMixtureExample/GaussianMixtureModel")

// output parameters of max-likelihood model
for (i <- 0 until gmm.k) {
  println("weight=%f\nmu=%s\nsigma=\n%s\n" format
    (gmm.weights(i), gmm.gaussians(i).mu, gmm.gaussians(i).sigma))
}

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

MLlib의 모든 메서드는 Java 친화적인 타입을 사용하므로 Scala에서와 같은 방식으로 import하고 호출할 수 있어요. 유일한 주의점은 메서드가 Scala RDD 객체를 받는데, Spark Java API는 별도의 JavaRDD 클래스를 사용한다는 점이에요. JavaRDD 객체에서 .rdd()를 호출하면 Java RDD를 Scala RDD로 변환할 수 있어요. Scala에서 제공된 예제와 동등한 독립 실행형 애플리케이션 예제는 아래에 있어요.

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

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.mllib.clustering.GaussianMixture;
import org.apache.spark.mllib.clustering.GaussianMixtureModel;
import org.apache.spark.mllib.linalg.Vector;
import org.apache.spark.mllib.linalg.Vectors;

// Load and parse data
String path = "data/mllib/gmm_data.txt";
JavaRDD<String> data = jsc.textFile(path);
JavaRDD<Vector> parsedData = data.map(s -> {
  String[] sarray = s.trim().split(" ");
  double[] values = new double[sarray.length];
  for (int i = 0; i < sarray.length; i++) {
    values[i] = Double.parseDouble(sarray[i]);
  }
  return Vectors.dense(values);
});
parsedData.cache();

// Cluster the data into two classes using GaussianMixture
GaussianMixtureModel gmm = new GaussianMixture().setK(2).run(parsedData.rdd());

// Save and load GaussianMixtureModel
gmm.save(jsc.sc(), "target/org/apache/spark/JavaGaussianMixtureExample/GaussianMixtureModel");
GaussianMixtureModel sameModel = GaussianMixtureModel.load(jsc.sc(),
  "target/org.apache.spark.JavaGaussianMixtureExample/GaussianMixtureModel");

// Output the parameters of the mixture model
for (int j = 0; j < gmm.k(); j++) {
  System.out.printf("weight=%f\nmu=%s\nsigma=\n%s\n",
    gmm.weights()[j], gmm.gaussians()[j].mu(), gmm.gaussians()[j].sigma());
}

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

파워 반복 클러스터링 (Power iteration clustering, PIC)

파워 반복 클러스터링(PIC)은 쌍별 유사점(pairwise similarities)을 에지 속성으로 주어졌을 때 그래프 정점들을 클러스터링하는 확장 가능하고 효율적인 알고리즘이에요. Lin and Cohen, Power Iteration Clustering에 설명되어 있어요. 그래프의 정규화된 친화도 행렬(normalized affinity matrix)의 의사 고유벡터(pseudo-eigenvector)를 파워 반복(power iteration)으로 계산하고, 이를 사용해 정점을 클러스터링해요. spark.mllib에는 GraphX를 백엔드로 사용하는 PIC 구현이 포함돼 있어요. (srcId, dstId, similarity) 튜플들의 RDD를 받아 클러스터링 할당을 담은 모델을 출력해요. 유사도(similarities)는 음수가 아니어야 해요. PIC는 유사도 측정치가 대칭이라고 가정해요. (srcId, dstId) 쌍은 순서와 무관하게 입력 데이터에 최대 한 번만 나타나야 해요. 입력에 쌍이 없으면 그 유사도는 0으로 간주돼요. spark.mllib의 PIC 구현은 다음 (하이퍼)파라미터를 사용해요:

  • k: 클러스터 수
  • maxIterations: 최대 파워 반복 횟수
  • initializationMode: 초기화 모델이에요. 기본값인 "random"(무작위 벡터를 정점 속성으로 사용) 또는 "degree"(정규화된 합 유사도 사용) 중 하나일 수 있어요.

예제 (Examples)

아래에서는 spark.mllib에서 PIC를 사용하는 방법을 보여주는 코드 조각을 제시해요.

PowerIterationClustering이 PIC 알고리즘을 구현해요. 친화도 행렬(affinity matrix)을 나타내는 (srcId: Long, dstId: Long, similarity: Double) 튜플들의 RDD를 받아요. PowerIterationClustering.run을 호출하면 계산된 클러스터링 할당을 담고 있는 PowerIterationClusteringModel을 반환해요.

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

from pyspark.mllib.clustering import PowerIterationClustering, PowerIterationClusteringModel

# Load and parse the data
data = sc.textFile("data/mllib/pic_data.txt")
similarities = data.map(lambda line: tuple([float(x) for x in line.split(' ')]))

# Cluster the data into two classes using PowerIterationClustering
model = PowerIterationClustering.train(similarities, 2, 10)

model.assignments().foreach(lambda x: print(str(x.id) + " -> " + str(x.cluster)))

# Save and load model
model.save(sc, "target/org/apache/spark/PythonPowerIterationClusteringExample/PICModel")
sameModel = PowerIterationClusteringModel\
    .load(sc, "target/org/apache/spark/PythonPowerIterationClusteringExample/PICModel")

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

PowerIterationClustering이 PIC 알고리즘을 구현해요. 친화도 행렬을 나타내는 (srcId: Long, dstId: Long, similarity: Double) 튜플들의 RDD를 받아요. PowerIterationClustering.run을 호출하면 계산된 클러스터링 할당을 담고 있는 PowerIterationClusteringModel을 반환해요.

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

import org.apache.spark.mllib.clustering.PowerIterationClustering

val circlesRdd = generateCirclesRdd(sc, params.k, params.numPoints)
val model = new PowerIterationClustering()
  .setK(params.k)
  .setMaxIterations(params.maxIterations)
  .setInitializationMode("degree")
  .run(circlesRdd)

val clusters = model.assignments.collect().groupBy(_.cluster).transform((_, v) => v.map(_.id))
val assignments = clusters.toList.sortBy { case (k, v) => v.length }
val assignmentsStr = assignments
  .map { case (k, v) =>
    s"$k -> ${v.sorted.mkString("[", ",", "]")}"
  }.mkString(", ")
val sizesStr = assignments.map {
  _._2.length
}.sorted.mkString("(", ",", ")")
println(s"Cluster assignments: $assignmentsStr\ncluster sizes: $sizesStr")

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

PowerIterationClustering이 PIC 알고리즘을 구현해요. 친화도 행렬을 나타내는 (srcId: Long, dstId: Long, similarity: Double) 튜플들의 JavaRDD를 받아요. PowerIterationClustering.run을 호출하면 계산된 클러스터링 할당을 담고 있는 PowerIterationClusteringModel을 반환해요.

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

import org.apache.spark.mllib.clustering.PowerIterationClustering;
import org.apache.spark.mllib.clustering.PowerIterationClusteringModel;

JavaRDD<Tuple3<Long, Long, Double>> similarities = sc.parallelize(Arrays.asList(
  new Tuple3<>(0L, 1L, 0.9),
  new Tuple3<>(1L, 2L, 0.9),
  new Tuple3<>(2L, 3L, 0.9),
  new Tuple3<>(3L, 4L, 0.1),
  new Tuple3<>(4L, 5L, 0.9)));

PowerIterationClustering pic = new PowerIterationClustering()
  .setK(2)
  .setMaxIterations(10);
PowerIterationClusteringModel model = pic.run(similarities);

for (PowerIterationClustering.Assignment a: model.assignments().toJavaRDD().collect()) {
  System.out.println(a.id() + " -> " + a.cluster());
}

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

잠재 디리클레 할당 (Latent Dirichlet allocation, LDA)

잠재 디리클레 할당(LDA)은 텍스트 문서 집합에서 주제를 추론하는 주제 모델(topic model)이에요. LDA는 다음과 같이 클러스터링 알고리즘으로 생각할 수 있어요:

  • 주제는 클러스터 중심에, 문서는 데이터셋의 예(행)에 해당해요.
  • 주제와 문서 모두 피처 공간에 존재하며, 피처 벡터는 단어 개수 벡터(bag of words)예요.
  • 전통적인 거리를 사용해 클러스터링을 추정하는 대신, LDA는 텍스트 문서가 어떻게 생성되는지에 대한 통계 모델에 기반한 함수를 사용해요.

LDA는 setOptimizer 함수를 통해 다양한 추론 알고리즘을 지원해요. EMLDAOptimizer는 우도 함수에 대한 기대값 최대화(expectation-maximization)를 사용해 클러스터링을 학습하고 포괄적인 결과를 산출하는 반면, OnlineLDAOptimizer는 온라인 변분 추론(variational inference)을 위해 반복적인 미니배치 샘플링을 사용하고 일반적으로 메모리 친화적이에요.

LDA는 문서의 집합을 단어 개수 벡터로 받고 다음 파라미터(빌더 패턴으로 설정)를 사용해요:

  • k: 주제(즉, 클러스터 중심)의 수
  • optimizer: LDA 모델 학습에 사용할 옵티마이저로, EMLDAOptimizer 또는 OnlineLDAOptimizer
  • docConcentration: 주제에 대한 문서 분포의 사전(prior)에 대한 Dirichlet 파라미터예요. 값이 클수록 더 매끄러운(smooth) 추론 분포를 유도해요.
  • topicConcentration: 단어(용어)에 대한 주제 분포의 사전에 대한 Dirichlet 파라미터예요. 값이 클수록 더 매끄러운 추론 분포를 유도해요.
  • maxIterations: 반복 횟수의 제한이에요.
  • checkpointInterval: 체크포인팅(Spark 구성에서 설정)을 사용한다면, 이 파라미터는 체크포인트가 생성될 빈도를 지정해요. maxIterations가 크다면 체크포인팅을 사용하는 것이 디스크의 shuffle 파일 크기를 줄이고 장애 복구에 도움이 될 수 있어요.

spark.mllib의 모든 LDA 모델은 다음을 지원해요:

  • describeTopics: 주제를 가장 중요한 용어들과 용어 가중치의 배열로 반환해요.
  • topicsMatrix: 각 컬럼이 주제인 vocabSize x k 행렬을 반환해요.

참고: LDA는 아직 활발히 개발 중인 실험적 기능이에요. 그 결과 일부 기능은 두 옵티마이저/모델 중 하나에서만 사용 가능해요. 현재 분산 모델(distributed model)은 로컬 모델로 변환될 수 있지만, 그 반대는 불가능해요.

아래 논의는 각 옵티마이저/모델 쌍을 별도로 설명할게요.

기대값 최대화 (Expectation Maximization)

EMLDAOptimizerDistributedLDAModel에 구현돼 있어요.

LDA에 제공되는 파라미터에 대해:

  • docConcentration: 대칭 사전(prior)만 지원되므로 제공된 k차원 벡터의 모든 값이 동일해야 해요. 모든 값은 $> 1.0$이어야 해요. Vector(-1)을 제공하면 기본 동작((50 / k) + 1 값을 가진 균일 k차원 벡터)이 돼요.
  • topicConcentration: 대칭 사전만 지원돼요. 값은 $> 1.0$이어야 해요. -1을 제공하면 $0.1 + 1$ 값이 기본으로 사용돼요.
  • maxIterations: 최대 EM 반복 횟수예요.

참고: 충분한 반복 횟수를 수행하는 것이 중요해요. 초기 반복에서 EM은 종종 쓸모없는 주제를 갖지만, 이 주제들은 반복이 더 진행된 후에 극적으로 개선돼요. 데이터셋에 따라 최소 20회, 어쩌면 50-100회 반복을 사용하는 것이 흔히 합리적이에요.

EMLDAOptimizerDistributedLDAModel을 생성하는데, 이는 추론된 주제뿐 아니라 전체 훈련 코퍼스와 훈련 코퍼스의 각 문서에 대한 주제 분포를 저장해요. DistributedLDAModel은 다음을 지원해요:

  • topTopicsPerDocument: 훈련 코퍼스의 각 문서에 대한 상위 주제와 그 가중치
  • topDocumentsPerTopic: 각 주제에 대한 상위 문서와 문서에서 해당 주제의 가중치
  • logPrior: 하이퍼파라미터 docConcentrationtopicConcentration이 주어졌을 때 추정된 주제와 문서-주제 분포의 로그 확률
  • logLikelihood: 추론된 주제와 문서-주제 분포가 주어졌을 때 훈련 코퍼스의 로그 우도

온라인 변분 베이즈 (Online Variational Bayes)

OnlineLDAOptimizerLocalLDAModel에 구현돼 있어요.

LDA에 제공되는 파라미터에 대해:

  • docConcentration: k차원 각각의 Dirichlet 파라미터와 같은 값을 가진 벡터를 전달해 비대칭 사전을 사용할 수 있어요. 값은 $>= 0$이어야 해요. Vector(-1)을 제공하면 기본 동작((1.0 / k) 값을 가진 균일 k차원 벡터)이 돼요.
  • topicConcentration: 대칭 사전만 지원돼요. 값은 $>= 0$이어야 해요. -1을 제공하면 (1.0 / k) 값이 기본으로 사용돼요.
  • maxIterations: 제출할 최대 미니배치 수예요.

또한 OnlineLDAOptimizer는 다음 파라미터를 받아들여요:

  • miniBatchFraction: 각 반복에서 샘플링·사용되는 코퍼스의 비율
  • optimizeDocConcentration: true로 설정하면 각 미니배치 후에 하이퍼파라미터 docConcentration(일명 alpha)의 최대우도 추정을 수행하고, 반환된 LocalLDAModel에 최적화된 docConcentration을 설정해요.
  • tau0kappa: 학습률 감쇠(learning-rate decay)에 사용되며, $(\tau_0 + iter)^{-\kappa}$로 계산돼요. 여기서 $iter$는 현재 반복 횟수예요.

OnlineLDAOptimizer는 추론된 주제만 저장하는 LocalLDAModel을 생성해요. LocalLDAModel은 다음을 지원해요:

  • logLikelihood(documents): 추론된 주제가 주어졌을 때 제공된 documents에 대한 하한(lower bound)을 계산해요.
  • logPerplexity(documents): 추론된 주제가 주어졌을 때 제공된 documents의 혼잡도(perplexity)에 대한 상한(upper bound)을 계산해요.

예제 (Examples)

다음 예제에서는 문서 코퍼스를 나타내는 단어 개수 벡터를 로드해요. 그런 다음 LDA를 사용해 문서에서 세 가지 주제를 추론해요. 원하는 클러스터 수를 알고리즘에 전달해요. 그 다음 주제를 단어에 대한 확률 분포로 출력해요.

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

from pyspark.mllib.clustering import LDA, LDAModel
from pyspark.mllib.linalg import Vectors

# Load and parse the data
data = sc.textFile("data/mllib/sample_lda_data.txt")
parsedData = data.map(lambda line: Vectors.dense([float(x) for x in line.strip().split(' ')]))
# Index documents with unique IDs
corpus = parsedData.zipWithIndex().map(lambda x: [x[1], x[0]]).cache()

# Cluster the documents into three topics using LDA
ldaModel = LDA.train(corpus, k=3)

# Output topics. Each is a distribution over words (matching word count vectors)
print("Learned topics (as distributions over vocab of " + str(ldaModel.vocabSize())
      + " words):")
topics = ldaModel.topicsMatrix()
for topic in range(3):
    print("Topic " + str(topic) + ":")
    for word in range(0, ldaModel.vocabSize()):
        print(" " + str(topics[word][topic]))

# Save and load model
ldaModel.save(sc, "target/org/apache/spark/PythonLatentDirichletAllocationExample/LDAModel")
sameModel = LDAModel\
    .load(sc, "target/org/apache/spark/PythonLatentDirichletAllocationExample/LDAModel")

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

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

import org.apache.spark.mllib.clustering.{DistributedLDAModel, LDA}
import org.apache.spark.mllib.linalg.Vectors

// Load and parse the data
val data = sc.textFile("data/mllib/sample_lda_data.txt")
val parsedData = data.map(s => Vectors.dense(s.trim.split(' ').map(_.toDouble)))
// Index documents with unique IDs
val corpus = parsedData.zipWithIndex().map(_.swap).cache()

// Cluster the documents into three topics using LDA
val ldaModel = new LDA().setK(3).run(corpus)

// Output topics. Each is a distribution over words (matching word count vectors)
println(s"Learned topics (as distributions over vocab of ${ldaModel.vocabSize} words):")
val topics = ldaModel.topicsMatrix
for (topic <- Range(0, 3)) {
  print(s"Topic $topic :")
  for (word <- Range(0, ldaModel.vocabSize)) {
    print(s"${topics(word, topic)}")
  }
  println()
}

// Save and load model.
ldaModel.save(sc, "target/org/apache/spark/LatentDirichletAllocationExample/LDAModel")
val sameModel = DistributedLDAModel.load(sc,
  "target/org/apache/spark/LatentDirichletAllocationExample/LDAModel")

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

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

import scala.Tuple2;

import org.apache.spark.api.java.JavaPairRDD;
import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.mllib.clustering.DistributedLDAModel;
import org.apache.spark.mllib.clustering.LDA;
import org.apache.spark.mllib.clustering.LDAModel;
import org.apache.spark.mllib.linalg.Matrix;
import org.apache.spark.mllib.linalg.Vector;
import org.apache.spark.mllib.linalg.Vectors;

// Load and parse the data
String path = "data/mllib/sample_lda_data.txt";
JavaRDD<String> data = jsc.textFile(path);
JavaRDD<Vector> parsedData = data.map(s -> {
  String[] sarray = s.trim().split(" ");
  double[] values = new double[sarray.length];
  for (int i = 0; i < sarray.length; i++) {
    values[i] = Double.parseDouble(sarray[i]);
  }
  return Vectors.dense(values);
});
// Index documents with unique IDs
JavaPairRDD<Long, Vector> corpus =
  JavaPairRDD.fromJavaRDD(parsedData.zipWithIndex().map(Tuple2::swap));
corpus.cache();

// Cluster the documents into three topics using LDA
LDAModel ldaModel = new LDA().setK(3).run(corpus);

// Output topics. Each is a distribution over words (matching word count vectors)
System.out.println("Learned topics (as distributions over vocab of " + ldaModel.vocabSize()
  + " words):");
Matrix topics = ldaModel.topicsMatrix();
for (int topic = 0; topic < 3; topic++) {
  System.out.print("Topic " + topic + ":");
  for (int word = 0; word < ldaModel.vocabSize(); word++) {
    System.out.print(" " + topics.apply(word, topic));
  }
  System.out.println();
}

ldaModel.save(jsc.sc(),
  "target/org/apache/spark/JavaLatentDirichletAllocationExample/LDAModel");
DistributedLDAModel sameModel = DistributedLDAModel.load(jsc.sc(),
  "target/org/apache/spark/JavaLatentDirichletAllocationExample/LDAModel");

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

이분 K-means (Bisecting k-means)

이분 K-means(Bisecting K-means)는 일반 K-means보다 훨씬 빠를 수 있지만, 보통은 다른 클러스터링 결과를 만들어요.

이분 k-means는 계층적 클러스터링(hierarchical clustering)의 한 종류예요. 계층적 클러스터링은 클러스터의 계층을 구축하려는 클러스터 분석에서 가장 흔히 사용되는 방법 중 하나예요. 계층적 클러스터링 전략은 일반적으로 두 가지 유형으로 나뉘어요:

  • 응집형(Agglomerative): 이는 "상향식(bottom up)" 접근 방식이에요. 각 관측치는 자신만의 클러스터에서 시작하며, 계층을 올라갈수록 클러스터 쌍이 병합돼요.
  • 분할형(Divisive): 이는 "하향식(top down)" 접근 방식이에요. 모든 관측치가 하나의 클러스터에서 시작하며, 계층을 내려갈수록 재귀적으로 분할돼요.

이분 k-means 알고리즘은 분할형 알고리즘의 한 종류예요. MLlib의 구현은 다음과 같은 파라미터를 가져요:

  • k: 원하는 리프 클러스터 수 (기본값: 4). 분할 가능한 리프 클러스터가 없으면 실제 수가 더 작을 수 있어요.
  • maxIterations: 클러스터를 분할하기 위한 최대 k-means 반복 횟수 (기본값: 20)
  • minDivisibleClusterSize: 분할 가능한 클러스터의 최소 포인트 수 (>= 1.0이면) 또는 최소 포인트 비율 (< 1.0이면) (기본값: 1)
  • seed: 랜덤 시드 (기본값: 클래스 이름의 해시 값)

예제 (Examples)

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

from numpy import array

from pyspark.mllib.clustering import BisectingKMeans

# Load and parse the data
data = sc.textFile("data/mllib/kmeans_data.txt")
parsedData = data.map(lambda line: array([float(x) for x in line.split(' ')]))

# Build the model (cluster the data)
model = BisectingKMeans.train(parsedData, 2, maxIterations=5)

# Evaluate clustering
cost = model.computeCost(parsedData)
print("Bisecting K-means Cost = " + str(cost))

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

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

import org.apache.spark.mllib.clustering.BisectingKMeans
import org.apache.spark.mllib.linalg.{Vector, Vectors}

// Loads and parses data
def parse(line: String): Vector = Vectors.dense(line.split(" ").map(_.toDouble))
val data = sc.textFile("data/mllib/kmeans_data.txt").map(parse).cache()

// Clustering the data into 6 clusters by BisectingKMeans.
val bkm = new BisectingKMeans().setK(6)
val model = bkm.run(data)

// Show the compute cost and the cluster centers
println(s"Compute Cost: ${model.computeCost(data)}")
model.clusterCenters.zipWithIndex.foreach { case (center, idx) =>
  println(s"Cluster Center ${idx}: ${center}")
}

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

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

import java.util.Arrays;
import java.util.List;

import org.apache.spark.api.java.JavaRDD;
import org.apache.spark.mllib.clustering.BisectingKMeans;
import org.apache.spark.mllib.clustering.BisectingKMeansModel;
import org.apache.spark.mllib.linalg.Vector;
import org.apache.spark.mllib.linalg.Vectors;

List<Vector> localData = Arrays.asList(
  Vectors.dense(0.1, 0.1),   Vectors.dense(0.3, 0.3),
  Vectors.dense(10.1, 10.1), Vectors.dense(10.3, 10.3),
  Vectors.dense(20.1, 20.1), Vectors.dense(20.3, 20.3),
  Vectors.dense(30.1, 30.1), Vectors.dense(30.3, 30.3)
);
JavaRDD<Vector> data = sc.parallelize(localData, 2);

BisectingKMeans bkm = new BisectingKMeans()
  .setK(4);
BisectingKMeansModel model = bkm.run(data);

System.out.println("Compute Cost: " + model.computeCost(data));

Vector[] clusterCenters = model.clusterCenters();
for (int i = 0; i < clusterCenters.length; i++) {
  Vector clusterCenter = clusterCenters[i];
  System.out.println("Cluster Center " + i + ": " + clusterCenter);
}

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

스트리밍 K-means (Streaming k-means)

데이터가 스트림으로 도착할 때 우리는 클러스터를 동적으로 추정하고, 새 데이터가 도착할 때마다 갱신하고 싶을 수 있어요. spark.mllib는 스트리밍 k-means 클러스터링을 지원하며, 추정치의 감쇠(decay, 또는 "망각(forgetfulness)")를 제어하는 파라미터를 제공해요. 이 알고리즘은 미니배치 k-means 갱신 규칙의 일반화를 사용해요. 각 배치 데이터에 대해 모든 포인트를 가장 가까운 클러스터에 할당하고, 새 클러스터 중심을 계산한 후 다음을 사용해 각 클러스터를 갱신해요:

\begin{equation}
    c_{t+1} = \frac{c_tn_t\alpha + x_tm_t}{n_t\alpha+m_t}
\end{equation}
\begin{equation}
    n_{t+1} = n_t + m_t
\end{equation}

여기서 $c_t$는 클러스터의 이전 중심, $n_t$는 지금까지 클러스터에 할당된 포인트 수, $x_t$는 현재 배치에서의 새 클러스터 중심, $m_t$는 현재 배치에서 클러스터에 추가된 포인트 수예요. 감쇠 인자 $\alpha$는 과거를 무시하는 데 사용할 수 있어요: $\alpha$=1이면 처음부터 모든 데이터가 사용되고, $\alpha$=0이면 가장 최근 데이터만 사용돼요. 이는 지수 가중 이동 평균(exponentially-weighted moving average)과 유사해요.

감쇠는 halfLife 파라미터로 지정할 수 있는데, 이는 시간 t에 획득된 데이터가 시간 t + halfLife까지 그 기여도가 0.5로 떨어지도록 하는 올바른 감쇠 인자 a를 결정해요. 시간 단위는 batches 또는 points로 지정할 수 있으며, 갱신 규칙이 그에 맞게 조정돼요.

예제 (Examples)

이 예제는 스트리밍 데이터에서 클러스터를 추정하는 방법을 보여줘요.

API에 대한 자세한 내용은 StreamingKMeans Python 문서를 참고하세요. StreamingContext에 대한 자세한 내용은 Spark Streaming Programming Guide를 참고하세요.

from pyspark.mllib.linalg import Vectors
from pyspark.mllib.regression import LabeledPoint
from pyspark.mllib.clustering import StreamingKMeans

# we make an input stream of vectors for training,
# as well as a stream of vectors for testing
def parse(lp):
    label = float(lp[lp.find('(') + 1: lp.find(')')])
    vec = Vectors.dense(lp[lp.find('[') + 1: lp.find(']')].split(','))

    return LabeledPoint(label, vec)

trainingData = sc.textFile("data/mllib/kmeans_data.txt")\
    .map(lambda line: Vectors.dense([float(x) for x in line.strip().split(' ')]))

testingData = sc.textFile("data/mllib/streaming_kmeans_data_test.txt").map(parse)

trainingQueue = [trainingData]
testingQueue = [testingData]

trainingStream = ssc.queueStream(trainingQueue)
testingStream = ssc.queueStream(testingQueue)

# We create a model with random clusters and specify the number of clusters to find
model = StreamingKMeans(k=2, decayFactor=1.0).setRandomCenters(3, 1.0, 0)

# Now register the streams for training and testing and start the job,
# printing the predicted cluster assignments on new data points as they arrive.
model.trainOn(trainingStream)

result = model.predictOnValues(testingStream.map(lambda lp: (lp.label, lp.features)))
result.pprint()

ssc.start()
ssc.stop(stopSparkContext=True, stopGraceFully=True)

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

API에 대한 자세한 내용은 StreamingKMeans Scala 문서를 참고하세요. StreamingContext에 대한 자세한 내용은 Spark Streaming Programming Guide를 참고하세요.

import org.apache.spark.mllib.clustering.StreamingKMeans
import org.apache.spark.mllib.linalg.Vectors
import org.apache.spark.mllib.regression.LabeledPoint
import org.apache.spark.streaming.{Seconds, StreamingContext}

val conf = new SparkConf().setAppName("StreamingKMeansExample")
val ssc = new StreamingContext(conf, Seconds(args(2).toLong))

val trainingData = ssc.textFileStream(args(0)).map(Vectors.parse)
val testData = ssc.textFileStream(args(1)).map(LabeledPoint.parse)

val model = new StreamingKMeans()
  .setK(args(3).toInt)
  .setDecayFactor(1.0)
  .setRandomCenters(args(4).toInt, 0.0)

model.trainOn(trainingData)
model.predictOnValues(testData.map(lp => (lp.label, lp.features))).print()

ssc.start()
ssc.awaitTermination()

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

데이터가 담긴 새 텍스트 파일을 추가하면 클러스터 중심이 갱신돼요. 각 학습 포인트는 [x1, x2, x3] 형식이어야 하고, 각 테스트 데이터 포인트는 (y, [x1, x2, x3]) 형식이어야 해요. 여기서 y는 유용한 레이블이나 식별자(예: 실제 카테고리 할당)예요. 텍스트 파일이 /training/data/dir에 놓일 때마다 모델이 갱신돼요. 텍스트 파일이 /testing/data/dir에 놓일 때마다 예측을 볼 수 있어요. 새 데이터로 클러스터 중심이 바뀔 거예요!

더 알아보기 (Learn more)