【问题标题】:Issue with using .take() function with spark 2+ pyspark使用带有 spark 2+ pyspark 的 .take() 函数的问题
【发布时间】:2020-04-19 20:32:18
【问题描述】:

这是我正在使用的代码。在没有 data.take 的情况下它运行良好,但在 pyspark python 中使用它时会出错

from pyspark.mllib.recommendation import ALS, MatrixFactorizationModel, Rating
data = sc.textFile("re_u.data")
pData=data.take(2000)
ratings = pData.map(lambda l: l.split(','))\
.map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))

给出错误

AttributeError                            Traceback (most recent call last)
<ipython-input-12-c9c51af1b2e9> in <module>
      2 data = sc.textFile("re_u.data")
      3 pData=data.take(2000)
----> 4 ratings = pData.map(lambda l: l.split(','))\
      5  .map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))

AttributeError: 'list' object has no attribute 'map'

更新: 使用您的更改@Hristo Iliev 后,它有所帮助,但遇到了另一个问题,随后是评级作为列表。感谢您的帮助!

from pyspark.mllib.recommendation import ALS, MatrixFactorizationModel, Rating
data = sc.textFile("re_u.data")
ratings = data.map(lambda l: l.split(','))\
  .map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))\
  .take(2000)
rank = 20
numIterations = 20
model = ALS.train(ratings, rank, numIterations)

报错

---------------------------------------------------------------------------
TypeError                                 Traceback (most recent call last)
<ipython-input-24-7e35afff970b> in <module>
      1 rank = 20
      2 numIterations = 20
----> 3 model = ALS.train(ratings, rank, numIterations)

C:\spark\spark-3.0.0-preview2-bin-hadoop2.7\python\pyspark\mllib\recommendation.py in train(cls, ratings, rank, iterations, lambda_, blocks, nonnegative, seed)
    271           (default: None)
    272         """
--> 273         model = callMLlibFunc("trainALSModel", cls._prepare(ratings), rank, iterations,
    274                               lambda_, blocks, nonnegative, seed)
    275         return MatrixFactorizationModel(model)

C:\spark\spark-3.0.0-preview2-bin-hadoop2.7\python\pyspark\mllib\recommendation.py in _prepare(cls, ratings)
    227         else:
    228             raise TypeError("Ratings should be represented by either an RDD or a DataFrame, "
--> 229                             "but got %s." % type(ratings))
    230         first = ratings.first()
    231         if isinstance(first, Rating):

TypeError: Ratings should be represented by either an RDD or a DataFrame, but got <class 'list'>.

请帮忙!

【问题讨论】:

  • 欢迎来到 Stack Overflow。您最初没有展示您是如何使用ratings 的。请始终提供完整的上下文,因为这会改变答案。
  • 请不要在您的原始问题得到回答后通过添加后续问题来更新您的帖子。接受答案,并在新标题下提出新问题。

标签: python-3.x apache-spark pyspark jupyter-notebook rdd


【解决方案1】:

take() 是一个动作,它从 RDD 的顶部获取指定数量的元素并将它们传输到驱动程序。你从中得到的是一个包含请求元素的 Python 列表,即:

  • 驱动程序本地,因此您不应使用太多元素
  • 没有map() 方法只是因为Python list 类没有map() 方法

您最可能想要做的是首先将转换应用到转换后的 RDD 中的 data RDD 和 take()

data = sc.textFile("re_u.data")
ratings = data.map(lambda l: l.split(','))\
  .map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))\
  .take(2000)

您将获得Rating 实例的列表。

由于您将数据进一步向下传递给 ALS,后者采用分布式数据,即 RDD,而不是驱动程序本地 list,因此您有三个选择:

  1. 再次并行化列表,将其变成 RDD:

    data = sc.textFile("re_u.data")
    ratings = data.map(lambda l: l.split(','))\
      .map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))\
      .take(2000)
    ratingsRDD = sc.parallelize(ratings)
    rank = 20
    numIterations = 20
    model = ALS.train(ratingsRDD, rank, numIterations)
    
  2. 使用sample()方法对RDD中的数据子集进行采样:

    data = sc.textFile("re_u.data")
    ratings = data.map(lambda l: l.split(','))\
      .map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))\
      .sample(False, 0.1, 42)
    rank = 20
    numIterations = 20
    model = ALS.train(ratings, rank, numIterations)
    

    这里的sample(False, 0.1, 42) 表示取大约10% 的原始数据并使用42 作为伪随机数生成器的种子。固定种子将允许测试时的可重复性。您应该将0.1 调整为适当的值,以便获得大约 2000 个样本。请注意,这些样本将取自 RDD 内的随机位置,并且很可能不是前 2000 个。

  3. 模拟take(),同时保持 RDD 领域:

    data = sc.textFile("re_u.data")
    ratings = data.map(lambda l: l.split(','))\
      .map(lambda l: Rating(int(l[0]), int(l[1]), float(l[2])))\
      .zipWithIndex()\
      .filter(lambda l: l[1] < 2000)\
      .map(lambda l: l[0])
    rank = 20
    numIterations = 20
    model = ALS.train(ratings, rank, numIterations)
    

    zipWithIndex() 创建 RDD 内容的元组,其中第一个元素来自 RDD,第二个元素是 RDD 中的索引(本质上是行号)。然后,您可以只过滤索引小于 2000 的元素,然后使用 map(lambda l: l[0]) 删除索引。

方法2可能是最好的。

【讨论】:

  • 感谢@Hristo Iliev 的帮助,现在我得到的错误评级是使用 ALS 功能的列表。 rank = 20 numIterations = 20 model = ALS.train(ratings, rank, numIterations) Ratings 应该由 RDD 或 DataFrame 表示,但得到 .
  • 谢谢!你会建议使用 list 然后转换为 rdd likeratings_rdd= spark.sparkContext.parallelize(rat)
  • @KaranModi 我已根据您修改后的问题更新了答案。
猜你喜欢
  • 2019-06-09
  • 1970-01-01
  • 1970-01-01
  • 2018-04-13
  • 1970-01-01
  • 2018-01-29
  • 2020-04-06
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多