Spark Scala UDF开发遇到Task not serializable序列化错误求助
解决Spark UDF的Task not serializable错误
你的代码出现这个错误,核心原因是ageCheck方法属于dataframes这个App对象的成员方法,Spark在将UDF分发到Executor节点时,需要序列化整个dataframes对象,但App类默认没有实现Serializable接口,导致序列化失败。下面是几种可行的解决办法:
方法1:将ageCheck定义为独立对象的静态方法
把判断方法放到单独的object里,方法作为静态存在,无需依赖外部对象的序列化:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.Dataset import org.apache.spark.sql.Row object AgeChecker { // 独立对象中的静态方法,可直接序列化 def ageCheck(age: Int): String = if (age > 18) "Y" else "N" } object dataframes extends App { val spark = SparkSession.builder().appName("testing").master("local[*]") .getOrCreate() val withOutHeaderDF = spark.read.format("csv") .option("inferSchema", true) .option("path","D:/1. Technologies/Big Data/sharedPath/dataset1.csv") .load() val withHeaderDF: Dataset[Row] = withOutHeaderDF.toDF("name","age","city") import spark.implicits._ // 使用独立对象中的方法创建UDF val definedAgeCheck = udf(AgeChecker.ageCheck(_: Int)) val finalResult = withHeaderDF.withColumn("Adult",definedAgeCheck(col("age"))) finalResult.show() spark.stop() }
方法2:直接用匿名函数创建UDF
省去单独的方法定义,直接把逻辑写入udf构造函数,避免依赖外部对象:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.Dataset import org.apache.spark.sql.Row object dataframes extends App { val spark = SparkSession.builder().appName("testing").master("local[*]") .getOrCreate() val withOutHeaderDF = spark.read.format("csv") .option("inferSchema", true) .option("path","D:/1. Technologies/Big Data/sharedPath/dataset1.csv") .load() val withHeaderDF: Dataset[Row] = withOutHeaderDF.toDF("name","age","city") import spark.implicits._ // 直接传入匿名函数构建UDF val definedAgeCheck = udf((age: Int) => if (age > 18) "Y" else "N") val finalResult = withHeaderDF.withColumn("Adult",definedAgeCheck(col("age"))) finalResult.show() spark.stop() }
方法3:用Spark内置函数替代UDF(推荐)
对于这种简单的判断逻辑,完全不需要自定义UDF,用Spark自带的when和otherwise函数更高效,还能彻底避免序列化问题:
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions._ import org.apache.spark.sql.Dataset import org.apache.spark.sql.Row object dataframes extends App { val spark = SparkSession.builder().appName("testing").master("local[*]") .getOrCreate() val withOutHeaderDF = spark.read.format("csv") .option("inferSchema", true) .option("path","D:/1. Technologies/Big Data/sharedPath/dataset1.csv") .load() val withHeaderDF: Dataset[Row] = withOutHeaderDF.toDF("name","age","city") // 使用Spark内置函数实现判断逻辑 val finalResult = withHeaderDF.withColumn("Adult", when(col("age") > 18, "Y").otherwise("N")) finalResult.show() spark.stop() }
内容的提问来源于stack exchange,提问作者sami phani
相关产品推荐
相关产品推荐

