[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
Open
[SPARK-59818][SQL] Handle a NULL sketch state in approx_top_k_estimate and approx_top_k_combine#59094jubins wants to merge 3 commits into
jubins wants to merge 3 commits into
Conversation
…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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this pull request?
Fixes SPARK-59818 — makes
approx_top_k_estimateandapprox_top_k_combinehandle a NULL sketch state instead of throwingNullPointerException. Two related gaps are addressed:approx_top_k_estimatedereferenced the state unconditionally.ApproxTopKEstimateoverridesBinaryExpression.evaldirectly rather than implementingnullSafeEval, so the inherited null handling is bypassed, andevalcallsstateEval.asInstanceOf[InternalRow].getBinary(0)with no null check. It also declaredoverride def nullable: Boolean = falseeven 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_combinehad the same unguarded dereference inupdate.ApproxTopKCombine.updatecallsinputState.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:
checkStateFieldAndTypevalidates 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, aUNIONwith 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.evalreturnsnullwhen the state evaluates tonull, andnullableis nowstate.nullableinstead of a hardcodedfalse, so the reported schema matches the values actually produced.sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/aggregate/ApproxTopKAggregates.scala:ApproxTopKCombine.updatereturns 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_estimateraises the NPE insideConstantFolding, which reports it to the user as an internal error:approx_top_k_combinefails on the executors and aborts the job:The new behaviour follows ordinary SQL semantics for a nullable input column:
approx_top_k_estimate(NULL, k)returns NULL, andapprox_top_k_combineskips NULL sketches.This is distinct from the existing
APPROX_TOP_K_NULL_ARGerror, which is raised for NULL foldable configuration arguments (k,maxItemsTracked). That check,ApproxTopK.checkExpressionNotNull, callsexpr.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
nullabledeclaration; SPARK-58095 concerns empty input (zero rows) in combine'seval/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:
approx_top_k_estimateon a NULL sketch state returns NULL instead of failing the query with[INTERNAL_ERROR].approx_top_k_combineskips NULL sketches instead of aborting the job with aNullPointerException.approx_top_k_estimatenow reports itself as nullable when its state argument is nullable. Previously it always reportednullable = false, which was incorrect for a nullable state. The reported schema is unchanged for a non-nullable state, which is the common case — thesql-expression-schema.mdgolden file is unaffected, as its example applies the function to an aggregate result.These are changes compared to released versions:
approx_top_kis annotatedsince = "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 NULLSPARK-59818: estimate of a non-foldable NULL state returns NULLSPARK-59818: estimate is nullable when its state is nullableSPARK-59818: combine skips NULL sketchesThe foldable and non-foldable state are covered separately because only the former is reached by
ConstantFolding; the non-foldable case is built with aUNION ALLso 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 nullfor 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
ApproxTopKSuitepasses (200 tests: 196 existing plus the 4 new).ExpressionInfoSuite,NullabilitySuite,CanonicalizeSuite, andExpressionsSchemaSuitealso pass, confirming thenullablechange leaves the expression schema golden file unchanged.Run:
Was this patch authored or co-authored using generative AI tooling?