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]

Reply via email to