【问题标题】:scala aggregate first function giving unexpected results [duplicate]scala聚合第一个函数给出了意想不到的结果[重复]
【发布时间】:2019-06-30 16:38:22
【问题描述】:

我在 scala spark 中使用了一个简单的 groupby 查询,其目标是在已排序的数据框中获取组中的第一个值。这是我的火花数据框

+---------------+------------------------------------------+
|ID             |some_flag |some_type  |  Timestamp        |
+---------------+------------------------------------------+
|      656565654|      true|     Type 1|2018-08-10 00:00:00|
|      656565654|     false|     Type 1|2017-08-02 00:00:00|
|      656565654|     false|     Type 2|2016-07-30 00:00:00|
|      656565654|     false|     Type 2|2016-05-04 00:00:00|
|      656565654|     false|     Type 2|2016-04-29 00:00:00|
|      656565654|     false|     Type 2|2015-10-29 00:00:00|
|      656565654|     false|     Type 2|2015-04-29 00:00:00|
+---------------+----------+-----------+-------------------+

这是我的汇总查询

val sampleDF = df.sort($"Timestamp".desc).groupBy("ID").agg(first("Timestamp"), first("some_flag"), first("some_type"))

预期的结果是

+---------------+-------------+---------+-------------------+
|ID             |some_falg    |some_type|  Timestamp        |
+---------------+-------------+---------+-------------------+
|      656565654|         true|   Type 1|2018-08-10 00:00:00|
+---------------+-------------+---------+-------------------+

但是得到以下奇怪的输出并且它像随机行一样不断变化

+---------------+-------------+---------+-------------------+
|ID             |some_falg    |some_type|  Timestamp        |
+---------------+-------------+---------+-------------------+
|      656565654|        false|   Type 2|2015-10-29 00:00:00|
+---------------+-------------+---------+-------------------+

另外请注意,数据框中没有空值。我在做错事时挠头。需要帮助!

【问题讨论】:

  • timestamp的数据类型是什么?
  • timestamp 是时间戳的数据类型

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


【解决方案1】:

您尝试获取所有第一个值的方式会返回不正确的结果。每列值可能来自不同的行。

相反,您应该只按每个组的降序排列order by 时间戳并获取第一行。一个简单的方法是使用像 row_number 这样的函数。

import org.apache.spark.sql.functions._
import org.apache.spark.sql.expressions.Window

val sampleDF = df.withColumn("rnum",row_number().over(Window.partitionBy(col("ID")).orderBy(col("Timestamp").desc)))

sampleDF.filter(col("rnum") == 1).show

【讨论】:

    【解决方案2】:

    只是为了补充 Vamsi 的答案;问题是groupBy 结果组中的值没有以任何特定顺序返回(特别是考虑到 Spark 操作的分布式特性),因此first 函数的命名可能会产生误导。它返回它为该列找到的第一个非空值,即组内该列的几乎所有非空值。

    groupBy 之前对行进行排序不会以任何可重现的方式影响组内的顺序。

    另请参阅此blog post,它解释说,由于上述行为,您从多个first 调用中获得的值甚至可能不是来自组内的同一行。

    输入数据有 3 列“k, t, v”

    z, 1, null
    z, 2, 1.5
    z, 3, 2.4
    

    代码:

    df.groupBy("k").agg(
      $"k",
      first($"t"),
      first($"v")
    )
    

    输出:

    z, 1, 1.5
    

    此结果是 2 条记录的混合!

    【讨论】:

      猜你喜欢
      • 2017-10-26
      • 2016-12-17
      • 2014-06-22
      • 2013-07-19
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-11-30
      • 2011-03-21
      相关资源
      最近更新 更多