Skip to content

[SPARK-59818][SQL] Handle a NULL sketch state in approx_top_k_estimate and approx_top_k_combine - #59094

Open
jubins wants to merge 3 commits into
apache:masterfrom
jubins:j-spark-59818-fix-null-sketch-state
Open

jubins wants to merge 3 commits into
apache:masterfrom
jubins:j-spark-59818-fix-null-sketch-state

Conversation

@jubins

@jubins jubins commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

Fixes SPARK-59818 — makes approx_top_k_estimate and approx_top_k_combine handle a NULL sketch state instead of throwing NullPointerException. Two related gaps are addressed:

approx_top_k_estimate dereferenced the state unconditionally. ApproxTopKEstimate overrides BinaryExpression.eval directly rather than implementing nullSafeEval, so the inherited null handling is bypassed, and eval calls stateEval.asInstanceOf[InternalRow].getBinary(0) with no null check. It also declared override def nullable: Boolean = false even though the result is NULL whenever the state is, so the reported output schema was wrong and the optimizer reasoned incorrectly about the expression.

approx_top_k_combine had the same unguarded dereference in update. ApproxTopKCombine.update calls inputState.getBinary(0) directly, so a NULL sketch fails on the executors and aborts the running job, rather than being ignored the way aggregates normally ignore NULL inputs.

Analysis does not prevent either case: checkStateFieldAndType validates only the schema of the state struct, not its nullness, so a correctly-typed nullable struct column passes the analyzer and then fails at optimization or execution time. A NULL sketch arises naturally from an outer join, a UNION with a missing side, or a partially-populated column.

Both are surfaced and locked down by new tests.

Brief change log

  • sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/ApproxTopKExpressions.scala: ApproxTopKEstimate.eval returns null when the state evaluates to null, and nullable is now state.nullable instead of a hardcoded false, so the reported schema matches the values actually produced.
  • sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/ApproxTopKAggregates.scala: ApproxTopKCombine.update returns the buffer untouched for a NULL sketch, so combine ignores NULL inputs like any other aggregate.
  • sql/core/src/test/scala/org/apache/spark/sql/ApproxTopKSuite.scala: added four tests covering the foldable NULL state, the non-foldable NULL state, the reported nullability, and combine skipping NULL sketches.

Why are the changes needed?

Both paths fail on a valid query.

approx_top_k_estimate raises the NPE inside ConstantFolding, which reports it to the user as an internal error:

SELECT approx_top_k_estimate(s, 5) FROM (
  SELECT CAST(NULL AS STRUCT<sketch: BINARY, maxItemsTracked: INT,
    itemDataType: INT, itemDataTypeDDL: STRING>) AS s
);
[INTERNAL_ERROR] The Spark SQL phase optimization failed with an internal error.
You hit a bug in Spark or the Spark plugins you use. SQLSTATE: XX000

Caused by: java.lang.NullPointerException: Cannot invoke
  "org.apache.spark.sql.catalyst.InternalRow.getBinary(int)" because "stateEval" is null
  at ApproxTopKEstimate.eval(ApproxTopKExpressions.scala:113)
  at org.apache.spark.sql.catalyst.optimizer.ConstantFolding$.tryFold(expressions.scala:65)

approx_top_k_combine fails on the executors and aborts the job:

SELECT approx_top_k_estimate(approx_top_k_combine(s), 5) FROM (
  SELECT CAST(NULL AS STRUCT<sketch: BINARY, maxItemsTracked: INT,
    itemDataType: INT, itemDataTypeDDL: STRING>) AS s
);
org.apache.spark.SparkException: Job aborted due to stage failure:
  Task 0 in stage 0.0 failed 1 times ...
Caused by: java.lang.NullPointerException
  at ApproxTopKCombine.update(ApproxTopKAggregates.scala:908)
  at TypedImperativeAggregate.update(interfaces.scala:583)

The new behaviour follows ordinary SQL semantics for a nullable input column: approx_top_k_estimate(NULL, k) returns NULL, and approx_top_k_combine skips NULL sketches.

This is distinct from the existing APPROX_TOP_K_NULL_ARG error, which is raised for NULL foldable configuration arguments (k, maxItemsTracked). That check, ApproxTopK.checkExpressionNotNull, calls expr.eval() with no input row, so it applies only to literals and is not applicable to a per-row data column.

It is also distinct from earlier NULL-related work on these functions. SPARK-53960 added handling for NULLs within the accumulated data (counting them as items) and did not touch the NULL-state path or the nullable declaration; SPARK-58095 concerns empty input (zero rows) in combine's eval/serialize, which is a different code path and trigger.

Does this PR introduce any user-facing change?

Yes, in three ways, all aligning the functions with ordinary SQL null semantics:

  1. approx_top_k_estimate on a NULL sketch state returns NULL instead of failing the query with [INTERNAL_ERROR].
  2. approx_top_k_combine skips NULL sketches instead of aborting the job with a NullPointerException.
  3. approx_top_k_estimate now reports itself as nullable when its state argument is nullable. Previously it always reported nullable = false, which was incorrect for a nullable state. The reported schema is unchanged for a non-nullable state, which is the common case — the sql-expression-schema.md golden file is unaffected, as its example applies the function to an aggregate result.

These are changes compared to released versions: approx_top_k is annotated since = "4.1.0".

How was this patch tested?

Four new tests in ApproxTopKSuite, covering both functions:

  • SPARK-59818: estimate of a foldable NULL state returns NULL
  • SPARK-59818: estimate of a non-foldable NULL state returns NULL
  • SPARK-59818: estimate is nullable when its state is nullable
  • SPARK-59818: combine skips NULL sketches

The foldable and non-foldable state are covered separately because only the former is reached by ConstantFolding; the non-foldable case is built with a UNION ALL so the state stays a per-row column.

All four tests were confirmed to fail without the fix. Reverting the two source commits and re-running them produces exactly the reported failures — Cannot invoke "InternalRow.getBinary(int)" because "stateEval" is null for the three estimate tests and ... because "inputState" is null, aborting the stage, for the combine test — and all four pass with the fix applied.

The full ApproxTopKSuite passes (200 tests: 196 existing plus the 4 new). ExpressionInfoSuite, NullabilitySuite, CanonicalizeSuite, and ExpressionsSchemaSuite also pass, confirming the nullable change leaves the expression schema golden file unchanged.

Run:

build/sbt 'sql/testOnly org.apache.spark.sql.ApproxTopKSuite'
build/sbt 'catalyst/testOnly *ExpressionInfoSuite *NullabilitySuite *CanonicalizeSuite'
build/sbt 'sql/testOnly *ExpressionsSchemaSuite'

Was this patch authored or co-authored using generative AI tooling?

  • Yes - Claude Code was used as a pair-programming assistant. All changes were reviewed and verified by the author. Generated-by: Claude Opus 5

…h 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.
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.
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.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant