【问题标题】:PySpark - undefined function collect_listPySpark - 未定义的函数 collect_list
【发布时间】:2020-07-01 18:22:14
【问题描述】:

我正在使用 Python 2.6.6 和 Spark 1.6.0。我有df 这样的:

id | name      |  number |
-------------------------- 
1  | joe       | 148590  |
2  | bob       | 148590  |
2  | steve     | 279109  |
3  | sue       | 382901  |
3  | linda     | 148590  |

每当我尝试运行类似 df2 = df.groupBy('id','length','type').pivot('id').agg(collect_list('name')),我收到以下错误 pyspark.sql.utils.AnalysisException: u'undefined function collect_list;'这是为什么呢?

我也试过: hive_context = HiveContext(sc) df2 = df.groupBy('id','length','type').pivot('id').agg(hive_context.collect_list('name')) 并得到错误:

AttributeError: 'HiveContext' object has no attribute 'collect_list'

【问题讨论】:

    标签: python dataframe apache-spark pyspark


    【解决方案1】:

    这里的collect_list 看起来像一个用户定义的函数。 PySpark API 仅支持少数预定义函数,如 sum、count 等

    如果您指的是任何其他代码,请确保您在某处定义了 collect_list 函数。 要导入集体主义功能,请在顶部添加以下行

    from pyspark.sql import functions as F
    

    然后将代码更改为:

     df.groupBy('id','length','type').pivot('id').agg(F.collect_list(name))
    

    如果你已经定义了,试试下面的 sn-p。

    df.groupBy('id','length','type').pivot('id').agg({'name':'collect_list'})
    

    【讨论】:

    • 你如何定义它?我不认为我是。
    • 检查我更新的答案。这是有用的链接stackoverflow.com/questions/41026178/…
    • 我收到错误:Aggregate expression required for pivot, found 'pythonUDF#33';
    • 这似乎是一个不同的错误。这也意味着您的 collecti_list 相关错误已解决。对于新的错误,我相信你需要改变agg和pivot的顺序。 df.groupBy('id','length','type').agg({'name':'collect_list'}).pivot('id')
    • 好的.. 看起来应该是第一个枢轴。所以这不是问题。也许您可以发布一个与此枢轴相关错误的新问题。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-10-01
    • 2019-11-29
    • 1970-01-01
    • 2015-09-03
    相关资源
    最近更新 更多