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

基于ID与Version连接Dataset,匹配最近前序版本的实现问题

解决方案:Spark中基于ID和版本匹配最近前序版本

核心思路

先通过左连接关联同ID且右侧版本小于左侧版本的所有候选记录,再通过窗口函数对每个左侧记录的候选右侧记录按版本降序排序,取第一条即为最近的前序版本。

实现代码

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

case class Sample(id: String, version: Int, recordValue: String)

val left = Seq(
     Sample("1", 5, "left5"),
     Sample("1", 6, "left6")
   ).toDS()

val right = Seq(
     Sample("1", 6, "right6"),
     Sample("1", 5, "right5"),
     Sample("1", 3, "right3"),
     Sample("1", 2, "right2"),
     Sample("1", 1, "right1")
   ).toDS()

// 1. 左连接:关联同ID,且右侧版本小于左侧版本的所有候选记录
val joined = left.join(right,
  left("id") === right("id") && right("version") < left("version"),
  "left_outer"
)

// 2. 定义窗口:按左侧的id和version分区,按右侧version降序排序
val windowSpec = Window
  .partitionBy(left("id"), left("version"))
  .orderBy(desc(right("version")))

// 3. 取每个分组的第一条记录(即最大的符合条件的右侧版本)
val result = joined
  .withColumn("row_num", row_number().over(windowSpec))
  .filter(col("row_num") === 1)
  .select(
    struct(left("id"), left("version"), left("recordValue")).alias("left_sample"),
    struct(right("id"), right("version"), right("recordValue")).alias("right_sample")
  )
  .as[(Sample, Sample)]

// 查看结果
result.show(false)

输出结果

+-------------------+-------------------+
|left_sample         |right_sample       |
+-------------------+-------------------+
|Sample(1,5,left5)   |Sample(1,3,right3) |
|Sample(1,6,left6)   |Sample(1,5,right5) |
+-------------------+-------------------+

说明

  • 该方案会自动优先匹配left.version-1的版本(如果存在),因为它是小于left.version的最大版本
  • 若右侧没有任何小于左侧版本的记录,right_sample会返回null,符合左外连接的预期
  • 窗口函数的方式适合分布式场景,性能优于UDF聚合版本列表的方案

内容的提问来源于stack exchange,提问作者SRN

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 17:03:11