From cc9c4806e919f40f94b9828cbad0e175c3bdf08f Mon Sep 17 00:00:00 2001 From: Jubin Soni Date: Sun, 27 Sep 2026 23:23:52 -0700 Subject: [PATCH 1/3] [SPARK-59818] Return NULL from approx_top_k_estimate for a NULL sketch state ApproxTopKEstimate overrides BinaryExpression.eval directly, bypassing the inherited null handling, and dereferences the evaluated state unconditionally. A NULL state therefore raises a NullPointerException, which surfaces to users as an INTERNAL_ERROR when ConstantFolding evaluates the expression. Return NULL when the state is NULL, and declare the expression nullable whenever its state argument is nullable so the reported output schema matches the values actually produced. --- .../sql/catalyst/expressions/ApproxTopKExpressions.scala | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/ApproxTopKExpressions.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/ApproxTopKExpressions.scala index fc647c8602086..20257002647a9 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/ApproxTopKExpressions.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/ApproxTopKExpressions.scala @@ -109,6 +109,9 @@ case class ApproxTopKEstimate(state: Expression, k: Expression) ApproxTopK.checkExpressionNotNull(k, "k") // eval val stateEval = left.eval(input) + if (stateEval == null) { + return null + } val kEval = right.eval(input) val dataSketchBytes = stateEval.asInstanceOf[InternalRow].getBinary(0) val maxItemsTrackedVal = stateEval.asInstanceOf[InternalRow].getInt(1) @@ -127,7 +130,9 @@ case class ApproxTopKEstimate(state: Expression, k: Expression) override protected def withNewChildrenInternal(newState: Expression, newK: Expression) : Expression = copy(state = newState, k = newK) - override def nullable: Boolean = false + // The sketch state is an ordinary nullable input column: `approx_top_k_estimate(NULL, k)` + // returns NULL rather than failing, so the result is nullable whenever the state is. + override def nullable: Boolean = state.nullable override def prettyName: String = getTagValue(FunctionRegistry.FUNC_ALIAS).getOrElse("approx_top_k_estimate") From 0e27a0136a75fbd427d30c555b45d919d5053afb Mon Sep 17 00:00:00 2001 From: Jubin Soni Date: Sun, 27 Sep 2026 23:24:00 -0700 Subject: [PATCH 2/3] [SPARK-59818] Skip NULL sketches in approx_top_k_combine ApproxTopKCombine.update dereferences the evaluated state unconditionally, so a NULL sketch raises a NullPointerException. Unlike the sibling issue in approx_top_k_estimate this fails on the executors and aborts the running job. Leave the buffer untouched for a NULL sketch, so combine ignores NULL inputs the way any other aggregate does. --- .../catalyst/expressions/aggregate/ApproxTopKAggregates.scala | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/ApproxTopKAggregates.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/ApproxTopKAggregates.scala index 8ed6590793daa..df07778f8f822 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/ApproxTopKAggregates.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/ApproxTopKAggregates.scala @@ -905,6 +905,10 @@ case class ApproxTopKCombine( */ override def update(buffer: CombineInternal[Any], input: InternalRow): CombineInternal[Any] = { val inputState = state.eval(input).asInstanceOf[InternalRow] + if (inputState == null) { + // A NULL sketch contributes nothing, like NULL inputs to any other aggregate. + return buffer + } val inputSketchBytes = inputState.getBinary(0) val inputMaxItemsTracked = inputState.getInt(1) val inputItemDataTypeDDL = inputState.getUTF8String(3).toString From 5c91274a32e2f3ab13b4f6350976d5404a97d8d3 Mon Sep 17 00:00:00 2001 From: Jubin Soni Date: Sun, 27 Sep 2026 23:26:50 -0700 Subject: [PATCH 3/3] [SPARK-59818] Add tests for a NULL sketch state Cover both the foldable and the non-foldable state, since only the former is reached by ConstantFolding, and assert the reported nullability along with the returned values. All four tests fail without the preceding two commits. --- .../apache/spark/sql/ApproxTopKSuite.scala | 45 +++++++++++++++++++ 1 file changed, 45 insertions(+) diff --git a/sql/core/src/test/scala/org/apache/spark/sql/ApproxTopKSuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/ApproxTopKSuite.scala index dbccbcce73933..10b3e7912248c 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/ApproxTopKSuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/ApproxTopKSuite.scala @@ -549,6 +549,37 @@ class ApproxTopKSuite extends SharedSparkSession { ) } + + private val nullSketchState = + """CAST(NULL AS STRUCT)""".stripMargin + + test("SPARK-59818: estimate of a foldable NULL state returns NULL") { + checkAnswer(sql(s"SELECT approx_top_k_estimate($nullSketchState, 5)"), Row(null)) + } + + test("SPARK-59818: estimate of a non-foldable NULL state returns NULL") { + withTempView("estimate_null_state") { + sql( + s"""SELECT approx_top_k_accumulate(expr) AS state + |FROM VALUES 0, 1, 1 AS tab(expr) + |UNION ALL + |SELECT $nullSketchState AS state""".stripMargin) + .createOrReplaceTempView("estimate_null_state") + val res = sql("SELECT approx_top_k_estimate(state, 2) FROM estimate_null_state") + checkAnswer(res, Seq(Row(Seq(Row(1, 2), Row(0, 1))), Row(null))) + } + } + + test("SPARK-59818: estimate is nullable when its state is nullable") { + withTempView("estimate_nullable_state") { + sql(s"SELECT $nullSketchState AS state") + .createOrReplaceTempView("estimate_nullable_state") + val res = sql("SELECT approx_top_k_estimate(state, 2) FROM estimate_nullable_state") + assert(res.schema.fields.head.nullable) + } + } + ///////////////////////////////// // approx_top_k_combine ///////////////////////////////// @@ -1296,4 +1327,18 @@ class ApproxTopKSuite extends SharedSparkSession { checkAnswer(est, Row(Seq(Row(null, 5)))) } } + + test("SPARK-59818: combine skips NULL sketches") { + withTempView("combine_null_state") { + sql( + s"""SELECT approx_top_k_accumulate(expr) AS state + |FROM VALUES 0, 1, 1 AS tab(expr) + |UNION ALL + |SELECT $nullSketchState AS state""".stripMargin) + .createOrReplaceTempView("combine_null_state") + val res = sql( + "SELECT approx_top_k_estimate(approx_top_k_combine(state), 2) FROM combine_null_state") + checkAnswer(res, Row(Seq(Row(1, 2), Row(0, 1)))) + } + } }