Scala SparkSQL创建UDF处理列(结构体/字符串)类型兼容问题
解决Spark中混合类型列的字段提取问题
我来帮你搞定这个棘手的混合类型列处理问题!首先得说清楚你原来的UDF为啥失败,再给你两种更靠谱的解决思路。
问题根源:参数类型不匹配
你写的UDF把参数声明成了Row,但问题是当annoyingCol.data是字符串类型时,Spark会把字符串值传入UDF,而不是Row,这就直接导致了类型不匹配的错误——毕竟字符串没法当成Row来处理嘛。
方案一:用Spark内置函数(推荐!)
Spark本身就提供了typeof函数,可以直接判断列的数据类型,完全不需要自己写UDF,性能还更好。
SQL写法
SELECT CASE WHEN typeof(annoyingCol.data) = 'struct' THEN annoyingCol.data.my_data ELSE NULL END FROM SomeDf
DataFrame API写法
import org.apache.spark.sql.functions._ val resultDf = someDf.withColumn( "extracted_my_data", when( typeof(col("annoyingCol.data")) === "struct", col("annoyingCol.data.my_data") ).otherwise(lit(null)) )
方案二:修正你的UDF(如果必须用UDF的话)
要是你确实需要用UDF来实现,那得把参数类型改成Any,然后通过模式匹配判断传入的值是不是Row(Spark里结构体类型的数据就是用Row来表示的):
import org.apache.spark.sql.functions.udf import org.apache.spark.sql.Row // 定义UDF val isStruct = udf((value: Any) => { value match { case _: Row => true case _ => false } }) // 注册成SQL可用的函数 spark.udf.register("isStruct", isStruct)
然后就可以像你原来想的那样用了:
SELECT CASE WHEN isStruct(annoyingCol.data) THEN annoyingCol.data.my_data ELSE NULL END FROM SomeDf
小提醒
尽量优先用Spark的内置函数,因为UDF需要进行数据的序列化和反序列化,性能比内置函数差不少,而且内置函数已经能完美解决这个问题啦~
内容的提问来源于stack exchange,提问作者Benny Elgazar
相关产品推荐
相关产品推荐

