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
相关产品推荐
相关产品推荐

