클러스터링
클러스터링 (Clustering)
MLlib에서 제공하는 클러스터링 알고리즘을 소개하는 페이지예요. 데이터 포인트를 몇 개의 군집(cluster)으로 묶는 데 쓰이는 K-means부터, 토픽 모델링에 쓰이는 LDA, 그래프 클러스터링 PIC까지 다양한 알고리즘을 예제 코드와 함께 설명해 드릴게요. RDD 기반 API의 클러스터링 가이드에도 이 알고리즘들에 대한 관련 정보가 있으니 함께 참고하면 좋아요.
출처: 문서
본문
이 페이지는 MLlib의 클러스터링 알고리즘을 설명해요. RDD 기반 API의 클러스터링 가이드에도 이 알고리즘들에 대한 관련 정보가 있습니다.
목차 (Table of Contents)
- K-means
- 잠재 디리클레 할당 LDA (Latent Dirichlet allocation)
- 이분 K-means (Bisecting k-means)
- 가우시안 혼합 모델 GMM (Gaussian Mixture Model)
- 전력 반복 클러스터링 PIC (Power Iteration Clustering)
K-means
k-means는 가장 널리 쓰이는 클러스터링 알고리즘 중 하나로, 데이터 포인트를 미리 정해진 개수의 군집으로 묶어요. MLlib 구현에는 k-means++ 방법의 병렬화 변형인 kmeans||이 포함되어 있습니다.
KMeans는 Estimator로 구현되며 기본 모델로 KMeansModel을 생성합니다.
입력 컬럼 (Input Columns)
| 파라미터 이름 | 타입 | 기본값 | 설명 |
|---|---|---|---|
| featuresCol | Vector | "features" | 피처 벡터 |
출력 컬럼 (Output Columns)
| 파라미터 이름 | 타입 | 기본값 | 설명 |
|---|---|---|---|
| predictionCol | Int | "prediction" | 예측된 클러스터 중심 |
예제 (Examples)
Python:
자세한 내용은 [Python API 문서](api/python/reference/api/pyspark.ml.clustering.KMeans.html)를 참고하세요.
from pyspark.ml.clustering import KMeans
from pyspark.ml.evaluation import ClusteringEvaluator
# Loads data.
dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
# Trains a k-means model.
kmeans = KMeans().setK(2).setSeed(1)
model = kmeans.fit(dataset)
# Make predictions
predictions = model.transform(dataset)
# Evaluate clustering by computing Silhouette score
evaluator = ClusteringEvaluator()
silhouette = evaluator.evaluate(predictions)
print("Silhouette with squared euclidean distance = " + str(silhouette))
# Shows the result.
centers = model.clusterCenters()
print("Cluster Centers: ")
for center in centers:
print(center)
전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/kmeans_example.py"에서 확인할 수 있어요.
Scala:
자세한 내용은 [Scala API 문서](api/scala/org/apache/spark/ml/clustering/KMeans.html)를 참고하세요.
import org.apache.spark.ml.clustering.KMeans
import org.apache.spark.ml.evaluation.ClusteringEvaluator
// Loads data.
val dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
// Trains a k-means model.
val kmeans = new KMeans().setK(2).setSeed(1L)
val model = kmeans.fit(dataset)
// Make predictions
val predictions = model.transform(dataset)
// Evaluate clustering by computing Silhouette score
val evaluator = new ClusteringEvaluator()
val silhouette = evaluator.evaluate(predictions)
println(s"Silhouette with squared euclidean distance = $silhouette")
// Shows the result.
println("Cluster Centers: ")
model.clusterCenters.foreach(println)
전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/KMeansExample.scala"에서 확인할 수 있어요.
Java:
자세한 내용은 [Java API 문서](api/java/org/apache/spark/ml/clustering/KMeans.html)를 참고하세요.
import org.apache.spark.ml.clustering.KMeansModel;
import org.apache.spark.ml.clustering.KMeans;
import org.apache.spark.ml.evaluation.ClusteringEvaluator;
import org.apache.spark.ml.linalg.Vector;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
// Loads data.
Dataset<Row> dataset = spark.read().format("libsvm").load("data/mllib/sample_kmeans_data.txt");
// Trains a k-means model.
KMeans kmeans = new KMeans().setK(2).setSeed(1L);
KMeansModel model = kmeans.fit(dataset);
// Make predictions
Dataset<Row> predictions = model.transform(dataset);
// Evaluate clustering by computing Silhouette score
ClusteringEvaluator evaluator = new ClusteringEvaluator();
double silhouette = evaluator.evaluate(predictions);
System.out.println("Silhouette with squared euclidean distance = " + silhouette);
// Shows the result.
Vector[] centers = model.clusterCenters();
System.out.println("Cluster Centers: ");
for (Vector center: centers) {
System.out.println(center);
}
전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaKMeansExample.java"에서 확인할 수 있어요.
R:
자세한 내용은 [R API 문서](api/R/reference/spark.kmeans.html)를 참고하세요.
# Fit a k-means model with spark.kmeans
t <- as.data.frame(Titanic)
training <- createDataFrame(t)
df_list <- randomSplit(training, c(7,3), 2)
kmeansDF <- df_list[[1]]
kmeansTestDF <- df_list[[2]]
kmeansModel <- spark.kmeans(kmeansDF, ~ Class + Sex + Age + Freq,
k = 3)
# Model summary
summary(kmeansModel)
# Get fitted result from the k-means model
head(fitted(kmeansModel))
# Prediction
kmeansPredictions <- predict(kmeansModel, kmeansTestDF)
head(kmeansPredictions)
전체 예제 코드는 Spark 저장소의 "examples/src/main/r/ml/kmeans.R"에서 확인할 수 있어요.
잠재 디리클레 할당 (Latent Dirichlet allocation, LDA)
LDA는 EMLDAOptimizer와 OnlineLDAOptimizer 둘 다 지원하는 Estimator로 구현되며, 기본 모델로 LDAModel을 생성해요. 고급 사용자는 필요한 경우 EMLDAOptimizer가 생성한 LDAModel을 DistributedLDAModel로 캐스팅할 수 있습니다.
예제 (Examples)
Python:
자세한 내용은 [Python API 문서](api/python/reference/api/pyspark.ml.clustering.LDA.html)를 참고하세요.
from pyspark.ml.clustering import LDA
# Loads data.
dataset = spark.read.format("libsvm").load("data/mllib/sample_lda_libsvm_data.txt")
# Trains a LDA model.
lda = LDA(k=10, maxIter=10)
model = lda.fit(dataset)
ll = model.logLikelihood(dataset)
lp = model.logPerplexity(dataset)
print("The lower bound on the log likelihood of the entire corpus: " + str(ll))
print("The upper bound on perplexity: " + str(lp))
# Describe topics.
topics = model.describeTopics(3)
print("The topics described by their top-weighted terms:")
topics.show(truncate=False)
# Shows the result
transformed = model.transform(dataset)
transformed.show(truncate=False)
전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/lda_example.py"에서 확인할 수 있어요.
Scala:
자세한 내용은 [Scala API 문서](api/scala/org/apache/spark/ml/clustering/LDA.html)를 참고하세요.
import org.apache.spark.ml.clustering.LDA
// Loads data.
val dataset = spark.read.format("libsvm")
.load("data/mllib/sample_lda_libsvm_data.txt")
// Trains a LDA model.
val lda = new LDA().setK(10).setMaxIter(10)
val model = lda.fit(dataset)
val ll = model.logLikelihood(dataset)
val lp = model.logPerplexity(dataset)
println(s"The lower bound on the log likelihood of the entire corpus: $ll")
println(s"The upper bound on perplexity: $lp")
// Describe topics.
val topics = model.describeTopics(3)
println("The topics described by their top-weighted terms:")
topics.show(false)
// Shows the result.
val transformed = model.transform(dataset)
transformed.show(false)
전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/LDAExample.scala"에서 확인할 수 있어요.
Java:
자세한 내용은 [Java API 문서](api/java/org/apache/spark/ml/clustering/LDA.html)를 참고하세요.
import org.apache.spark.ml.clustering.LDA;
import org.apache.spark.ml.clustering.LDAModel;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.SparkSession;
// Loads data.
Dataset<Row> dataset = spark.read().format("libsvm")
.load("data/mllib/sample_lda_libsvm_data.txt");
// Trains a LDA model.
LDA lda = new LDA().setK(10).setMaxIter(10);
LDAModel model = lda.fit(dataset);
double ll = model.logLikelihood(dataset);
double lp = model.logPerplexity(dataset);
System.out.println("The lower bound on the log likelihood of the entire corpus: " + ll);
System.out.println("The upper bound on perplexity: " + lp);
// Describe topics.
Dataset<Row> topics = model.describeTopics(3);
System.out.println("The topics described by their top-weighted terms:");
topics.show(false);
// Shows the result.
Dataset<Row> transformed = model.transform(dataset);
transformed.show(false);
전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaLDAExample.java"에서 확인할 수 있어요.
R:
자세한 내용은 [R API 문서](api/R/reference/spark.lda.html)를 참고하세요.
# Load training data
df <- read.df("data/mllib/sample_lda_libsvm_data.txt", source = "libsvm")
training <- df
test <- df
# Fit a latent dirichlet allocation model with spark.lda
model <- spark.lda(training, k = 10, maxIter = 10)
# Model summary
summary(model)
# Posterior probabilities
posterior <- spark.posterior(model, test)
head(posterior)
# The log perplexity of the LDA model
logPerplexity <- spark.perplexity(model, test)
print(paste0("The upper bound bound on perplexity: ", logPerplexity))
전체 예제 코드는 Spark 저장소의 "examples/src/main/r/ml/lda.R"에서 확인할 수 있어요.
이분 K-means (Bisecting k-means)
이분 K-means는 분할적(divisive, 'top-down') 접근 방식을 사용하는 계층적 클러스터링의 일종이에요. 모든 관측치는 하나의 클러스터에서 시작하고, 계층을 따라 내려가면서 재귀적으로 분할이 수행됩니다.
이분 K-means는 일반 K-means보다 훨씬 빠른 경우가 많지만, 일반적으로는 다른 클러스터링 결과를 만들어요.
BisectingKMeans는 Estimator로 구현되며 기본 모델로 BisectingKMeansModel을 생성합니다.
예제 (Examples)
Python:
자세한 내용은 [Python API 문서](api/python/reference/api/pyspark.ml.clustering.BisectingKMeans.html)를 참고하세요.
from pyspark.ml.clustering import BisectingKMeans
from pyspark.ml.evaluation import ClusteringEvaluator
# Loads data.
dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
# Trains a bisecting k-means model.
bkm = BisectingKMeans().setK(2).setSeed(1)
model = bkm.fit(dataset)
# Make predictions
predictions = model.transform(dataset)
# Evaluate clustering by computing Silhouette score
evaluator = ClusteringEvaluator()
silhouette = evaluator.evaluate(predictions)
print("Silhouette with squared euclidean distance = " + str(silhouette))
# Shows the result.
print("Cluster Centers: ")
centers = model.clusterCenters()
for center in centers:
print(center)
전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/bisecting_k_means_example.py"에서 확인할 수 있어요.
Scala:
자세한 내용은 [Scala API 문서](api/scala/org/apache/spark/ml/clustering/BisectingKMeans.html)를 참고하세요.
import org.apache.spark.ml.clustering.BisectingKMeans
import org.apache.spark.ml.evaluation.ClusteringEvaluator
// Loads data.
val dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
// Trains a bisecting k-means model.
val bkm = new BisectingKMeans().setK(2).setSeed(1)
val model = bkm.fit(dataset)
// Make predictions
val predictions = model.transform(dataset)
// Evaluate clustering by computing Silhouette score
val evaluator = new ClusteringEvaluator()
val silhouette = evaluator.evaluate(predictions)
println(s"Silhouette with squared euclidean distance = $silhouette")
// Shows the result.
println("Cluster Centers: ")
val centers = model.clusterCenters
centers.foreach(println)
전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/BisectingKMeansExample.scala"에서 확인할 수 있어요.
Java:
자세한 내용은 [Java API 문서](api/java/org/apache/spark/ml/clustering/BisectingKMeans.html)를 참고하세요.
import org.apache.spark.ml.clustering.BisectingKMeans;
import org.apache.spark.ml.clustering.BisectingKMeansModel;
import org.apache.spark.ml.evaluation.ClusteringEvaluator;
import org.apache.spark.ml.linalg.Vector;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
// Loads data.
Dataset<Row> dataset = spark.read().format("libsvm").load("data/mllib/sample_kmeans_data.txt");
// Trains a bisecting k-means model.
BisectingKMeans bkm = new BisectingKMeans().setK(2).setSeed(1);
BisectingKMeansModel model = bkm.fit(dataset);
// Make predictions
Dataset<Row> predictions = model.transform(dataset);
// Evaluate clustering by computing Silhouette score
ClusteringEvaluator evaluator = new ClusteringEvaluator();
double silhouette = evaluator.evaluate(predictions);
System.out.println("Silhouette with squared euclidean distance = " + silhouette);
// Shows the result.
System.out.println("Cluster Centers: ");
Vector[] centers = model.clusterCenters();
for (Vector center : centers) {
System.out.println(center);
}
전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaBisectingKMeansExample.java"에서 확인할 수 있어요.
R:
자세한 내용은 [R API 문서](api/R/reference/spark.bisectingKmeans.html)를 참고하세요.
t <- as.data.frame(Titanic)
training <- createDataFrame(t)
# Fit bisecting k-means model with four centers
model <- spark.bisectingKmeans(training, Class ~ Survived, k = 4)
# get fitted result from a bisecting k-means model
fitted.model <- fitted(model, "centers")
# Model summary
head(summary(fitted.model))
# fitted values on training data
fitted <- predict(model, training)
head(select(fitted, "Class", "prediction"))
전체 예제 코드는 Spark 저장소의 "examples/src/main/r/ml/bisectingKmeans.R"에서 확인할 수 있어요.
가우시안 혼합 모델 (Gaussian Mixture Model, GMM)
가우시안 혼합 모델은 포인트들이 각각 고유한 확률을 가진 k개의 가우시안 하위 분포 중 하나에서 추출되는 복합 분포를 나타내요. spark.ml 구현은 주어진 샘플 집합에서 최대우도 모델을 유도하기 위해 기대-최대화(expectation-maximization) 알고리즘을 사용합니다.
GaussianMixture는 Estimator로 구현되며 기본 모델로 GaussianMixtureModel을 생성합니다.
입력 컬럼 (Input Columns)
| 파라미터 이름 | 타입 | 기본값 | 설명 |
|---|---|---|---|
| featuresCol | Vector | "features" | 피처 벡터 |
출력 컬럼 (Output Columns)
| 파라미터 이름 | 타입 | 기본값 | 설명 |
|---|---|---|---|
| predictionCol | Int | "prediction" | 예측된 클러스터 중심 |
| probabilityCol | Vector | "probability" | 각 클러스터의 확률 |
예제 (Examples)
Python:
자세한 내용은 [Python API 문서](api/python/reference/api/pyspark.ml.clustering.GaussianMixture.html)를 참고하세요.
from pyspark.ml.clustering import GaussianMixture
# loads data
dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
gmm = GaussianMixture().setK(2).setSeed(538009335)
model = gmm.fit(dataset)
print("Gaussians shown as a DataFrame: ")
model.gaussiansDF.show(truncate=False)
전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/gaussian_mixture_example.py"에서 확인할 수 있어요.
Scala:
자세한 내용은 [Scala API 문서](api/scala/org/apache/spark/ml/clustering/GaussianMixture.html)를 참고하세요.
import org.apache.spark.ml.clustering.GaussianMixture
// Loads data
val dataset = spark.read.format("libsvm").load("data/mllib/sample_kmeans_data.txt")
// Trains Gaussian Mixture Model
val gmm = new GaussianMixture()
.setK(2)
val model = gmm.fit(dataset)
// output parameters of mixture model model
for (i <- 0 until model.getK) {
println(s"Gaussian $i:\nweight=${model.weights(i)}\n" +
s"mu=${model.gaussians(i).mean}\nsigma=\n${model.gaussians(i).cov}\n")
}
전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/GaussianMixtureExample.scala"에서 확인할 수 있어요.
Java:
자세한 내용은 [Java API 문서](api/java/org/apache/spark/ml/clustering/GaussianMixture.html)를 참고하세요.
import org.apache.spark.ml.clustering.GaussianMixture;
import org.apache.spark.ml.clustering.GaussianMixtureModel;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
// Loads data
Dataset<Row> dataset = spark.read().format("libsvm").load("data/mllib/sample_kmeans_data.txt");
// Trains a GaussianMixture model
GaussianMixture gmm = new GaussianMixture()
.setK(2);
GaussianMixtureModel model = gmm.fit(dataset);
// Output the parameters of the mixture model
for (int i = 0; i < model.getK(); i++) {
System.out.printf("Gaussian %d:\nweight=%f\nmu=%s\nsigma=\n%s\n\n",
i, model.weights()[i], model.gaussians()[i].mean(), model.gaussians()[i].cov());
}
전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaGaussianMixtureExample.java"에서 확인할 수 있어요.
R:
자세한 내용은 [R API 문서](api/R/reference/spark.gaussianMixture.html)를 참고하세요.
# Load training data
df <- read.df("data/mllib/sample_kmeans_data.txt", source = "libsvm")
training <- df
test <- df
# Fit a gaussian mixture clustering model with spark.gaussianMixture
model <- spark.gaussianMixture(training, ~ features, k = 2)
# Model summary
summary(model)
# Prediction
predictions <- predict(model, test)
head(predictions)
전체 예제 코드는 Spark 저장소의 "examples/src/main/r/ml/gaussianMixture.R"에서 확인할 수 있어요.
전력 반복 클러스터링 (Power Iteration Clustering, PIC)
전력 반복 클러스터링(PIC)은 Lin and Cohen이 개발한 확장 가능한 그래프 클러스터링 알고리즘이에요. 요약에 따르면, PIC는 데이터의 정규화된 쌍별 유사도 행렬에 대해 절단된 전력 반복(truncated power iteration)을 사용해 데이터의 매우 저차원 임베딩을 찾아냅니다.
spark.ml의 PowerIterationClustering 구현은 다음 파라미터를 받아요:
k: 만들 클러스터 수initMode: 초기화 알고리즘용 파라미터maxIter: 최대 반복 횟수 파라미터srcCol: 소스 정점 ID를 위한 입력 컬럼 이름 파라미터dstCol: 목적지 정점 ID를 위한 입력 컬럼 이름weightCol: 가중치 컬럼 이름 파라미터
예제 (Examples)
Python:
자세한 내용은 [Python API 문서](api/python/reference/api/pyspark.ml.clustering.PowerIterationClustering.html)를 참고하세요.
from pyspark.ml.clustering import PowerIterationClustering
df = spark.createDataFrame([
(0, 1, 1.0),
(0, 2, 1.0),
(1, 2, 1.0),
(3, 4, 1.0),
(4, 0, 0.1)
], ["src", "dst", "weight"])
pic = PowerIterationClustering(k=2, maxIter=20, initMode="degree", weightCol="weight")
# Shows the cluster assignment
pic.assignClusters(df).show()
전체 예제 코드는 Spark 저장소의 "examples/src/main/python/ml/power_iteration_clustering_example.py"에서 확인할 수 있어요.
Scala:
자세한 내용은 [Scala API 문서](api/scala/org/apache/spark/ml/clustering/PowerIterationClustering.html)를 참고하세요.
import org.apache.spark.ml.clustering.PowerIterationClustering
val dataset = spark.createDataFrame(Seq(
(0L, 1L, 1.0),
(0L, 2L, 1.0),
(1L, 2L, 1.0),
(3L, 4L, 1.0),
(4L, 0L, 0.1)
)).toDF("src", "dst", "weight")
val model = new PowerIterationClustering().
setK(2).
setMaxIter(20).
setInitMode("degree").
setWeightCol("weight")
val prediction = model.assignClusters(dataset).select("id", "cluster")
// Shows the cluster assignment
prediction.show(false)
전체 예제 코드는 Spark 저장소의 "examples/src/main/scala/org/apache/spark/examples/ml/PowerIterationClusteringExample.scala"에서 확인할 수 있어요.
Java:
자세한 내용은 [Java API 문서](api/java/org/apache/spark/ml/clustering/PowerIterationClustering.html)를 참고하세요.
import java.util.Arrays;
import java.util.List;
import org.apache.spark.ml.clustering.PowerIterationClustering;
import org.apache.spark.sql.Dataset;
import org.apache.spark.sql.Row;
import org.apache.spark.sql.RowFactory;
import org.apache.spark.sql.SparkSession;
import org.apache.spark.sql.types.DataTypes;
import org.apache.spark.sql.types.Metadata;
import org.apache.spark.sql.types.StructField;
import org.apache.spark.sql.types.StructType;
List<Row> data = Arrays.asList(
RowFactory.create(0L, 1L, 1.0),
RowFactory.create(0L, 2L, 1.0),
RowFactory.create(1L, 2L, 1.0),
RowFactory.create(3L, 4L, 1.0),
RowFactory.create(4L, 0L, 0.1)
);
StructType schema = new StructType(new StructField[]{
new StructField("src", DataTypes.LongType, false, Metadata.empty()),
new StructField("dst", DataTypes.LongType, false, Metadata.empty()),
new StructField("weight", DataTypes.DoubleType, false, Metadata.empty())
});
Dataset<Row> df = spark.createDataFrame(data, schema);
PowerIterationClustering model = new PowerIterationClustering()
.setK(2)
.setMaxIter(10)
.setInitMode("degree")
.setWeightCol("weight");
Dataset<Row> result = model.assignClusters(df);
result.show(false);
전체 예제 코드는 Spark 저장소의 "examples/src/main/java/org/apache/spark/examples/ml/JavaPowerIterationClusteringExample.java"에서 확인할 수 있어요.
R:
자세한 내용은 [R API 문서](api/R/reference/spark.powerIterationClustering.html)를 참고하세요.
df <- createDataFrame(list(list(0L, 1L, 1.0), list(0L, 2L, 1.0),
list(1L, 2L, 1.0), list(3L, 4L, 1.0),
list(4L, 0L, 0.1)),
schema = c("src", "dst", "weight"))
# assign clusters
clusters <- spark.assignClusters(df, k = 2L, maxIter = 20L,
initMode = "degree", weightCol = "weight")
showDF(arrange(clusters, clusters$id))
전체 예제 코드는 Spark 저장소의 "examples/src/main/r/ml/powerIterationClustering.R"에서 확인할 수 있어요.
더 알아보기 (Learn more)
- MLlib 주요 가이드 (MLlib: Main Guide): MLlib의 전체 구조와 개념을 살펴보세요.
- RDD 기반 클러스터링 가이드: RDD 기반 API에서의 클러스터링 알고리즘을 알아봐요.
- 클러스터링 평가 (ClusteringEvaluator): Silhouette 점수 등 클러스터링 품질을 평가하는 방법을 확인해요.