Scala Spark窗口排名函数Task未序列化问题排查与优化求助
Spark单元测试触发Task Not Serializable异常排查及优化方案
生产环境运行代码时未出现Task not serializable问题,但单元测试执行以下代码时触发该异常,需排查原因差异,或提供更优序列化方案获取Hive表最新行。
val distinctBy = Window.partitionBy("id").orderBy(desc("updated_at")); val uniqueSellerDf = enrichedDf.withColumn("rank", rank().over(distinctBy))//row_number().over(distinctBy)) // rank和row_number函数均出现相同问题 .where($"rank" === 1).drop("rank") uniqueSellerDf.show() // 执行show动作时触发Task not serializable异常
异常原因排查
- 环境序列化上下文差异:生产环境Spark作业多运行在集群模式,序列化逻辑由集群统一管理;单元测试多为本地模式,若测试类未实现
Serializable,或引用了外部非序列化成员变量/闭包,Window算子序列化时会捕获到这些不可序列化对象。 - 关联对象序列化问题:
enrichedDf构建过程中若依赖了未实现Serializable的资源(如自定义UDF、配置对象),会连带触发整个算子链的序列化失败。 - 本地模式检查更严格:Spark本地模式在Driver端直接序列化任务,不像集群模式有容错机制,一些生产环境中被隐式处理的非序列化对象,会在单元测试中直接暴露问题。
解决方案及优化方案
1. 修复单元测试序列化问题
- 确保测试类实现
Serializable接口:
class YourTestClass extends FunSuite with Serializable { // 测试代码逻辑 }
- 检查
enrichedDf构建时引用的所有对象,确保它们都实现Serializable,或改用局部变量存储这些对象,避免引用测试类成员。 - 闭包中避免引用测试类成员:将依赖变量声明为局部变量,减少对测试类非序列化成员的引用。
2. 更优的获取最新行方案(替代Window函数)
如果Window函数的序列化问题难以快速定位,可改用以下两种轻量方式获取每个id的最新行:
方式一:groupBy + agg 结合first函数
按updated_at降序排序后,取每个分组的第一行:
import org.apache.spark.sql.functions.{first, desc} val uniqueSellerDf = enrichedDf .orderBy(desc("updated_at")) .groupBy("id") .agg( first("col1").alias("col1"), first("col2").alias("col2"), // 依次列出所有需要保留的列 first("updated_at").alias("updated_at") )
适合列数较少的场景,需显式声明所有要保留的列。
方式二:orderBy + dropDuplicates
先按updated_at降序排序,再根据id去重,Spark会保留排序后的第一行:
val uniqueSellerDf = enrichedDf .orderBy(desc("updated_at")) .dropDuplicates("id")
代码简洁,适合无需处理并列排名的场景(若存在相同updated_at的行,会随机保留其中一行)。
内容的提问来源于stack exchange,提问作者Deepak Singh
相关产品推荐
相关产品推荐

