【问题标题】:Summing up DenseVectors after a groupByKey() works in Pyspark Shell but not in a Spark-Submit在 groupByKey() 在 Pyspark Shell 中有效但在 Spark-Submit 中无效后总结 DenseVectors
【发布时间】:2016-06-28 19:22:27
【问题描述】:

下面是我正在尝试做的一些示例代码:

首先,我正在使用 Word2Vec 构建句子特征向量:

from pyspark.ml.feature import Word2Vec

# Input data: Each row is a bag of words from a sentence or document.
documentDF = sqlContext.createDataFrame([
    ("Hi I heard about Spark".split(" "), ),
    ("I wish Java could use case classes".split(" "), ),
    ("Logistic regression models are neat".split(" "), )
], ["text"])
# Learn a mapping from words to Vectors.
word2Vec = Word2Vec(vectorSize=3, minCount=0, inputCol="text", outputCol="result")
model = word2Vec.fit(documentDF)
result = model.transform(documentDF)

Converting output result to an RDD:
result_rdd=result.select("result").rdd
rdd_with_sample_ids_attached = result_rdd.map(lambda x: (1, x[0]))
rdd_with_sample_ids_attached.collect()

输出: [(1, DenseVector([0.0472, -0.0078, 0.0377])), (1, DenseVector([-0.0253, -0.0171, 0.0664])), (1, DenseVector([0.0101, 0.0324, 0.0158]))]

现在,我执行 groupByKey() 并找到每个组中 DenseVectors 的总和如下:

rdd_sum = rdd_with_sample_ids_attached.groupByKey().map(lambda x: (x[0], sum(x[1])))
rdd_sum.collect()

输出: [(1, DenseVector([0.0319, 0.0075, 0.1198]))]

如图所示,此代码在 pyspark shell 中完美运行。但是,当我提交与 spark-submit 相同的代码时,我收到以下错误:

File "/mnt1/yarn/usercache/hadoop/appcache/application_1465567204576_0170/container_1465567204576_0170_01_000002/pyspark.zip/pyspark/sql/functions.py", line 39, in _
   jc = getattr(sc._jvm.functions, name)(col._jc if isinstance(col, Column) else col)
AttributeError: 'NoneType' object has no attribute '_jvm'

我尝试将 RDD 重新分区到单个分区,同样的错误。 请帮忙?

【问题讨论】:

  • 错误提示scNoneType。也许您正在处理一个大型数据集并且您的集群已经死了?这是很有可能的,尤其是。您正在使用groupByKey,它需要足够大的内存来保存任何键及其值。
  • groupByKey 之后的 lambda 函数也不给出总和
  • 嘿!不,我尝试使用上面发布的相同示例数据集。同样的错误。在 pyspark shell 中工作,当我在 .py 文件中提交相同的代码时不起作用。此外,lambda 函数确实进行了求和 - 这是问题中的一个错字。我已经在问题中编辑了我的代码。我能够展开与分组 ID 关联的 DenseVectors 列表,并且还可以执行 len() 操作。它只是失败的 sum() 。令人沮丧的是,它可以在 pyspark shell 和 ipython notebook 中运行,所以我觉得我在这里遗漏了一些东西。

标签: python apache-spark pyspark


【解决方案1】:

想通了! 问题是我的脚本中有一个导入函数,如下所示:

from pyspark.sql.functions import *

这导入了 sum() 函数,它取代了内置的 pythonic sum()。当我删除此导入功能时,它可以正常工作。当 pythonic 内置 sum() 函数能够添加 DenseVectors 时,从 pyspark.sql.functions 导入的 sum() 不能这样做。

【讨论】:

  • 显然这不是你的问题:)
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-01-22
  • 1970-01-01
  • 1970-01-01
  • 2019-07-30
  • 2012-04-02
相关资源
最近更新 更多