【发布时间】:2019-03-05 18:15:08
【问题描述】:
我是 pyspark 的新手,我有一个脚本如下;
joinedRatings=ratings.join(ratings)
joinedRatings.take(4)
输出是;
[(196, ((242, 3.0), (242, 3.0))), (196, ((242, 3.0), (393, 4.0))), (196, ((242, 3.0), (381, 4.0))), (196, ((242, 3.0), (251, 3.0)))]
之后我的功能就是 ;
def filterDuplicates(userRatings):
ratings = userRatings[1]
(movie1, rating1) = ratings[0]
(movie2, rating2) = ratings[1]
return movie1 < movie2
比我有这个 RDD
uniqueJoinedRatings = joinedRatings.filter(filterDuplicates)
我的问题是能够理解我编写的这个函数是如何运行的
joinedRatings[1]
我收到的错误是;
Fail to execute line 1: joinedRatings[1]
Traceback (most recent call last):
File "/tmp/zeppelin_pyspark-240579357005199320.py", line 380, in <module>
exec(code, _zcUserQueryNameSpace)
File "<stdin>", line 1, in <module>
TypeError: 'PipelinedRDD' object does not support indexing
但它在“def filterDuplicates(userRatings):”函数下运行没有任何问题,请告诉我如何学习“joinedRatings[1]”的值?
【问题讨论】:
-
print type(userRatings)infilterDuplicates打印什么?print type(joinedRatings)打印什么? -
@Tobias Brösamle ,类型(joinedRatings)是
-
type(userRatings) 怎么样?
-
@Tobias Brösamle (userRatings) 是 def 函数下的一个变量,我使用joinedRatings 来执行这个函数
-
但显然 userRatings 与joinedRatings 的类型不同。这就是为什么我希望您打印 userRatings 的类型,所以您可以看到它不一样,并且“但它适用于 filterDuplicates 中的 userRatings”的论点是无效的。
标签: python apache-spark pyspark