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

Spark 3.3转UDAF为TypedColumn调用报NotSerializableException

Spark UDAF强类型调用报TypedColumn序列化异常问题

问题复现

复现Spark SQL用户自定义聚合函数(UDAF)官方示例时,仅调整了测试DataFrame的创建逻辑,核心实现代码如下:

import org.apache.spark.sql.{Encoder, Encoders, SparkSession}
import org.apache.spark.sql.expressions.Aggregator

case class Employee(name: String, salary: Long)
case class Average(var sum: Long, var count: Long)

object MyAverage extends Aggregator[Employee, Average, Double] {
  // 聚合初始零值,满足任意值b + zero = b的规则
  def zero: Average = Average(0L, 0L)
  // 分区内聚合逻辑,直接修改buffer提升性能
  def reduce(buffer: Average, employee: Employee): Average = {
    buffer.sum += employee.salary
    buffer.count += 1
    buffer
  }
  // 分区间聚合合并逻辑
  def merge(b1: Average, b2: Average): Average = {
    b1.sum += b2.sum
    b1.count += b2.count
    b1
  }
  // 最终结果计算逻辑
  def finish(reduction: Average): Double = reduction.sum.toDouble / reduction.count
  // 中间缓冲类型编码器
  def bufferEncoder: Encoder[Average] = Encoders.product
  // 最终输出类型编码器
  def outputEncoder: Encoder[Double] = Encoders.scalaDouble
}

val originalDF = Seq(
     ("Michael", 3000),
     ("Andy", 4500),
     ("Justin", 3500),
     ("Berta", 4000)
   ).toDF("name", "salary")

测试DataFrame内容如下:

+-------+------+
|name   |salary|
+-------+------+
|Michael|3000  |
|Andy   |4500  |
|Justin |3500  |
|Berta  |4000  |
+-------+------+

两种调用方式的执行结果存在差异:

  • 采用SQL注册方式调用可正常运行,代码如下:
    spark.udf.register("myAverage", functions.udaf(MyAverage))
    originalDF.createOrReplaceTempView("employees")
    val result = spark.sql("SELECT myAverage(salary) as average_salary FROM employees")
    result.show()
    
    执行后返回预期平均薪资3750.0,无任何报错。
  • 采用TypedColumn强类型方式调用时作业失败,代码如下:
    val averageSalary = MyAverage.toColumn.name("average_salary")
    val result = originalDF.as[Employee].select(averageSalary)
    result.show()
    
    抛出NotSerializableException: org.apache.spark.sql.TypedColumn异常,提示TypedColumn实例无法序列化。当前运行环境为DBR 11.0、Spark 3.3.0、Scala 2.12。

问题原因

该问题是Spark 3.3.0版本的已知bug:Aggregator.toColumn方法生成的TypedColumn实例未做序列化适配,当实例作为独立变量在Driver端定义、需要序列化分发到Executor执行时,就会触发序列化异常。SQL注册调用路径因为走的是UDF注册的独立序列化逻辑,不会触发该问题。DBR 11.0版本绑定的是原生Spark 3.3.0,没有对该bug做补丁修复。

修复方案

可任选以下一种方案解决:

  • 内联TypedColumn定义,不要单独将其赋值给外层变量持有,避免跨作用域序列化触发问题:
    // 直接在select方法内生成TypedColumn实例
    val result = originalDF.as[Employee].select(MyAverage.toColumn.name("average_salary"))
    result.show()
    
  • 调整Spark配置,开启Typed UDF序列化兼容:在SparkSession初始化配置中添加spark.sql.typed.udf.trySerializable = true,强制Spark对生成的TypedColumn做序列化包装。
  • 升级运行环境:将Spark升级到3.3.1及以上版本,DBR环境对应升级到11.3 LTS及以上版本,官方已在后续小版本中修复了该序列化缺陷。

内容的提问来源于stack exchange,提问作者Ignacio Alorre

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.29 15:30:05