【问题标题】:Spark-SQL Window functions on Dataframe - Finding first timestamp in a groupDataframe 上的 Spark-SQL 窗口函数 - 查找组中的第一个时间戳
【发布时间】:2016-05-20 20:13:12
【问题描述】:

我有以下数据框(比如 UserData)。

uid region  timestamp
a   1   1
a   1   2
a   1   3
a   1   4
a   2   5
a   2   6
a   2   7
a   3   8
a   4   9
a   4   10
a   4   11
a   4   12
a   1   13
a   1   14
a   3   15
a   3   16
a   5   17
a   5   18
a   5   19
a   5   20

这些数据只不过是用户(uid)在不同时间(时间戳)穿越不同区域(区域)的用户(uid)。目前,为简单起见,时间戳显示为“int”。请注意,上述数据帧不一定按时间戳的递增顺序。此外,不同用户之间可能存在一些行。为简单起见,我仅以时间戳的单调递增顺序显示了单个用户的数据框。

我的目标是——找出用户“a”在每个区域花费了多少时间以及以什么顺序?所以我最终的预期输出看起来像

uid region  regionTimeStart regionTimeEnd
a   1   1   5
a   2   5   8
a   3   8   9
a   4   9   13
a   1   13  15
a   3   15  17
a   5   17  20

根据我的发现,Spark SQL Window 函数可用于此目的。 我已经尝试过以下事情,

val w = Window
  .partitionBy("region")
  .partitionBy("uid")
  .orderBy("timestamp")

val resultDF = UserData.select(
  UserData("uid"), UserData("timestamp"),
  UserData("region"), rank().over(w).as("Rank"))

但从这里开始,我不确定如何获取 regionTimeStartregionTimeEnd 列。 regionTimeEnd 列只不过是regionTimeStart 的“领先”,除了组中的最后一个条目。

我看到聚合操作具有“第一”和“最后”功能,但为此我需要根据 ('uid','region') 对数据进行分组,这会破坏遍历路径的单调递增顺序,即在时间 13,14 用户已回到区域“1”,我希望保留它,而不是在时间 1 将其与初始区域“1”混为一谈。

如果有人可以指导我,那将非常有帮助。我是 Spark 新手,与 Python/JAVA Spark API 相比,我对 Scala Spark API 有更好的理解。

【问题讨论】:

    标签: sql apache-spark dataframe apache-spark-sql window-functions


    【解决方案1】:

    窗口函数确实很有用,尽管您的方法只有在您假设用户只访问给定区域一次时才有效。您使用的窗口定义也不正确 - 多次调用 partitionBy 只会返回具有不同窗口定义的新对象。如果您想按多列进行分区,您应该在一次调用中传递它们 (.partitionBy("region", "uid"))。

    让我们从标记每个区域的连续访问开始:

    import org.apache.spark.sql.functions.{lag, sum, not}
    import org.apache.spark.sql.expressions.Window 
    
    val w = Window.partitionBy($"uid").orderBy($"timestamp")
    
    val change = (not(lag($"region", 1).over(w) <=> $"region")).cast("int")
    val ind = sum(change).over(w)
    
    val dfWithInd = df.withColumn("ind", ind)
    

    接下来,我们只需汇总组并找到潜在客户:

    import org.apache.spark.sql.functions.{lead, coalesce}
    
    val regionTimeEnd = coalesce(lead($"timestamp", 1).over(w), $"max_")
    
    val result = dfWithInd
      .groupBy($"uid", $"region", $"ind")
      .agg(min($"timestamp").alias("timestamp"), max($"timestamp").alias("max_"))
      .drop("ind")
      .withColumn("regionTimeEnd", regionTimeEnd)
      .withColumnRenamed("timestamp", "regionTimeStart")
      .drop("max_")
    
    result.show
    
    // +---+------+---------------+-------------+
    // |uid|region|regionTimeStart|regionTimeEnd|
    // +---+------+---------------+-------------+
    // |  a|     1|              1|            5|
    // |  a|     2|              5|            8|
    // |  a|     3|              8|            9|
    // |  a|     4|              9|           13|
    // |  a|     1|             13|           15|
    // |  a|     3|             15|           17|
    // |  a|     5|             17|           20|
    // +---+------+---------------+-------------+
    

    【讨论】:

    • 您好 zero323,它按预期工作。非常感谢您提供快速详细的解决方案。我一直坚持使用没有给我正确方法的“第一个”和“最后一个”API。我可以扩展您的方法来解决我要解决的许多其他问题。谢谢。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-04-21
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多