apache / apache/datafusion

EPIC: fix nullability report for spark expression

Open
#19,144 0 comments 0 reactions 0 assignees View on GitHub
EPIC
Dominant language
Rust
Stars
9.3k
Forks
2.4k
Avg merge
3d 7h
Merged PRs (30d)
344

Description

almost all of spark expression do not report the nullability correctly and they should implement `return_field_from_args` rather than `return_type`.

in some spark expressions source code, the expression is marked as null intolerant which means that if any of the arguments are null, than everything is null
In previous spark versions it was a trait, in later versions it became a field, **this does not mean it's nullable depend on the children**.

here `bitwise_not` marked as `nullIntolerant: true` which means if the child is null than the output is null as well
https://github.com/apache/spark/blob/98058da21e8a341eca10207f0ca458671220ca94/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/bitwiseExpressions.scala#L186

in others, the expression explicitly set nullability (make sure to see if it set null intolerant as well).

for example this `isNull` expression always marked as non nullable:
https://github.com/apache/spark/blob/98058da21e8a341eca10207f0ca458671220ca94/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/nullExpressions.scala#L411

for others who don't specify but extends from UnaryExpression, than the nullability by default is based on the child nullability:
https://github.com/apache/spark/blob/98058da21e8a341eca10207f0ca458671220ca94/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/expressions/Expression.scala#L590

Contributor guide

Open the contributing guide

Research direction

Start by comparing DataFusion's Spark expression implementations with the linked Spark sources: bitwiseExpressions.scala, nullExpressions.scala, and Expression.scala. Identify expressions whose nullability should follow their arguments or be explicitly non-null, then update them to use return_field_from_args rather than return_type. Done means the affected expressions report nullability consistently with the referenced Spark behavior.

Written by the indexing model from the issue text.

Assessment

Tech stack
rust, spark
Domain
data-engineering
Issue type
Bug
Difficulty
4/5
Estimated time
3-5 days
Activity status
Stale
Clarity
Mostly clear
Newbie friendliness
35/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.