andygrove opened a new pull request, #6277: URL: https://github.com/apache/datafusion-comet/pull/6277
Backport of #6219 to `branch-1.0`. Cherry-picked from `646ff181a1d52bb6a918b88f7dd0657d9f5d546c`. The production change matches upstream except for one argument, and the tests needed a new home; see "What changes are included" below. ## Which issue does this PR close? Closes #6091 on `branch-1.0`. Listed in #6201, as #6092, the earlier PR for the same fix that #6219 picked up. ## Rationale for this change The bug ships in 1.0.0. `executePlan`'s path for plans with a JVM-fed input is the same loop on `branch-1.0` as on `main` before #6219. It polls the plan's stream, and on `Pending` pulls from the JVM iterators and polls again. Once every JVM-fed `ScanExec` or `ShuffleScanExec` holds a batch or has reached EOF, that pull does nothing. Both return `Pending` without registering a waker. So a plan that is also waiting on native I/O spins the Spark task thread at 100% CPU for as long as the read takes. The common shape is a broadcast hash join, whose build side is a JVM-fed `ScanExec`, over a native Parquet or Iceberg scan on S3, HDFS or ABFS. Broadcast joins and native scans are both on by default in 1.0. #6091 reports 87.3 core-hours against 6.6 for Spark on the same stages of a production job, with stages 3 to 6 times slower. ## What changes are included in this PR? The fix is the original one; see #6219 for the details. `ScanStream` and `ShuffleScanStream` register the poll's waker when their buffer is empty, and a refill wakes it. `next_batch` parks the task until a waker fires instead of polling again at once, using a flag of its own so a nested `block_on` on the same thread cannot take its wake-up. The adaptations: - `ShuffleScanExec::get_next_batch` calls `get_next` without a `requires_validation` argument. On `main` that argument comes from the Celeborn shuffle reader (#5531), which is not on `branch-1.0`. The rest of `shuffle_scan.rs` matches upstream. - `scan.rs` imports `AtomicWaker` into `branch-1.0`'s imports. The `decode_string_arrays` import next to it on `main` comes from #5310. - `jni_api.rs` has no test module on `branch-1.0`. This PR adds one holding only #6219's three `next_batch` tests and their two helpers. On `main` they were appended to an existing module whose other tests cover code that is not on `branch-1.0`. The nested case matters on `branch-1.0` too. A JNI pull there can run another Comet plan on the same thread, and every `block_on` on a thread shares one parker. `next_batch_polls_again_when_a_nested_block_on_took_the_wake_up` covers it. As on `main`, the park has no timeout, so a lost wake-up would hang a task rather than spin it. ## How are these changes tested? Run locally on `branch-1.0` with the default Spark 4.1 profile and JDK 17: - The four new tests pass: the three `next_batch` tests in `jni_api.rs` and `refill_wakes_the_pending_poll_and_eof_stays_buffered` in `shuffle_scan.rs`. The `next_batch` tests drive the function this PR introduces, so they cannot run against the old loop. On `main`, the busy-poll test's closure ran 61,139 times in one 50 ms wait before the fix. - All `datafusion-comet` lib tests pass: 186 passed, 4 ignored. - `CometJoinSuite`, `CometExecSuite`, `CometNativeShuffleSuite` and `CometIcebergNativeSuite` pass against the new native library, 299 tests, and none of them hangs. - `cargo fmt --all -- --check` and `cargo clippy --all-targets --workspace -- -D warnings` pass. I added `run-iceberg-tests`, which #6219 did not have, so CI also runs the Iceberg Spark suites against the new loop. On `branch-1.0` the Spark 3.5 and 4.1 SQL suites run on pull requests without a label. ## Are there any user-facing changes? A task whose native plan has a JVM-fed input no longer holds a core at 100% while its native scans wait on storage. There are no config or API changes. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
