如何将类型<class'pyspark.sql.types.Row'>转换为Vector

hpn*_*xwn 5 python machine-learning k-means apache-spark pyspark

我是Spark的新手,目前我正在尝试使用Python编写对一组数据执行KMeans的简单代码。

from pyspark import SparkContext, SparkConf
from pyspark.sql import SQLContext
import re
from pyspark.mllib.clustering import KMeans, KMeansModel
from pyspark.mllib.linalg import DenseVector
from pyspark.mllib.linalg import SparseVector
from numpy import array
from pyspark.ml.feature import VectorAssembler
from pyspark.ml.feature import MinMaxScaler

import pandas as pd
import numpy
df = pd.read_csv("/<path>/Wholesale_customers_data.csv")
sql_sc = SQLContext(sc)
cols = ["Channel", "Region", "Fresh", "Milk", "Grocery", "Frozen", "Detergents_Paper", "Delicassen"]
s_df = sql_sc.createDataFrame(df)
vectorAss = VectorAssembler(inputCols=cols, outputCol="feature")
vdf = vectorAss.transform(s_df)
km = KMeans.train(vdf, k=2, maxIterations=10, runs=10, initializationMode="k-means||")
model = kmeans.fit(vdf)
cluster = model.clusterCenters()
print(cluster)
Run Code Online (Sandbox Code Playgroud)

我在pyspark shell中键入了它们,当它运行model = kmeans.fit(vdf)时,出现以下错误:

TypeError:无法将类型转换为Vector

在org.apache.spark.api.python.PythonRunner $$ anon $ 1.read(PythonRDD.scala:166)在org.apache.spark.api.python.PythonRunner $$ anon $ 1。(PythonRDD.scala:207)在位于org.apache.spark.rdd的org.apache.spark.api.python.PythonRDD.compute(PythonRDD.scala:70)的org.apache.spark.api.python.PythonRunner.compute(PythonRDD.scala:125)。 org.org.apache.spark.rdd上的RDD.computeOrReadCheckpoint(RDD.scala:313).org org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38)上的RD.iterator(RDD.scala:277) org.apache.spark.CacheManager.getOrCompute(CacheManager.scala:69)上的.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:313)org.apache.spark.rdd.RDD.iterator(RDD.scala) :275),位于org.apache.spark.rdd.ZippedPartitionsRDD2.compute(ZippedPartitionsRDD.scala:88),位于org.apache.spark.rdd.RDD.computeOrReadCheckpoint(RDD.scala:313),位于org.apache.spark.rdd.MapPartitionsRDD.compute(MapPartitionsRDD.scala:38),位于org.apache.spark.rdd.RDD.iterator(RDD.scala:277),位于org.apache.spark.rdd.RDD org.apache.spark.rdd上的.computeOrReadCheckpoint(RDD.scala:313).org上org.apache.spark.scheduler.ResultTask.runTask(ResultTask.scala:66)上的RDD.iterator(RDD.scala:277)。 org.apache.spark.executor.Executor $ TaskRunner.run(Executor.scala:227)的apache.spark.scheduler.Task.run(Task.scala:89),java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor。 java:1142)at java.util.concurrent.ThreadPoolExecutor $ Worker.run(ThreadPoolExecutor.java:617)at java.lang.Thread.run(Thread.java:745)17/02/26 23:31:58错误执行器:阶段23.0(TID 113)org.apache.spark.api.python.Python.6.0中任务6.0中的异常:追溯(最近一次调用为最新):“文件”/usr/hdp/2.5.0.0-1245/spark/python/lib/pyspark.zip/pyspark/worker.py”,第111行,位于主进程中()文件“ /usr/hdp/2.5.0.0-1245/spark” /python/lib/pyspark.zip/pyspark/worker.py”,第106行,在进程serializer.dump_stream(func(split_index,iterator),outfile)文件“ /usr/hdp/2.5.0.0-1245/spark/python”中/lib/pyspark.zip/pyspark/serializers.py“,第263行,位于dump_stream vs = list(itertools.islice(迭代器,批处理))文件” /usr/hdp/2.5.0.0-1245/spark/python/lib /pyspark.zip/pyspark/mllib/linalg/init.py“,第77行,在_convert_to_vector中,引发TypeError(”无法将类型%s转换为Vector“%type(l))TypeError:无法将类型转换为Vector0-1245 / spark / python / lib / pyspark.zip / pyspark / worker.py”,行106,在进程serializer.dump_stream(func(split_index,iterator),outfile)中文件“ /usr/hdp/2.5.0.0- 1245 / spark / python / lib / pyspark.zip / pyspark / serializers.py”,第263行,位于dump_stream vs = list(itertools.islice(iterator,batch))文件“ /usr/hdp/2.5.0.0-1245/ spark / python / lib / pyspark.zip / pyspark / mllib / linalg / init.py“,第77行,在_convert_to_vector中,引发TypeError(”无法将类型%s转换为Vector“%type(l))TypeError:无法将类型转换为矢量0-1245 / spark / python / lib / pyspark.zip / pyspark / worker.py”,行106,在进程serializer.dump_stream(func(split_index,iterator),outfile)中文件“ /usr/hdp/2.5.0.0- 1245 / spark / python / lib / pyspark.zip / pyspark / serializers.py”,第263行,位于dump_stream vs = list(itertools.islice(iterator,batch))文件“ /usr/hdp/2.5.0.0-1245/ spark / python / lib / pyspark.zip / pyspark / mllib / linalg / init.py“,第77行,在_convert_to_vector中,引发TypeError(”无法将类型%s转换为Vector“%type(l))TypeError:无法将类型转换为矢量批处理))文件“ /usr/hdp/2.5.0.0-1245/spark/python/lib/pyspark.zip/pyspark/mllib/linalg/init.py”,第77行,在_convert_to_vector中引发TypeError(“无法转换类型% s转换为Vector“%type(l))TypeError:无法将类型转换为Vector批处理))文件“ /usr/hdp/2.5.0.0-1245/spark/python/lib/pyspark.zip/pyspark/mllib/linalg/init.py”,第77行,在_convert_to_vector中引发TypeError(“无法转换类型% s转换为Vector“%type(l))TypeError:无法将类型转换为Vector

