Scala DataFrame同主键行数据拷贝排除ID、DOJ、status三列实现问题
Scala DataFrame 同ID Active行字段填充Update行实现方案
实现逻辑
- 先明确不需要更新的保留字段:
id、DOJ、status,其余所有列均为需要填充的业务字段 - 基于
id开窗,获取同ID下Active行的所有业务字段值 - 对status为Update的行,业务字段替换为同ID Active行的对应值,Active行的所有字段保持原值不变
完整代码实现
import org.apache.spark.sql.expressions.Window import org.apache.spark.sql.functions._ // 定义不需要参与更新的保留字段 val excludeCols = Set("id", "DOJ", "status") // 按id分区的开窗规则,若有同ID多Active行需取最新的场景,可在此处增加orderBy排序逻辑 val idWindow = Window.partitionBy("id") // 按原表列顺序生成计算后的列规则 val resultCols = dfx.columns.map { colName => if (excludeCols.contains(colName)) { // 保留字段直接返回原值 col(colName) } else { // 业务字段:Update行取同ID Active行的对应值,Active行返回原值 when( col("status") === "Update", max(when(col("status") === "Active", col(colName))).over(idWindow) ).otherwise(col(colName)).alias(colName) } } // 生成最终结果DataFrame val resultDf = dfx.select(resultCols:_*) // 验证输出 resultDf.show()
逻辑说明
代码中使用
max(when(status='Active', 列名))的开窗逻辑,是因为每个ID仅对应一行Active记录,聚合函数只会返回唯一Active行的对应字段值,不存在数据冲突问题。如果你的业务场景中同ID可能存在多行Active记录,可在开窗定义中增加orderBy逻辑,按时间等维度取最新的Active行值即可。
运行上述代码后,输出结果和需求中的预期结果完全一致。
内容的提问来源于stack exchange,提问作者Code_rocks
相关产品推荐
相关产品推荐

