andygrove opened a new pull request, #6219:
URL: https://github.com/apache/datafusion-comet/pull/6219

   ## Which issue does this PR close?
   
   Closes #6091.
   
   ## Rationale for this change
   
   This picks up #6092 by @mixermt, who diagnosed the busy-poll in #6091 and 
wrote the fix. Their three commits are here unchanged apart from a rebase onto 
main, and the last commit adds the one change still open from review there.
   
   When a plan has a JVM-fed input (a broadcast build side, 
`CometSparkRowToColumnar`, a shuffle read), `executePlan` polls the stream and, 
on `Pending`, pulls the next batches from the JVM. Once every JVM-fed scan 
holds a batch or has reached EOF, that pull is a no-op, so a stream that is 
also waiting on native I/O (a Parquet or Iceberg scan reading from S3 or HDFS) 
re-polls at full speed for the whole read and pins one core per task. #6091 has 
the numbers from a production workload.
   
   ## What changes are included in this PR?
   
   From #6092, by @mixermt: `ScanStream` and `ShuffleScanStream` register the 
poll's waker when their buffer is empty, `get_next_batch` wakes it after a 
refill, and EOF stays buffered so a drained reader isn't pulled again. The loop 
moves into `next_batch(stream, on_pending)`, which parks the `block_on` task 
until a waker fires instead of polling again at once, and the metrics update 
and the tracing memory sample run on the metrics interval instead of every 100 
polls. `development.md` describes the loop.
   
   Added here: `get_next_batch` reports again whether it made the JNI call, as 
in the first version of #6092, and `next_batch` parks only when no scan did. A 
JNI pull can run another Comet plan on the same thread, for example when a 
native writer's input is itself native and comes in through `CometArrowStream`. 
Every `block_on` on a thread shares tokio's thread-local parker, which holds a 
single wake-up. If the outer plan's I/O completes while the nested plan is 
parked, the nested park consumes that wake-up, and a park after the pull waits 
for a wake that has already fired. The park has no timeout, so that is a hang. 
I raised this in [review on 
#6092](https://github.com/apache/datafusion-comet/pull/6092#pullrequestreview-5293229643),
 and @sunchao reproduced it there independently. A no-op pull still parks, 
which is the #6091 case.
   
   ## How are these changes tested?
   
   `next_batch_does_not_park_after_a_pull_that_ran_a_nested_block_on` is new. 
Its pull closure runs a nested `block_on` of a 100 ms sleep while the stream 
waits on a 20 ms sleep. It asserts on elapsed time, because the test's 
ten-second timeout would otherwise rescue the lost wake and let it pass. With 
the loop parking after every pull it fails at 10.0 s, and with this change it 
returns in about 100 ms.
   
   The tests from #6092 cover the rest. 
`next_batch_parks_while_the_stream_waits_on_native_io` guards #6091: with the 
park removed, the pull closure ran 183,478 times in one 50 ms wait. 
`next_batch_resumes_on_a_refill_and_stops_pulling_after_eof` and 
`refill_wakes_the_pending_poll_and_eof_stays_buffered` cover the wakers and the 
buffered EOF.
   
   The core unit tests pass locally (465 passed, 5 ignored), and clippy and 
rustfmt are clean. I haven't rerun the JVM suites locally for the last commit. 
@mixermt ran the join, shuffle, task-metrics and Iceberg suites on the earlier 
version, and CI runs the Comet suites here plus the Spark 4.1 SQL suites 
through `run-spark-4.1-tests`.
   


-- 
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