Scala DataFrame如何将前一ID的值赋值给后一ID对应行?
解决Scala DataFrame中按ID传递前一个ID的value1到当前ID的value2问题
你的问题核心是同一ID的所有行需要共享前一个ID的value1值作为value2,而不是简单地按行偏移取前N行的值——这也是你之前用lag($"value1", 2)失败的原因:固定偏移量只适用于每个ID行数固定的场景,一旦ID的行数变化(比如示例中的id3有3行),后续行的value2就会取到当前ID自己的值,不符合预期。
正确解决思路
因为同一ID的value1值是相同的,我们可以先提取每个ID对应的唯一value1,再为每个ID关联前一个ID的value1,最后将结果回关联到原始DataFrame中,确保同一ID的所有行都拿到正确的value2。
具体代码实现
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions.lag // 1. 提取每个ID对应的唯一value1(因为同一ID的value1完全相同) val idUniqueValue = df.select("id", "value1").distinct().orderBy("id") // 2. 为每个ID生成前一个ID的value1作为value2 val windowSpec = Window.orderBy("id") val idWithPrevValue = idUniqueValue.withColumn("value2", lag($"value1", 1).over(windowSpec)) // 3. 关联回原始DataFrame,得到最终结果 val resultDF = df.join(idWithPrevValue, Seq("id", "value1"), "left") // 查看结果 resultDF.show()
结果验证
执行后输出的DataFrame会完全符合你的预期:
+---+------+------+ | id|value1|value2| +---+------+------+ |id1| ab| null| |id1| ab| null| |id2| ac| ab| |id2| ac| ab| |id3| abc| ac| |id3| abc| ac| |id3| abc| ac| +---+------+------+
为什么这个方法可行?
- 先聚合去重:避免了重复处理同一ID的多行数据,提升效率;
- 针对ID级别的lag:在去重后的ID列表上使用
lag(1),确保每个ID拿到的是前一个ID的value1,不受当前ID行数影响; - 左关联回原表:保证原始DataFrame的所有行都能匹配到对应的
value2,不会丢失数据。
内容的提问来源于stack exchange,提问作者user9601368
相关产品推荐
相关产品推荐

