Skip to content

feat: support AtLeastNNonNulls natively - #6180

Draft
rich7420 wants to merge 9 commits into
apache:mainfrom
rich7420:feat/6093-at-least-n-non-nulls
Draft

rich7420 wants to merge 9 commits into
apache:mainfrom
rich7420:feat/6093-at-least-n-non-nulls

Conversation

@rich7420

@rich7420 rich7420 commented Sep 24, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6093.

Rationale for this change

Spark uses AtLeastNNonNulls for DataFrame.na.drop. Native evaluation avoids falling back to Spark while preserving the count of non-NULL, non-NaN values and Spark's per-row short-circuit and error behavior.

What changes are included in this PR?

  • Add the Scala serde, protobuf representation, Rust planner wiring and native expression.
  • Use bitmap OR/AND for the any/all-valid cases, and row or bitmap counters for general thresholds.
  • Expose spark.comet.exec.atLeastNNonNulls.smallBatchThreshold, default 64. It selects the general counting strategy by input batch row count and accepts any positive integer. It changes neither the fixed 64-bit word size nor the na.drop threshold.
  • Cover logical dictionary NULLs, floating-point NaNs including signed payloads, sliced buffers, counter boundaries, scalar inputs, outer nullability and per-row child errors.
  • Add configuration/planner tests, tuning documentation, Criterion crossover controls and whole-query benchmarks comparing native execution, Comet fallback and Spark.

Shared CASE/IF evaluation (#6350), filter-output pruning (#6396), and independent IF semantic tests (#6397) have merged. The branch uses these from main. Generic projection code is no longer part of this PR's diff. The projection result regression is in operators/filter_projection.sql; its Scala test retains native output-row metrics checks.

How are these changes tested?

Current published head: 1f0c6478ecde67d31463bb23249f334bc26fa35c. Current-head upstream CI completed successfully: 25 checks passed, with 15 conditional jobs skipped.

  • Native CI passed all eight at_least_n_non_nulls unit tests and the planner threshold round-trip test. The workspace run completed with 2,228 passing tests and 5 skipped tests.
  • Spark 4.1 / JDK 17 expressions CI passed all four AtLeastNNonNulls JVM semantic tests and the query-serde threshold test. The job completed with 2,232 passing tests, 1 cancellation and 12 ignored tests, with no failures or aborted suites.
  • The default exec, shuffle and scan suites, strict Spark 3.5 warnings, supported-profile Java lint checks, and TPC-H/TPC-DS result verification passed.

Historical combined-candidate tests and timings remain in fork #42 at 4528a2a67cff87c06179337ce4e1b35de2355839 and its completed CI. They do not establish current-head performance. The original issue author's workload remains unverified.

Keep this PR Draft pending the remaining Ready requirements: relevant Spark SQL/other-profile CI and current-head benchmark controls. Separate count-only output pruning from full-column conditional-expression costs, compare both counting strategies on identical inputs, and report repeated-run variation and slower cases. Default CI establishes the native/JVM test verdict; it does not provide those broader suite or performance results.

@github-actions github-actions Bot added enhancement New feature or request area:expressions Expression evaluation labels Sep 24, 2026
@rich7420
rich7420 marked this pull request as draft September 24, 2026 12:57
@rich7420
rich7420 force-pushed the feat/6093-at-least-n-non-nulls branch from 1bf4f08 to 6ad5c0d Compare September 28, 2026 14:50
@rich7420

rich7420 commented Sep 28, 2026 •

Copy link
Copy Markdown
Contributor Author

Status updated 2026-10-07: #6350, #6396 and #6397 are merged. Feature refresh c34480068 onto main c8e6553da is now published and mergeable; its CI is running. See the current delivery plan. The results below belong to historical candidate 4528a2a67, not the new feature head.

Historical integration progress and validation

#6350 has merged. #6396 remains open with passing PR CI and is the remaining implementation prerequisite. #6397 contributes independent IF semantic tests and does not block this feature.

The refreshed integration candidate 4528a2a67 is published against actual main 62ed9d16d. It preserves the previous fork candidate through a normal merge. CASE/IF and float normalization now match actual main, including the UTC fix and regression. Comparing against a prerequisite-only tree still leaves 13 AtLeastNNonNulls files.

Local validation of this exact tree passed: release native build; 24 Spark 4.1.3 tests; six Spark 3.5 tests with strict warnings (both on JDK 21); 16 conditional, seven AtLeastNNonNulls and 50 planner Rust tests; workspace/all-target release Clippy, cargo fmt, Spotless and RAT. Neither Spark run had failures, cancellations or skipped tests.

Two independent processes passed all full-query answers and native/dispatcher-path assertions on the M3 Pro/Spark 4.1.3 workload: 1,048,576 Parquet rows, 32 Double columns, no NULLs, na.drop(16), count and all 32 cleaned sums. Twelve rounds alternated mode order; medians exclude the first three. The second process reversed initial mode and NaN-rate order. At 5% NaN, native/fallback medians were 329.92/333.53 ms and 319.96/360.13 ms (1.1% and 11.2% lower native median time). At 0% NaN, they were 290.71/319.72 ms and 302.48/329.58 ms. The permanent focus-double benchmark also passed; its 5% NaN native/fallback averages were 325/354 ms, with standard deviations of 11/37 ms. Variance is material, so this is workload-specific combined-candidate evidence, not a fixed speedup. The original issue workload remains unverified.

The refreshed candidate is published in fork draft #42. Current-head fork CI passed at 4528a2a67: Required Checks, all Comet suites, Rust, strict Spark 3.5/JDK17 compilation, benchmark checks, TPC-H/TPC-DS and all nine Spark 4.1 SQL shards. The earlier same-head run cancelled during this push is superseded by the completed successful run. Other Spark runtime profiles, Iceberg and macOS suites were skipped by this run. This is the combined candidate's verdict; the final narrowed upstream head still needs independent verification after #6396 merges.

Upstream #6180 remains Draft at 22c04aaba. After #6396 merges, I will update this branch from actual main, apply the review follow-ups, narrow the feature diff and repeat validation on its resulting head. The delivery plan is also recorded in the original issue.

}
}

