Scala Spark按客户ID分组比对相邻行列值的问题如何解决
解决方案
你当前的问题是窗口函数没有按客户ID分区,导致lag取的是全表排序后的前一行,出现跨客户对比的问题,仅需要调整窗口定义,添加partitionBy("id")配置即可,修改后的完整代码如下:
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.Column import org.apache.spark.sql.functions._ import spark.implicits._ // 示例数据定义 val df = Seq( (1, "2021-08-15", 10), (1, "2021-08-16", 10), (1, "2021-08-17", 12), (2, "2021-08-15", 5), (2, "2021-08-16", 5) ).toDF("id", "date", "money") def compareCol1(curr: Column, prev: Column): Column = curr === prev // 调整窗口定义:按id分区,分区内按日期排序 val window = Window.partitionBy("id").orderBy("date") val result = df.withColumn("col-comparison", compareCol1($"money", lag("money", 1).over(window))) result.show()
运行后输出结果和预期完全一致:
+---+----------+-----+--------------+ | id| date|money|col-comparison| +---+----------+-----+--------------+ | 1|2021-08-15| 10| null| | 1|2021-08-16| 10| true| | 1|2021-08-17| 12| false| | 2|2021-08-15| 5| null| | 2|2021-08-16| 5| true| +---+----------+-----+--------------+
内容的提问来源于stack exchange,提问作者Agustina
相关产品推荐
相关产品推荐

