如何在Spark Scala DataFrame中按分区获取最近的两个值
在Spark Scala中按分组查找DataFrame中最近的两个值
需求是按id分组,为每组的value1找到value2中最接近的两个值:一个是不大于value1的最大值,另一个是不小于value1的最小值,最终保留这两个值对应的行。
方法一:分组聚合关联法
这种方法先通过分组聚合计算每组的上下边界值,再关联原数据筛选目标行,逻辑直观易理解。
1. 创建示例DataFrame
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ val spark = SparkSession.builder().appName("NearestValues").master("local[*]").getOrCreate() import spark.implicits._ // 构造输入数据 val inputDF = Seq( (1, 3, 1), (1, 3, 2), (1, 3, 7), (2, 4, 2), (2, 4, 3), (2, 4, 8), (3, 5, 3), (3, 5, 6), (3, 5, 7), (3, 5, 8) ).toDF("id", "value1", "value2")
2. 计算每组的上下边界值
按id和value1分组,分别计算每组中小于等于value1的最大value2(lower_val),以及大于等于value1的最小value2(upper_val):
val boundaryDF = inputDF.groupBy("id", "value1") .agg( max(when(col("value2") <= col("value1"), col("value2"))).alias("lower_val"), min(when(col("value2") >= col("value1"), col("value2"))).alias("upper_val") )
3. 关联筛选目标行
将原数据与边界值DataFrame关联,筛选出value2等于lower_val或upper_val的行:
val resultDF = inputDF.join(boundaryDF, Seq("id", "value1"), "inner") .filter(col("value2") === col("lower_val") || col("value2") === col("upper_val")) .select("id", "value1", "value2") .orderBy("id", "value2") // 查看结果 resultDF.show()
输出结果与需求一致:
+---+------+------+ | id|value1|value2| +---+------+------+ | 1| 3| 2| | 1| 3| 7| | 2| 4| 3| | 2| 4| 8| | 3| 5| 3| | 3| 5| 6| +---+------+------+
方法二:窗口函数法
利用窗口函数在分组内排序后,直接标记符合条件的行,适合需要保留更多中间逻辑的场景。
import org.apache.spark.sql.expressions.Window // 定义按id、value1分区,按value2排序的窗口 val windowSpec = Window.partitionBy("id", "value1").orderBy("value2") // 添加标记列,标记是否是最后一个<=value1的行,或是第一个>=value1的行 val flaggedDF = inputDF .withColumn("is_lower", last(when(col("value2") <= col("value1"), true).otherwise(false), ignoreNulls = true).over(windowSpec)) .withColumn("is_upper", first(when(col("value2") >= col("value1"), true).otherwise(false), ignoreNulls = true).over(windowSpec)) // 筛选标记为true的行 val resultWindowDF = flaggedDF .filter(col("is_lower") || col("is_upper")) .select("id", "value1", "value2") .orderBy("id", "value2") resultWindowDF.show()
这个方法的输出结果和方法一完全相同。
注意事项
- 如果分组内存在
value2等于value1的情况,两个方法都会自动将其识别为同时满足lower和upper条件,最终只保留一行(避免重复)。 - 确保分组键是
id和value1,因为同一id下可能存在不同的value1(如果有这种场景的话)。
内容的提问来源于stack exchange,提问作者AmidOV
相关产品推荐
相关产品推荐

