【问题标题】:Look Back N Months and add the them as columns spark sql aggregate回顾 N 个月并将它们添加为列 spark sql 聚合
【发布时间】:2019-06-23 05:41:25
【问题描述】:

我必须每三个月回顾一次,并使用 with 列添加上个月的金额。

val data = Seq(("1","201706","5"),("1","201707","10"),("2","201604","12"),("2","201601","15")).toDF("id","yyyyMM","amount")

+---+------+------+
| id|yyyyMM|amount|
+---+------+------+
|  1|201706|     5|
|  1|201707|    10|
|  2|201604|    12|
|  2|201601|    15|
+---+------+------+

所需的输出应如下所示。对于每个月我们必须回顾三个月,我可以通过使用火花窗口滞后功能来做到这一点。我们应该如何包含添加附加记录的功能

+---+---------+------+-----------+-------+-----------+-------+
| id|yearmonth|amount|yearmonth-1|amount2|yearmonth-2|amount3|
+---+---------+------+-----------+-------+-----------+-------+
|  1|   201709|     0|     201708|      0|     201707|     10|
|  1|   201708|     0|     201707|     10|     201706|      5|
|  1|   201707|    10|     201706|      5|     201705|      0|
|  1|   201706|     5|     201705|      0|     201706|      0|
|  2|   201606|     0|     201605|      0|     201604|     12|
|  2|   201605|     0|     201604|     12|     201603|      0|
|  2|   201604|    12|     201603|      0|     201602|      0|
|  2|   201603|     0|     201602|      0|     201601|     15|
|  2|   201602|     0|     201601|     15|     201512|      0|
|  2|   201601|    15|     201512|      0|     201511|      0|
+---+---------+------+-----------+-------+-----------+-------+

我的意思是表中的第一条记录就像期待。就像增加几个月一样。采取以下记录。

+---+---------+------+-----------+-------+-----------+-------+
| id|yearmonth|amount|yearmonth-1|amount2|yearmonth-2|amount3|
+---+---------+------+-----------+-------+-----------+-------+
|  1|   201709|     0|     201708|      0|     201707|     10|
|  1|   201708|     0|     201707|     10|     201706|      5|

【问题讨论】:

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


    【解决方案1】:

    我不知道是否有更好的方法,但您需要在某处创建记录。 Lag 不会那样做。 因此,首先您需要根据当前记录生成新记录。 然后你可以使用滞后功能。

    可能是这样的:

    data
      // convert the string to an actual date
      .withColumn("yearmonth", to_date('yyyyMM, "yyyyMM"))
      // for each record create 2 additional in the future (with 0 amount)
      .select(
      explode(array(
        // org record
        struct('id, date_format('yearmonth, "yyyyMM").as("yearmonth"), 'amount),
        // 1 month in future
        struct('id, date_format(add_months('yearmonth, 1), "yyyyMM").as("yearmonth"), lit(0).as("amount")),
        // 2 months in future
        struct('id, date_format(add_months('yearmonth, 2), "yyyyMM").as("yearmonth"), lit(0).as("amount"))
      )).as("record"))
      // keep 1 record per month
      .groupBy($"record.yearmonth")
      .agg(
        min($"record.id").as("id"),
        sum($"record.amount").as("amount")
      )
      // final structure (with lag fields)
      .select(
        'id,
        'yearmonth,
        'amount,
         lag('yearmonth, 1).over(orderByWindow).as("yearmonth-1"),
         lag('amount, 1, 0).over(orderByWindow).as("amount2"),
         lag('yearmonth, 2).over(orderByWindow).as("yearmonth-2"),
         lag('amount, 2, 0).over(orderByWindow).as("amount3")
      )
      .orderBy('yearmonth.desc)
    

    这并不完美,但这是一个开始

    +---+---------+------+-----------+-------+-----------+-------+
    |id |yearmonth|amount|yearmonth-1|amount2|yearmonth-2|amount3|
    +---+---------+------+-----------+-------+-----------+-------+
    |1  |201709   |0.0   |201708     |0.0    |201707     |10.0   |
    |1  |201708   |0.0   |201707     |10.0   |201706     |5.0    |
    |1  |201707   |10.0  |201706     |5.0    |201606     |0.0    |
    |1  |201706   |5.0   |201606     |0.0    |201605     |0.0    |
    |2  |201606   |0.0   |201605     |0.0    |201604     |12.0   |
    |2  |201605   |0.0   |201604     |12.0   |201603     |0.0    |
    |2  |201604   |12.0  |201603     |0.0    |201602     |0.0    |
    |2  |201603   |0.0   |201602     |0.0    |201601     |15.0   |
    |2  |201602   |0.0   |201601     |15.0   |null       |0.0    |
    |2  |201601   |15.0  |null       |0.0    |null       |0.0    |
    +---+---------+------+-----------+-------+-----------+-------+
    

    【讨论】:

    • 我认为我们应该为前几个月再增加两个月以使其连续。
    • 我在程序中使用时遇到结构、数组和转换问题的错误。我想我错过了一些进口。
    猜你喜欢
    • 2015-08-11
    • 1970-01-01
    • 2013-08-09
    • 1970-01-01
    • 2010-12-09
    • 2018-07-24
    • 2015-05-23
    • 2015-01-26
    相关资源
    最近更新 更多