Scala Spark 3.0.2根据Map内容更新DataFrame Weight列实现求和的方法
Spark Scala 实现按Map更新DataFrame列的方案
以下两种方案均适配Spark 3.0.2版本,可直接运行:
方案1:内置函数实现(性能最优,推荐)
利用Spark内置的typedLit将Scala原生Map转换为Spark列类型,配合when/otherwise完成匹配更新,全程走Catalyst优化,适合大数据量场景:
import org.apache.spark.sql.functions.{typedLit, when} // 这里用df指代你的原始DataFrame val map: Map[String, Double] = Map("BR1" -> 70.0, "BR300" -> 90.0) // 将Scala Map转为Spark可识别的map列 val weightAddMap = typedLit(map) val resultDF = df.withColumn("Weight", // 匹配到Map键就相加,匹配不到保留原值 when(weightAddMap($"brand").isNotNull, $"Weight" + weightAddMap($"brand")) .otherwise($"Weight") )
如果要更简洁的写法,可以用coalesce替换when逻辑,效果完全一致:
import org.apache.spark.sql.functions.coalesce val resultDF = df.withColumn("Weight", coalesce($"Weight" + weightAddMap($"brand"), $"Weight") )
方案2:UDF实现(写法灵活,适合复杂逻辑)
如果后续映射逻辑会扩展,用UDF写法更易读易维护:
// 定义UDF逻辑 val updateWeight = udf((brand: String, originWeight: Double) => { map.get(brand) match { case Some(addVal) => originWeight + addVal case None => originWeight } }) // 调用UDF更新Weight列 val resultDF = df.withColumn("Weight", updateWeight($"brand", $"Weight"))
执行完成后调用resultDF.show()即可得到你预期的输出结果。
内容的提问来源于stack exchange,提问作者Nab
相关产品推荐
相关产品推荐

