Spark中按分区获取前一范围的对应字段值
Spark中按分区获取前一范围的对应字段值
嘿,我来帮你搞定这个问题!从你的描述和示例来看,你想要的是针对每个orderid下的不同range,获取上一个range对应的value1/2/3值,而不是同一个range内的上一行数据对吧?
你之前踩的坑我懂——用partitionBy(orderid, range)的话,Spark会把同一个orderid+range的行分到同一个窗口组里,这时候lag()自然只能取同组内的上一行,完全达不到“跨range取前值”的目的。
下面给你一套可行的解决方案,分三步来:
第一步:先聚合同orderid+range的value值
因为从你的示例输入能看到,同一个orderid+range下的value1/2/3都是相同的,所以我们先把每个分组的value值聚合起来(用first()或者max()都可以,结果一致),确保每个orderid+range只有一行数据:
import org.apache.spark.sql.functions._ import org.apache.spark.sql.expressions.Window // 聚合同orderid+range的value,得到每个range的唯一值 val rangeAggDF = df.groupBy("orderid", "range") .agg( first("value1").alias("value1"), first("value2").alias("value2"), first("value3").alias("value3") )
第二步:用窗口函数获取前一个range的value
这时候我们只需要按orderid分区,按range排序,再用lag()就能拿到上一个range的对应值了:
// 定义窗口:按orderid分区,按range升序排列 val windowSpec = Window.partitionBy("orderid").orderBy("range") // 生成previousvalue列 val prevValueDF = rangeAggDF .withColumn("previousvalue1", lag("value1", 1).over(windowSpec)) .withColumn("previousvalue2", lag("value2", 1).over(windowSpec)) .withColumn("previousvalue3", lag("value3", 1).over(windowSpec))
第三步:关联回原表,补全所有行的数据
因为原表可能有多个同orderid+range的行,我们把聚合后的previous值关联回去,就能让每一行都拿到对应的前range值:
// 关联回原表,得到最终结果 val finalDF = df.join(prevValueDF, Seq("orderid", "range"), "left")
这样处理后,就能得到你想要的结果啦——比如orderid=1、range=2的行,previousvalue就是range=1的对应值;range=1的行因为没有前一个range,所以previous值都是null,完全符合你的预期输出。
备注:内容来源于stack exchange,提问作者Shivansh Mehra
相关产品推荐
相关产品推荐