我得到的数据来自:https : //archive.ics.uci.edu/ml/machine-learning-databases/00292/Wholesale%20customers%20data.csv

有人可以告诉我这里出了什么问题以及我错过了什么吗?感谢您的帮助。

谢谢!

更新:@Garren我得到的错误是:

我得到的错误是:>>> kmm = kmeans.fit(s_df)17/03/02 21:58:01 INFO BlockManagerInfo:在内存中的本地主机上删除了broadcast_1_piece0:56193(大小:5.8 KB,可用空间:511.1 MB)17 / 03/02 21:58:01 INFO ContextCleaner:已清除累加器5 17/03/02 21:58:01 INFO BlockManagerInfo:删除了本地主机上的广播_0_piece0:56193在内存中(大小:5.8 KB,可用空间:511.1 MB)17/03 / 02 21:58:01 INFO ContextCleaner:清理了累加器4

追溯(最近一次通话最近):在文件“ /usr/hdp/2.5.0.0-1245/spark/python/pyspark/ml/pipeline.py”中的文件“”,行1,行69,适合返回自身。 _fit(数据集)文件“ /usr/hdp/2.5.0.0-1245/spark/python/pyspark/ml/wrapper.py”,第133行,在_fit java_model = self._fit_java(dataset)文件“ / usr / hdp / 2.5.0.0-1245 / spark / python / pyspark / ml / wrapper.py“,行130,在_fit_java中返回self._java_obj.fit(dataset._jdf)文件“ /usr/hdp/2.5.0.0-1245/spark/ python / lib / py4j-0.9-src.zip / py4j / java_gateway.py”,行813,正在调用 在装饰性引发AnalysisException(s.split(':',1)[1],stackTrace)中的文件“ /usr/hdp/2.5.0.0-1245/spark/python/pyspark/sql/utils.py”,第51行pyspark.sql.utils.AnalysisException:u“无法解析给定输入列的“功能”:[通道,杂货,新鲜,冷冻,洗涤剂_纸,区域,熟食,牛奶];”

Gar*_*n S 3

仅使用 Spark 2.x ML 包而不是[即将弃用]spark mllib 包:

from pyspark.ml.clustering import KMeans
from pyspark.ml.feature import VectorAssembler
df = spark.read.option("inferSchema", "true").option("header", "true").csv("whole_customers_data.csv")
cols = df.columns
vectorAss = VectorAssembler(inputCols=cols, outputCol="features")
vdf = vectorAss.transform(df)
kmeans = KMeans(k=2, maxIter=10, seed=1)
kmm = kmeans.fit(vdf)
kmm.clusterCenters()
Run Code Online (Sandbox Code Playgroud)