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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 11:50:24