【问题标题】:spark aggregation for array column数组列的火花聚合
【发布时间】:2019-02-27 19:51:43
【问题描述】:

我有一个带有数组列的数据框。

val json = """[
{"id": 1, "value": [11, 12, 18]},
{"id": 2, "value": [23, 21, 29]}
]"""

val df = spark.read.json(Seq(json).toDS)

scala> df.show
+---+------------+
| id|       value|
+---+------------+
|  1|[11, 12, 18]|
|  2|[23, 21, 29]|
+---+------------+

现在我需要对值列应用不同的聚合函数。 我可以打电话给explodegroupBy,例如

df.select($"id", explode($"value").as("value")).groupBy($"id").agg(max("value"), avg("value")).show

+---+----------+------------------+
| id|max(value)|        avg(value)|
+---+----------+------------------+
|  1|        18|13.666666666666666|
|  2|        29|24.333333333333332|
+---+----------+------------------+

这里让我困扰的是,我将我的 DataFrame 分解为一个更大的,然后将其缩减为原始调用 groupBy

有没有更好(即更有效)的方法来调用数组列上的聚合函数?可能我可以实现 UDF,但我不想自己实现所有聚合 UDF。

编辑。有人引用了this SO question,但在我的情况下它不起作用。 size 工作正常

scala> df.select($"id", size($"value")).show
+---+-----------+
| id|size(value)|
+---+-----------+
|  1|          3|
|  2|          3|
+---+-----------+

但是avgmax 不起作用。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql aggregate-functions


    【解决方案1】:

    简短的回答是否定的,您必须实现自己的 UDF 来聚合数组列。至少在最新版本的 Spark 中(撰写本文时为 2.3.1)。正如您正确断言的那样,这不是很有效,因为它会迫使您分解行或支付在 Dataset API 中工作的序列化和反序列化成本。

    对于可能会发现此问题的其他人,要使用 Datasets 以类型安全的方式编写聚合,您可以使用 Aggregator API,该 API 诚然没有很好的文档记录,并且随着类型签名变得相当混乱,使用起来非常混乱详细。

    更长的答案是,此功能即将在 Apache Spark 2.4 中推出(?)

    父问题SPARK-23899补充:

    • array_max
    • array_min
    • 聚合
    • 地图
    • array_distinct
    • array_remove
    • array_join

    还有很多其他

    这个演讲“Extending Spark SQL API with Easier to Use Array Types Operations”在 2018 年 6 月的 Spark + AI 峰会上发表,涵盖了新功能。

    如果它已发布,将允许您像在您的示例中一样使用 max 函数,但是 average 有点棘手。 奇怪的是,array_sum 不存在,但它可以从 aggregate 函数构建。它可能看起来像:

    def sum_array(array_col: Column) = aggregate($"my_array_col", 0, (s, x) => s + x, s => s) df.select(sum_array($"my_array_col") 其中零值是聚合缓冲区的初始状态。

    正如你所指出的size已经可以获取数组的长度,这意味着可以计算平均值。

    【讨论】:

    • 很好且解释清楚的答案
    猜你喜欢
    • 1970-01-01
    • 2016-06-04
    • 1970-01-01
    • 1970-01-01
    • 2018-01-16
    • 2020-08-06
    • 2021-11-19
    • 2016-04-29
    相关资源
    最近更新 更多