【问题标题】:How to aggregate over rolling time window with groups in Spark如何在 Spark 中与组聚合滚动时间窗口
【发布时间】:2017-06-02 09:02:39
【问题描述】:

我有一些数据要按某个列分组,然后根据该组的滚动时间窗口聚合一系列字段。

以下是一些示例数据:

df = spark.createDataFrame([Row(date='2016-01-01', group_by='group1', get_avg=5, get_first=1),
                            Row(date='2016-01-10', group_by='group1', get_avg=5, get_first=2),
                            Row(date='2016-02-01', group_by='group2', get_avg=10, get_first=3),
                            Row(date='2016-02-28', group_by='group2', get_avg=20, get_first=3),
                            Row(date='2016-02-29', group_by='group2', get_avg=30, get_first=3),
                            Row(date='2016-04-02', group_by='group2', get_avg=8, get_first=4)])

我想按group_by 进行分组,然后创建从最早日期开始的时间窗口,并一直延伸到该组没有条目的 30 天。在这 30 天结束后,下一个时间窗口将从不属于上一个窗口的下一行的日期开始。

然后我想聚合,例如得到get_avg 的平均值,以及get_first 的第一个结果。

所以这个例子的输出应该是:

group_by    first date of window    get_avg  get_first
group1      2016-01-01              5        1
group2      2016-02-01              20       3
group2      2016-04-02              8        4

编辑:抱歉,我意识到我的问题没有正确指定。我实际上想要一个在 30 天不活动后结束的窗口。我已经相应地修改了示例的 group2 部分。

【问题讨论】:

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


    【解决方案1】:

    修改后的答案

    您可以在这里使用一个简单的窗口函数技巧。一堆进口:

    from pyspark.sql.functions import coalesce, col, datediff, lag, lit, sum as sum_
    from pyspark.sql.window import Window
    

    窗口定义:

    w = Window.partitionBy("group_by").orderBy("date")
    

    date 转换为DateType

    df_ = df.withColumn("date", col("date").cast("date"))
    

    定义以下表达式:

    # Difference from the previous record or 0 if this is the first one
    diff = coalesce(datediff("date", lag("date", 1).over(w)), lit(0))
    
    # 0 if diff <= 30, 1 otherwise
    indicator = (diff > 30).cast("integer")
    
    # Cumulative sum of indicators over the window
    subgroup = sum_(indicator).over(w).alias("subgroup")
    

    subgroup 表达式添加到表中:

    df_.select("*", subgroup).groupBy("group_by", "subgroup").avg("get_avg")
    
    +--------+--------+------------+
    |group_by|subgroup|avg(get_avg)|
    +--------+--------+------------+
    |  group1|       0|         5.0|
    |  group2|       0|        20.0|
    |  group2|       1|         8.0|
    +--------+--------+------------+
    

    first 对聚合没有意义,但如果列单调递增,您可以使用min。否则你也必须使用窗口函数。

    使用 Spark 2.1 测试。与较早的 Spark 版本一起使用时,可能需要子查询和 Window 实例。

    原答案(与指定范围无关)

    从 Spark 2.0 开始你应该可以使用a window function:

    在给定时间戳指定列的情况下,将行分桶到一个或多个时间窗口中。窗口开始是包含的,但窗口结束是排除的,例如12:05 将在窗口 [12:05,12:10) 中,但不在 [12:00,12:05) 中。

    from pyspark.sql.functions import window
    
    df.groupBy(window("date", windowDuration="30 days")).count()
    

    但你可以从结果中看到,

    +---------------------------------------------+-----+
    |window                                       |count|
    +---------------------------------------------+-----+
    |[2016-01-30 01:00:00.0,2016-02-29 01:00:00.0]|1    |
    |[2015-12-31 01:00:00.0,2016-01-30 01:00:00.0]|2    |
    |[2016-03-30 02:00:00.0,2016-04-29 02:00:00.0]|1    |
    +---------------------------------------------+-----+
    

    在时区方面你必须小心一点。

    【讨论】:

    • 你能解释一下总和是如何在这里工作的吗?计算这个总和时窗口的框架是什么?如果一组之间有 90 天的间隔,那么 sub 的数量将是 0、1、2?
    • @G_cy 它实际上是一个指标函数的总和 - 应用指标产生一个 {0, 1} 序列,其中每个 1 表示差距 > 阈值。 sum 只是这个序列的累积和。所以你的问题的答案是否定的。它不区分间隙 1 * 阈值、2 * 阈值 ... n * 阈值。
    • 我一步一步检查了。在一组中,我知道该指标产生一个序列 {0, 1}。在一组中,它显示为:(g1, 0), (g1, 0).....(g1, 1)...(g1, 0)...(g1, 1)。然后我将 sum 放入show,它显示 (g1,0), (g1, 0), (g1,0) ...(g1, 1), (g1, 1), (g1, 1).. .. 看起来sum 只会总结当前行之前的所有 {0, 1} 序列。是真的吗?
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2019-03-16
    • 2021-09-05
    • 1970-01-01
    • 1970-01-01
    • 2019-10-15
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多