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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 12:39:04