【问题标题】:Method for all rows of a PySpark DataFramePySpark DataFrame 的所有行的方法
【发布时间】:2018-10-06 21:48:22
【问题描述】:

我在为 PySpark 上的任务(python=2.7,pyspark=1.6)设计一个有效的 udf 时遇到了麻烦

我有一个data DataFrame,看起来像这样:

+-----------------+
|         sequence|
+-----------------+
|         idea can|  
|        fuel turn|          
|      found today|  
|           a cell|         
|administration in|           
+-----------------+

对于data 中的每一行,我想在另一个名为ggrams 的DataFrame 中查找信息(基于属性sequence),计算聚合并将其作为data 中的新列返回。

我的感觉是我应该这样做:

from pyspark.sql.types import IntegerType
from pyspark.sql.functions import udf

def compute_aggregates(x):
    res = ggrams.filter((ggrams.ngram.startswith(x)) \
                              ).groupby("ngram").sum("match_count")
    return res.collect()[0]['sum(match_count)']

aggregate = udf(compute_aggregates, IntegerType())
result = data.withColumn('aggregate', compute_aggregates('sequence'))

但这会返回 PicklingError。

PicklingError: Could not serialize object: Py4JError: An error occurred while calling o1759.__getstate__. Trace:
py4j.Py4JException: Method __getstate__([]) does not exist

【问题讨论】:

  • 你的 spark 版本是什么?
  • 对不起,我忘了提。 Python 2.7、pyspark 1.6 -- 添加到帖子中
  • 您不能在 UDF 中使用 collect。你不能用join 做到这一点吗?
  • 即使删除 collect 也会出现同样的错误。不幸的是,我设置了相等条件 ggrams.ngram == x 以使其更容易,但它也可能是 ggrams.ngram.startswith(x)ggrams.ngram.contains(x) 所以它不是完全符合加入条件的,除非我错过了。让我修改我的问题以消除歧义
  • 你的数据框有多大?如果不是那么大,请考虑进行交叉连接,然后应用 udf 过滤行。

标签: python pyspark user-defined-functions


【解决方案1】:

抛出错误是因为您无法访问另一个 udf 中的不同数据帧。解决此问题的最简单方法是收集您要检查的数据框。

另一个选项是交叉连接,但我可以根据经验说收集其他数据帧更快。 (虽然我没有数学/统计数据来支持这一点)

所以:
1. 收集你想在udf中使用的dataframe
2. 在你的 udf 中调用这个收集的数据框(现在是一个列表),你现在可以/必须使用 python 逻辑,因为你正在与一个对象列表交谈

注意:尽量将您收集的数据框限制在最低限度,只选择您需要的列

更新:
如果您正在处理一个非常大的集合,这将使收集变得不可能,那么交叉连接现在很可能会起作用(至少对我来说不是)。问题是两个数据帧的巨大叉积将花费太多时间来创建与工作节点的连接将超时,这会导致广播错误。因此,请查看是否有任何方法可以限制您正在使用的列,或者是否有可能过滤掉您可以确定它们不会被使用的行。

如果这一切都失败了,看看你是否可以创建一些批处理方法*,所以只运行前 X 行收集的数据,如果这样做了,加载接下来的 X 行。这很可能会非常慢,但至少不会超时(我想,我没有亲自尝试过,因为我可以收集)

*批处理你正在运行 udf 的数据帧和另一个数据帧,因为你仍然无法在 udf 内收集,因为你无法从那里访问数据帧

【讨论】:

  • 显然,df 是 gbs ......你需要一台大电脑,然后如果你想收集
  • 嗯,这很糟糕,但是没有办法访问另一个 udf 中的一个数据帧。然后交叉连接很可能会引发广播异常,因为工人连接将在如此大的集合上超时(这发生在我身上)
  • 我还是可以试一试,这是在一个相当大的集群上完成的
猜你喜欢
  • 2016-12-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-04-29
  • 2022-01-16
相关资源
最近更新 更多