【发布时间】:2017-10-01 00:38:45
【问题描述】:
使用 zeppelin 笔记本中的 spark,我从昨天开始收到此错误。 这是我的代码:
from pyspark.ml.clustering import KMeans
from pyspark.ml.feature import VectorAssembler
df = sqlContext.table("rfmdata_clust")
k = 4
# Set Kmeans input/output columns
vecAssembler = VectorAssembler(inputCols=["v1_clust", "v2_clust", "v3_clust"], outputCol="features")
featuresDf = vecAssembler.transform(df)
# Run KMeans
kmeans = KMeans().setInitMode("k-means||").setK(k)
model = kmeans.fit(featuresDf)
resultDf = model.transform(featuresDf)
# KMeans WSSSE
wssse = model.computeCost(featuresDf)
print("Within Set Sum of Squared Errors = " + str(wssse))
这是错误:
Traceback (most recent call last):
File "/tmp/zeppelin_pyspark-8890997346928959256.py", line 346, in <module>
raise Exception(traceback.format_exc())
Exception: Traceback (most recent call last):
File "/tmp/zeppelin_pyspark-8890997346928959256.py", line 334, in <module>
exec(code)
File "<stdin>", line 8, in <module>
File "/usr/lib/spark/python/pyspark/ml/base.py", line 64, in fit
return self._fit(dataset)
File "/usr/lib/spark/python/pyspark/ml/wrapper.py", line 236, in _fit
java_model = self._fit_java(dataset)
File "/usr/lib/spark/python/pyspark/ml/wrapper.py", line 233, in _fit_java
return self._java_obj.fit(dataset._jdf)
File "/usr/lib/spark/python/lib/py4j-0.10.4-src.zip/py4j/java_gateway.py", line 1133, in __call__
answer, self.gateway_client, self.target_id, self.name)
File "/usr/lib/spark/python/pyspark/sql/utils.py", line 79, in deco
raise IllegalArgumentException(s.split(': ', 1)[1], stackTrace)
IllegalArgumentException: u'requirement failed'
引发错误的行是 kmeans.fit() 行。 我检查了 rfmdata_clust 数据框,它似乎一点也不奇怪。
df.printSchema() 给出:
root
|-- id: string (nullable = true)
|-- v1_clust: double (nullable = true)
|-- v2_clust: double (nullable = true)
|-- v3_clust: double (nullable = true)
featuresDf.printSchema() 给出:
root
|-- id: string (nullable = true)
|-- v1_clust: double (nullable = true)
|-- v2_clust: double (nullable = true)
|-- v3_clust: double (nullable = true)
|-- features: vector (nullable = true)
另一个有趣的点是在 featuresDf 的定义下面添加featuresDf = featuresDf.limit(10000) 可以使代码无错误地运行。可能和数据的大小有关?
【问题讨论】:
标签: apache-spark dataframe pyspark k-means apache-spark-mllib