Skip to content

fix: preserve Comet execution with AQE, subqueries and memory pressure - #1

Draft
notimesea wants to merge 2 commits into
mainfrom
feature/native-workaround-fixes
Draft

notimesea wants to merge 2 commits into
mainfrom
feature/native-workaround-fixes

Conversation

@notimesea

@notimesea notimesea commented Sep 25, 2026 •

Copy link
Copy Markdown

Which issue does this PR close?

No issue is closed automatically. This addresses the failures behind the feature overrides in datasci's SparkSessionProducer.

Rationale for this change

Warehouse jobs disable native sort, windows, codegen expressions, timestamp truncation and AQE coalescing to avoid failures. Those overrides cause Spark operator fallback and prevent jobs from benefiting fully from Comet. This patch fixes the reproduced failures while retaining Comet execution and the existing memory limits.

For example, a scalar subquery used as a regexp pattern must produce both 123 and 456 for a two-row batch; it must neither serialize an unresolved subquery nor read past its single-value Arrow argument. A whole-partition aggregate over a large, wide input must spill its rows instead of retaining the entire partition in memory.

What changes are included in this PR?

  • Defer native scan DPP partition enumeration until execution. AQE can inspect partitioning and coalesce a sibling shuffle without executing an unresolved adaptive broadcast placeholder. Execution also preserves zero partitions when file pruning removes every file.
  • Bind codegen scalar subqueries as runtime arguments. Generated Arrow getters and the null fast path broadcast length-one arguments without copying them. These expressions continue to use the JVM codegen bridge inside the Comet operator.
  • Canonicalize the Etc/UTC alias in native timestamp truncation so its output type matches native scan timestamps.
  • Cap the eager sort merge reservation at 1/32 of the configured per-task off-heap budget, retaining DataFusion's default when it is smaller. Memory-pool limits and allocation failures remain enforced.
  • Add spill support for whole-partition sum, avg, count, min and max windows using the existing native accumulators. Other window shapes keep their existing implementation.
  • Add focused regressions and document the memory behavior. Native concat_ws(array<string>, ...) already works in the current upstream base and needs no additional source change.

The branch is based on aaa33eb203422fa2c1015740e5d889825dc48985. It does not change datasci configuration or deploy a new JAR. Window spilling is scoped to the aggregate frames above; this is not general spill support for every unbounded window function. Production cost savings have not been measured.

How are these changes tested?

All execution is local in Docker using Colima, with Spark 3.5.3, Java 17 and Linux ARM64.

  • All five affected Scala suites passed on the updated patch: 240 tests passed, zero failed; one TIME test canceled because it requires Spark 4.1. This includes a new zero-file scan regression with AQE on and off, empty projection results and a zero count.
  • Native tests cover timezone aliases, sort reservation sizing, global/partitioned window results, nulls, ordering, payload preservation, three memory budgets, spill/no-spill, cancellation and reservation cleanup.
  • Clippy with -D warnings, Rustfmt and git diff --check passed.
  • An independent answer/operator verifier passed 13 local runs on the updated patch: six Spark/Comet pairs and a separate mixed Spark/native AQE+DPP case. The paired runs enable all six target features plus conversion from Spark plans.
  • Each repro container has 3 GiB total memory, a 1 GiB driver heap and 64 MiB Spark off-heap. Sort validates 200,000 rows, every 1 KiB payload and ordering within/across partitions. The native whole-partition window processes 2,000,000 rows with 1 KiB payloads, verifies totals, and spills approximately 1,977 MB without an OOM kill.

Spark SQL gate passed: dev/local-ci.sh spark 3.5 sql_core-1 on Spark 3.5.9 / Java 17 completed against 8a75822fe with 9,205 passed, 0 failed, 590 suites completed, 0 aborted; 11 canceled and 246 ignored. Total runtime was 24m13s, including preparation. Previously failing zero-file scan cases pass in the broad suite too. The local run used a 6-CPU / 7-GiB dev container and capped only the SBT coordinator heap at 1 GiB; SQL test selection and test-JVM settings were unchanged. No additional cgroup OOM occurred.

Canceled-test follow-up: six of the 11 cancellations came from missing pandas in the local test environment. After installing numpy 1.26.4, pandas 2.2.3 and pyarrow 12.0.1 within Spark 3.5.9's dependency constraints, those exact six tests were rerun with Comet: 6 passed, 0 failed, 0 canceled across QueryCompilationErrorsSuite, PythonUDTFSuite and PythonUDFSuite. This is a separate targeted run, in addition to the 9,205-test broad gate.

The other five cancellations are the ORC/Parquet/CSV/JSON/text variants of ignoreMissingFiles. They reach an unconditional assume(false) in the pre-existing dev/diffs/3.5.9.diff patch, before an assertion requiring SparkFileNotFoundException. The upstream comment concerns native Parquet, but the shared test loop cancels all five formats. This PR does not modify that patch, and those five exception checks have not been established as passing.

The PR remains draft as requested. Validation covers the reproduced failures and this Spark SQL matrix row, not every supported Spark version or production cost savings.

Defer native scan DPP partition enumeration until execution so AQE can
coalesce a sibling shuffle without executing an adaptive placeholder.

Pass codegen scalar subqueries as runtime inputs and broadcast length-one
Arrow arguments, including the null fast path. Normalize Etc/UTC in native
timestamp truncation to match scan timestamp types.

Scale sort's eager spill reserve with small off-heap task budgets. Spill
whole-partition aggregate window rows while retaining the existing native
accumulators and their Spark semantics.

Add regressions for planning, multi-row scalar inputs, timezone aliases,
reservation sizing, window spilling, ordering and cancellation cleanup.
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