Spark中如何使用非UDF方法?调用普通函数报错原因及修复方案
我有如下Scala代码,运行正常:
myDF.map{ x => val inp = MyUtils.doSomething(x.value) //接受Int值并返回Int MyInfo(inp) }
其中MyUtils.doSomething是Spark Scala中的普通函数(非UDF)。
但改成下面的写法时出现错误:
val DF = myDF.withColumn("value", lit(MyUtils.doSomething(col("value").asInstanceOf[Int].toInt)))
报错信息:
class org.apache.spark.sql.Column cannot be cast to class java.lang.Integer
请问:
- 报错原因是什么?
- 如何修复?
- 有没有办法获取
col("value")的底层值,以便在doSomething函数中使用? - 为何
col("value").asInstanceOf[Int].toInt无法得到Int值?
1. 报错原因
col("value")返回的是Spark Column对象,它不是实际的Int数值,只是用来描述DataFrame列的逻辑表达式——代表后续要在分布式数据上执行的操作,而非内存中的具体值。你强行用asInstanceOf[Int]把Column对象转成Int,属于类型不匹配,自然会抛出类转换异常。
另外,lit()函数的作用是把本地内存中的常量转换成Column对象,但你这里试图把Column对象当成本地Int传入MyUtils.doSomething,本质上混淆了Spark的分布式计算逻辑和本地单机代码的执行时机:doSomething是本地函数,会在Driver端立即执行,但此时col("value")还没有实际数据,只有执行计划的抽象描述。
2. 修复方法
有两种常用方案:
方案一:将普通函数转为UDF
把MyUtils.doSomething包装成Spark UDF,这样就能在Column API中调用:import org.apache.spark.sql.functions.udf val doSomethingUdf = udf((num: Int) => MyUtils.doSomething(num)) val DF = myDF.withColumn("value", doSomethingUdf(col("value")))UDF会被分发到Executor端,对每一行的实际Int值执行计算,符合Spark的分布式计算模型。
方案二:继续使用map算子(DataFrame转Dataset)
如果你的DataFrame可以转成强类型Dataset,原来的map写法是可行的,因为map会对每一行的具体对象(包含实际的Int值)操作:// 假设myDF对应的case类是MyCaseClass(value: Int) val myDS = myDF.as[MyCaseClass] val resultDS = myDS.map { x => val inp = MyUtils.doSomething(x.value) MyInfo(inp) } // 若需要转回DataFrame val resultDF = resultDS.toDF()
3. 能不能直接获取col("value")的底层值?
不能直接获取。Spark是分布式计算框架,DataFrame的数据分散在各个Executor节点的分区中,Driver端无法直接拿到某一列的所有底层值(除非用collect()把所有数据拉到Driver,但这只适合小数据集,大数据集会导致内存溢出)。
如果必须在Driver端处理数据,只能先把数据收集到本地:
// 仅适合小数据集! val values = myDF.select("value").as[Int].collect() val processedValues = values.map(MyUtils.doSomething) // 再把处理后的数据转成DataFrame val processedDF = spark.createDataFrame(processedValues.map(MyInfo), classOf[MyInfo])
但这种方法不适合大数据场景,因为collect()会把所有数据拉到Driver节点,风险很高。
4. 为何col("value").asInstanceOf[Int].toInt无效?
如前所述,col("value")是Column对象,它的作用是构建SQL执行计划,不是实际的数值。asInstanceOf[Int]是Scala的强制类型转换,只能在类型兼容的情况下生效,但Column和Integer完全是不同的类,转换必然失败。toInt是Int类型的方法,Column对象根本没有这个方法,所以这行代码从逻辑上就不成立。
内容的提问来源于stack exchange,提问作者Oxana Grey