test("filter output projection preserves rows, aliases and metrics") {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's see if you could port this test to Comet SQL tests: https://datafusion.apache.org/comet/contributor-guide/sql-file-tests.html

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Moved the projection result regression to spark/src/test/resources/sql-tests/operators/filter_projection.sql in #6396, which is now merged. It covers count-only output, aliases/order/duplicates, full projections, empty results and a rejected row that must not reach an ANSI cast. The Scala test retains the native output-row metrics checks, which the SQL fixture cannot assert. These generic projection changes are no longer in this PR’s diff.

})
.collect();

let mut exprs = ProjectionExprs::from(exprs?);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it possible to fix this issue in a separate PR?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Split out as #6396 and merged. This branch now takes the projection implementation from main, so this PR’s diff is limited to AtLeastNNonNulls and its configuration, tests, benchmarks and documentation. The shared CASE/IF implementation (#6350) and independent IF tests (#6397) are also merged.

result, None,
))));
}
if batch.num_rows() < 64 {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why 64?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

64 was a one-word starting point for the general counting path: bitmap counters process 64 rows per word, while small batches may not amortize their bookkeeping. It is a heuristic, not a required or universal crossover. I made it configurable as spark.comet.exec.atLeastNNonNulls.smallBatchThreshold, default 64; any positive value is accepted, independently of the bitmap word size and the na.drop threshold. Tests cover both strategies and non-word-aligned thresholds. Criterion now compares row, bitmap and default strategies on identical inputs around the boundary. Current-head timings are still pending, so I am not claiming that 64 is optimal.

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

Labels

area:expressions Expression evaluation enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support Spark's AtLeastNNonNulls natively

2 participants