【问题标题】:How label properly original observations with predicted clusters using kmeans in Pyspark?如何使用 Pyspark 中的 kmeans 用预测的集群正确标记原始观测值?
【发布时间】:2018-04-23 15:03:48
【问题描述】:

我想了解 k-means 方法在 PySpark 中的工作原理。 为此,我做了这个小例子:

In [120]: entry = [ [1,1,1],[2,2,2],[3,3,3],[4,4,4],[5,5,5],[5,5,5],[5,5,5],[1,1,1],[5,5,5]]

In [121]: rdd_entry = sc.parallelize(entry)

In [122]: clusters = KMeans.train(rdd_entry, k=5, maxIterations=10, initializationMode="random")

In [123]:  rdd_labels = clusters.predict(rdd_entry)

In [125]: rdd_labels.collect()
Out[125]: [3, 1, 0, 0, 2, 2, 2, 3, 2]

In [126]: entry
Out[126]:
[[1, 1, 1],
 [2, 2, 2],
 [3, 3, 3],
 [4, 4, 4],
 [5, 5, 5],
 [5, 5, 5],
 [5, 5, 5],
 [1, 1, 1],
 [5, 5, 5]]

乍一看,rdd_labels 似乎返回每个观察所属的集群,尊重原始 rdd 的顺序。尽管在此示例中很明显,我如何确定在我将使用 800 万个观测值的情况下工作?

另外,我想知道如何加入 rdd_entry 和 rdd_labels,尊重该顺序,以便 rdd_entry 的每个观察都正确地用其集群标记。 我试图做一个.join(),但它跳转错误

In [127]: rdd_total = rdd_entry.join(rdd_labels)

In [128]: rdd_total.collect()

TypeError: 'int' object has no attribute '__getitem__'

【问题讨论】:

  • 您是否仅限于使用pyspark.mllib(即将被弃用),或者您可能想要基于pyspark.ml 的解决方案(即首选的基于数据帧的API)?

标签: pyspark cluster-analysis apache-spark-mllib


【解决方案1】:

希望对您有所帮助! (本方案基于pyspark.ml

from pyspark.ml.clustering import KMeans
from pyspark.ml.feature import VectorAssembler

#sample data
df = sc.parallelize([[1,1,1],[2,2,2],[3,3,3],[4,4,4],[5,5,5],[5,5,5],[5,5,5],[1,1,1],[5,5,5]]).\
    toDF(('col1','col2','col3'))

vecAssembler = VectorAssembler(inputCols=df.columns, outputCol="features")
vector_df = vecAssembler.transform(df)

#kmeans clustering
kmeans=KMeans(k=3, seed=1)
model=kmeans.fit(vector_df)
predictions=model.transform(vector_df)
predictions.show()

输出是:

+----+----+----+-------------+----------+
|col1|col2|col3|     features|prediction|
+----+----+----+-------------+----------+
|   1|   1|   1|[1.0,1.0,1.0]|         0|
|   2|   2|   2|[2.0,2.0,2.0]|         0|
|   3|   3|   3|[3.0,3.0,3.0]|         2|
|   4|   4|   4|[4.0,4.0,4.0]|         1|
|   5|   5|   5|[5.0,5.0,5.0]|         1|
|   5|   5|   5|[5.0,5.0,5.0]|         1|
|   5|   5|   5|[5.0,5.0,5.0]|         1|
|   1|   1|   1|[1.0,1.0,1.0]|         0|
|   5|   5|   5|[5.0,5.0,5.0]|         1|
+----+----+----+-------------+----------+

虽然pyspark.ml 有更好的方法,但我想编写代码以使用pyspark.mllib 实现相同的结果(触发是来自@Muhammad 的评论)。所以这里是基于pyspark.mllib的解决方案...

from pyspark.mllib.clustering import KMeans
from pyspark.sql.functions import monotonically_increasing_id, row_number
from pyspark.sql.window import Window
from pyspark.sql.types import IntegerType

#sample data
rdd = sc.parallelize([[1,1,1],[2,2,2],[3,3,3],[4,4,4],[5,5,5],[5,5,5],[5,5,5],[1,1,1],[5,5,5]])

#K-Means example
model = KMeans.train(rdd, k=3, seed=1)
labels = model.predict(rdd)

#add cluster label to the original data
df1 = rdd.toDF(('col1','col2','col3')) \
         .withColumn('row_index', row_number().over(Window.orderBy(monotonically_increasing_id())))
df2 = spark.createDataFrame(labels, IntegerType()).toDF(('label')) \
           .withColumn('row_index', row_number().over(Window.orderBy(monotonically_increasing_id())))
df = df1.join(df2, on=["row_index"]).drop("row_index")
df.show()

【讨论】:

  • @CarmenPérezCarrillo 如果它可以帮助您解决问题,也许您应该accept the answer
  • 对不起,我无法证明这一点。我只是做了它并且它有效。谢谢你:)
  • @Prem pyspark.mllib.clustering有什么解决办法
  • @Muhammad 我添加了代码以使用pyspark.mllib 获得相同的结果。希望对您有所帮助!
猜你喜欢
  • 2017-05-25
  • 2016-07-08
  • 2020-09-08
  • 2013-12-23
  • 2019-01-04
  • 2016-07-20
  • 1970-01-01
  • 2013-12-13
  • 2020-11-05
相关资源
最近更新 更多