【发布时间】:2021-07-13 00:23:45
【问题描述】:
我有一个如下所示的带有 n 列的数据框。
+---+------------+--------+--------+--------+
|id | date|signal01|signal02|signal03|......signal(n)
+---+------------+--------+--------+--------+
|050|2021-01-14 |1 |3 |1 |
|050|2021-01-15 |null |4 |2 |
|050|2021-02-02 |2 |5 |3 |
|051|2021-01-14 |1 |3 |0 |
|051|2021-01-15 |null |null |null |
|051|2021-02-02 |3 |3 |2 |
|051|2021-02-03 |4 |3 |3 |
|052|2021-03-03 |1 |3 |0 |
|052|2021-03-05 |null |3 |null |
|052|2021-03-06 |null |null |2 |
|052|2021-03-16 |3 |5 |5 |.......value(n)
+-------------------------------------------+
我必须为每个信号添加一个信号差异值列,如下所示,不包括空记录并将第一个差异值保持为 0。
+---+------------+--------+-------------+--------+-------------+--------+-------------+
|id | date|signal01|signal01_diff|signal02|signal02_diff|signal03|signal03_diff|......signal(n)
+---+------------+--------+-------------+--------+-------------+--------+-------------+
|050|2021-01-14 |1 |0 |3 |0 |1 |0 |
|050|2021-01-15 |null |null |4 |1 |2 |1 |
|050|2021-02-02 |2 |1 |5 |1 |3 |1 |
|051|2021-01-14 |1 |0 |3 |0 |0 |0 |
|051|2021-01-15 |null |null |null |null |null |null |
|051|2021-02-02 |3 |2 |3 |0 |2 |2 |
|051|2021-02-03 |4 |1 |3 |0 |3 |1 |
|052|2021-03-03 |1 |0 |3 |0 |0 |0 |
|052|2021-03-05 |null |null |3 |0 |null |null |
|052|2021-03-06 |null |null |null |null |2 |2 |
|052|2021-03-16 |3 |2 |5 |2 |5 |3 |.......value(n)
+-----------------------------------------------------------------------+--------------
我尝试了延迟和窗口函数,但由于空值而没有得到预期的输出。
val w = org.apache.spark.sql.expressions.Window.orderBy("id")
val dfWithLag = df.withColumn("signal01_lag", lag("signal01", 1, 0).over(w))
以上是单列的代码,我必须为其余 n 列执行相同的代码。
有没有最佳的方法来实现这一点?
【问题讨论】:
标签: scala apache-spark apache-spark-sql