You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.14 08:37:33