Spark SQL含Join、OrderBy与UUID()的查询触发NullPointerException问题排查
Spark SQL 含Join、OrderBy与UUID()查询执行失败问题
问题详情
尝试运行一条包含Join、OrderBy子句,且最外层SELECT使用UUID()函数的Spark SQL查询,但执行失败。
代码示例
val query = spark.sql(" select name, uuid() as _iid from (select s.name from titanic s join titanic t on s.name = t.name order by name) ;")
触发异常
Exception in thread "main" java.lang.NullPointerException at org.apache.spark.sql.catalyst.expressions.GeneratedClass$SpecificUnsafeProjection.apply(Unknown Source) at org.apache.spark.sql.execution.TakeOrderedAndProjectExec.$anonfun$executeCollect$2(limit.scala:207) at scala.collection.TraversableLike.$anonfun$map$1(TraversableLike.scala:237) at scala.collection.IndexedSeqOptimized.foreach(IndexedSeqOptimized.scala:36) at scala.collection.IndexedSeqOptimized.foreach$(IndexedSeqOptimized.scala:33) at scala.collection.mutable.ArrayOps$ofRef.foreach(ArrayOps.scala:198) at scala.collection.TraversableLike.map(TraversableLike.scala:237) at scala.collection.TraversableLike.map$(TraversableLike.scala:230) at scala.collection.mutable.ArrayOps$ofRef.map(ArrayOps.scala:198) at org.apache.spark.sql.execution.TakeOrderedAndProjectExec.executeCollect(limit.scala:207) at org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.$anonfun$executeCollect$1(AdaptiveSparkPlanExec.scala:338) at org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.withFinalPlanUpdate(AdaptiveSparkPlanExec.scala:366) at org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec.executeCollect(AdaptiveSparkPlanExec.scala:338) at org.apache.spark.sql.Dataset.collectFromPlan(Dataset.scala:3715) at org.apache.spark.sql.Dataset.$anonfun$head$1(Dataset.scala:2728) at org.apache.spark.sql.Dataset.$anonfun$withAction$1(Dataset.scala:3706) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$5(SQLExecution.scala:103) at org.apache.spark.sql.execution.SQLExecution$.withSQLConfPropagated(SQLExecution.scala:163) at org.apache.spark.sql.execution.SQLExecution$.$anonfun$withNewExecutionId$1(SQLExecution.scala:90) at org.apache.spark.sql.SparkSession.withActive(SparkSession.scala:775) at org.apache.spark.sql.execution.SQLExecution$.withNewExecutionId(SQLExecution.scala:64) at org.apache.spark.sql.Dataset.withAction(Dataset.scala:3704) at org.apache.spark.sql.Dataset.head(Dataset.scala:2728) at org.apache.spark.sql.Dataset.take(Dataset.scala:2935) at org.apache.spark.sql.Dataset.getRows(Dataset.scala:287) at org.apache.spark.sql.Dataset.showString(Dataset.scala:326) at org.apache.spark.sql.Dataset.show(Dataset.scala:808) at org.apache.spark.sql.Dataset.show(Dataset.scala:785) at hyperspace2.sparkPlan$.delayedEndpoint$hyperspace2$sparkPlan$1(sparkPlan.scala:14) at hyperspace2.sparkPlan$delayedInit$body.apply(sparkPlan.scala:6) at scala.Function0.apply$mcV$sp(Function0.scala:39) at scala.Function0.apply$mcV$sp$(Function0.scala:39) at scala.runtime.AbstractFunction0.apply$mcV$sp(AbstractFunction0.scala:17) at scala.App.$anonfun$main$1$adapted(App.scala:80) at scala.collection.immutable.List.foreach(List.scala:392) at scala.App.main(App.scala:80) at scala.App.main$(App.scala:78) at hyperspace2.sparkPlan$.main(sparkPlan.scala:6) at hyperspace2.sparkPlan.main(sparkPlan.scala)
关键现象
- 移除OrderBy子句后,查询可正常输出结果;
- 将
UUID()替换为monotonically_increasing_id()也能正常运行,但不符合需求。
错误原因
这是Spark执行计划优化的已知问题:
当子查询包含OrderBy时,Spark会将外层的UUID()函数下推到排序后的执行阶段,但UUID()属于非确定性函数(每次调用返回不同值)。Spark在TakeOrderedAndProjectExec算子处理排序结果时,会尝试复用已生成的投影对象,此时非确定性函数的状态未正确初始化,最终触发空指针异常。
而monotonically_increasing_id()是确定性函数(基于分区和行号生成固定值),不会触发该问题;移除OrderBy后,执行计划不会进入TakeOrderedAndProjectExec的特殊处理逻辑,因此也不会报错。
解决方法
推荐以下两种可靠修复方式:
1. 给子查询添加别名并设置大LIMIT
强制Spark将子查询作为独立逻辑单元执行,避免函数错误下推:
val query = spark.sql(" select name, uuid() as _iid from (select s.name from titanic s join titanic t on s.name = t.name order by name) tmp limit 1000000000;")
2. 将UUID()封装为自定义UDF
自定义非确定性UDF包装UUID(),让Spark正确识别其特性,规避错误优化:
import org.apache.spark.sql.functions.udf import java.util.UUID val uuidUdf = udf(() => UUID.randomUUID().toString).asNonNullable() spark.udf.register("uuid_udf", uuidUdf) val query = spark.sql(" select name, uuid_udf() as _iid from (select s.name from titanic s join titanic t on s.name = t.name order by name) ;")
3. 临时关闭自适应执行计划(不推荐)
通过关闭自适应执行计划规避问题,但会影响整体查询性能:
spark.conf.set("spark.sql.adaptive.enabled", "false")
内容的提问来源于stack exchange,提问作者Chhavi Bansal
相关产品推荐
相关产品推荐

