UDF内ProductEncoder工作异常:该现象是否为预期行为?
问题:Spark UDF内部使用case class及toDF的异常现象分析
场景一:UDF外部定义case class,内部调用toDF触发运行时异常
在Databricks笔记本中运行以下Scala代码:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.sql.{DataFrame, Dataset, SparkSession} case class DataRow(field1: String) val sparkSession = SparkSession.builder.getOrCreate() import sparkSession.implicits._ val udf1 = udf((x: String) => { val testData = Seq(DataRow("test1"), DataRow("test2")).toDF("test") // 运行时失败 3 }) val df1 = Seq(DataRow("test1"), DataRow("test2")).toDF("test").withColumn("udf", udf1($"test")) // 正常运行 display(df1)
运行时抛出ScalaReflectionException异常,错误栈信息如下:
ScalaReflectionException: class $linedb3da7b2933d4b63a62b3d2a21c675f2141.$read in JavaMirror with com.databricks.backend.daemon.driver.DriverLocal$DriverLocalClassLoader@717115ad of type class com.databricks.backend.daemon.driver.DriverLocal$DriverLocalClassLoader with classpath [] and parent being com.databricks.backend.daemon.driver.ClassLoaders$ReplWrappingClassLoader@1670897 of type class com.databricks.backend.daemon.driver.ClassLoaders$ReplWrappingClassLoader with classpath [<unknown>] and parent being com.databricks.backend.daemon.driver.ClassLoaders$LibraryClassLoader@1a42da0a of type class com.databricks.backend.daemon.driver.ClassLoaders$LibraryClassLoader with classpath [file:/local_disk0/tmp/repl/spark-4071811259476162981-8e526ae9-25fb-4545-8d3f-963a8661cd2b/] and parent being sun.misc.Launcher$AppClassLoader@43ee72e6 of type class sun.misc.Launcher$AppClassLoader with classpath [...] not found. at scala.reflect.internal.Mirrors$RootsBase.staticClass(Mirrors.scala:141) at scala.reflect.internal.Mirrors$RootsBase.staticClass(Mirrors.scala:29) at $linedb3da7b2933d4b63a62b3d2a21c675f2196.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$$$c35395a9233a7197629c985e87dce75$$$$w$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$typecreator6$1.apply(command-2013905963200886:11) at scala.reflect.api.TypeTags$WeakTypeTagImpl.tpe$lzycompute(TypeTags.scala:237) at scala.reflect.api.TypeTags$WeakTypeTagImpl.tpe(TypeTags.scala:237) at org.apache.spark.sql.catalyst.ScalaReflection$.encoderFor(ScalaReflection.scala:848) at org.apache.spark.sql.catalyst.encoders.ExpressionEncoder$.apply(ExpressionEncoder.scala:55) at org.apache.spark.sql.Encoders$.product(Encoders.scala:312) at org.apache.spark.sql.LowPrioritySQLImplicits.newProductEncoder(SQLImplicits.scala:302) at org.apache.spark.sql.LowPrioritySQLImplicits.newProductEncoder$(SQLImplicits.scala:302) at org.apache.spark.sql.SQLImplicits.newProductEncoder(SQLImplicits.scala:34) at $linedb3da7b2933d4b63a62b3d2a21c675f2196.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$$$c35395a9233a7197629c985e87dce75$$$$w$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.$anonfun$udf1$1(command-2013905963200886:11)
场景二:UDF内部定义case class触发编译错误
若将case class定义在UDF内部:
import org.apache.spark.{SparkConf, SparkContext} import org.apache.spark.sql.{DataFrame, Dataset, SparkSession} case class DataRow(field1: String) val sparkSession = SparkSession.builder.getOrCreate() import sparkSession.implicits._ val udf1 = udf((x: String) => { case class DataRow(field1: String) val testData = Seq(DataRow("test1"), DataRow("test2")).toDF("test") // 运行时失败 3 }) val df1 = Seq(DataRow("test1"), DataRow("test2")).toDF("test").withColumn("udf", udf1($"test")) // 正常运行 display(df1)
则会出现编译错误:
error: value toDF is not a member of Seq[DataRow] val testData = Seq(DataRow("test1"), DataRow("test2")).toDF("test")
核心疑问
上述两种现象属于Spark的预期行为还是Bug?已向Apache Spark提交Bug工单SPARK-44706。
内容的提问来源于stack exchange,提问作者Juan Carlos Blanco Martínez
相关产品推荐
相关产品推荐

