You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.06 23:51:04