Spark中如何实现值映射?州名转缩写时类型不匹配如何解决
解决方法
错误原因
你遇到的类型不匹配错误是因为自定义的普通Scala函数只能接收字符串类型参数,但你传入的$"col.someKey.keyStoringValue"是Spark的ColumnName(Column类的子类)类型,无法直接传入普通Scala函数执行,需要将自定义逻辑封装为Spark支持的UDF,或者使用内置函数实现映射逻辑。
方案1:使用UDF封装自定义逻辑
UDF(用户自定义函数)可以适配Spark的列类型参数,实现代码如下:
// 导入udf工具类 import org.apache.spark.sql.functions.udf // 定义并注册UDF val getStateAbbrUDF = udf((state: String) => { val stateMap = Map("Texas" -> "TX", "Virginia" -> "VA") // getOrElse避免匹配不到的场景抛出异常,匹配失败默认返回原字符串 stateMap.getOrElse(state, state) }) // 在withColumn中调用UDF .withColumn("NewColName", getStateAbbrUDF($"col.someKey.keyStoringValue"))
方案2:使用Spark内置函数实现(性能更优)
如果映射关系数量不多,推荐使用内置create_map函数实现映射,避免UDF带来的序列化、反序列化性能开销:
import org.apache.spark.sql.functions.{create_map, lit, col, coalesce} // 构造映射关系列 val stateMapCol = create_map( lit("Texas"), lit("TX"), lit("Virginia"), lit("VA") ) // 取值映射,coalesce用于匹配失败时返回原字段值作为默认值 .withColumn("NewColName", coalesce(stateMapCol(col("col.someKey.keyStoringValue")), col("col.someKey.keyStoringValue")))
内容的提问来源于stack exchange,提问作者Bien
相关产品推荐
相关产品推荐

