基于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
相关产品推荐
相关产品推荐

