Spark UDF返回运行时定义Schema的Row遇阻,求解决方案
背景
正在构建一个ETL管道,从Kafka队列获取GNMI Protobuf更新消息,最终基于值路径的前缀和参数拆分至多个Delta表(运行于Databricks Runtime)。每个前缀对应一个表的Schema,但上游路径可能新增子树,因此Schema并非固定,类似嵌套JSON结构。
已按前缀拆分更新,确保同批更新Schema一致,并定义了转换逻辑可将不匹配Schema强制转为通用Schema,但在创建带通用Schema的struct列时遇到问题。
尝试1:废弃的UDF指定Schema方式
先尝试从UDF返回Array[Any]并在UDF定义中指定Schema(已知该方式已废弃),代码如下:
import org.apache.spark.sql.{functions => F, Row, types => T} def mapToRow(deserialized: Map[String, ParsedValueV2]): Array[Any] = { def getValue(key: String): Any = { deserialized.get(key) match { case Some(value) => value.asType(columns(key)) case None => None } } columns.keys.toArray.map(getValue).toArray } spark.conf.set("spark.sql.legacy.allowUntypedScalaUDF", "true") def mapToStructUdf = F.udf(mapToRow _, account.sparkSchemas(prefix))
代码生成了包含所需类型值的Array对象,但执行时抛出以下错误:
java.lang.ClassCastException: org.apache.spark.sql.catalyst.expressions.GenericRowWithSchema cannot be cast to $line8760b7c10da04d2489451bb90ca42c6535.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$ParsedValueV2
推测问题可能与值为Java类型而非Scala类型有关,但不确定具体原因。
尝试2:运行时创建Case Class的Typed UDF思路
考虑使用Typed UDF接口,尝试在运行时为每个Schema创建case class并作为UDF返回值,代码如下:
import scala.reflect.runtime.universe import scala.tools.reflect.ToolBox val tb = universe.runtimeMirror(getClass.getClassLoader).mkToolBox() val test = tb.eval(tb.parse("object Test; Test"))
但无法获取test实例,也不知道如何将其作为UDF返回值。推测需要泛型实现,但Scala能力不足无法推进。
核心问题
请指明应采用哪种方法实现动态Schema的struct列,并提供具体步骤。
更新:疑似Spark Bug?
将问题简化后复现异常:
import org.apache.spark.sql.{functions => F, Row, types => T} val spark = SparkSession.builder .master ("local") .appName ("Spark app") .getOrCreate () spark.conf.set("spark.sql.legacy.allowUntypedScalaUDF", "true") def simpleFn(foo: Any): Seq[Any] = List("hello world", "Another String", 42L) // def simpleFn(foo: Any): Seq[Any] = List("hello world", "Another String") def simpleUdf = F.udf( simpleFn(_), dataType = T.StructType( List( T.StructField("a_string", T.StringType), T.StructField("another_string", T.StringType), T.StructField("an_int", T.IntegerType), ) ) ) Seq(("bar", "foo")) .toDF("column", "input") .withColumn( "array_data", simpleUdf($"input") ) .show(truncate=false)
运行抛出错误:
IllegalArgumentException: The value (List(Another String, 42)) of the type (scala.collection.immutable.$colon$colon) cannot be converted to the string type
奇怪的是,当struct仅含一个字段时可正常运行:
def simpleFn(foo: Any): Seq[Any] = List("hello world") def simpleUdf = F.udf( simpleFn(_), dataType = T.StructType( List( T.StructField("a_string", T.StringType), ) ) )
查询结果:
+------+-----+-------------+ |column|input|array_data | +------+-----+-------------+ |bar |foo |{hello world}| +------+-----+-------------+
看起来Sequence的第一个元素被当作struct第一个字段,剩余元素被当作第二个字段,第三个字段为null从而引发异常。请问是否有人遇到过这种动态构建Schema的UDF问题?
使用环境:Spark 3.3.1、Scala 2.12、DBR 12.0
反射困境
一种繁琐的实现方式是根据推断Schema生成Scala代码实现case class,编译打包为JAR后加载到Databricks Runtime,再将case class作为UDF返回值,但过于复杂。希望直接生成case class并按如下方式使用:
def myUdf[CaseClass](input: SomeInputType): CaseClass = CaseClass(input.giveMeResults: _*)
但无法将通过eval创建的类型纳入当前上下文。运行以下代码:
import scala.reflect.runtime.universe import scala.tools.reflect.ToolBox val tb = universe.runtimeMirror(getClass.getClassLoader).mkToolBox() val test = tb.eval(tb.parse("object Test; Test"))
得到结果:
... test: Any = __wrapper$1$bb89c0cde37c48929fa9d8cdabeeb0f8.__wrapper$1$bb89c0cde37c48929fa9d8cdabeeb0f8$Test$1$@492531c0
test应为Test的实例,但REPL类型系统无法识别Test类型,无法使用test.asInstanceOf[Test]。
内容的提问来源于stack exchange,提问作者user961826

