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

Spark UDF返回运行时定义Schema的Row遇阻,求解决方案

ETL管道动态Schema构建瓶颈求助

背景

正在构建一个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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 15:20:23