【发布时间】: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