mbutrovich commented on code in PR #6219:
URL: https://github.com/apache/datafusion-comet/pull/6219#discussion_r4105889083


##########
native/core/src/execution/jni_api.rs:
##########
@@ -1008,17 +1005,62 @@ fn prepare_output(
 /// Because the input source could be another native execution stream, which
 /// will be executed in another tokio blocking thread. It causes JNI throw
 /// Java exception. So we pull input batches here and insert them into scan
-/// operators before polling the stream,
+/// operators before polling the stream. Returns whether any scan made a JNI 
call.
 #[inline]
-fn pull_input_batches(exec_context: &mut ExecutionContext) -> Result<(), 
CometError> {
-    exec_context.scans.iter_mut().try_for_each(|scan| {
-        scan.get_next_batch()?;
-        Ok::<(), CometError>(())
-    })?;
-    exec_context.shuffle_scans.iter_mut().try_for_each(|scan| {
-        scan.get_next_batch()?;
-        Ok::<(), CometError>(())
+fn pull_input_batches(exec_context: &mut ExecutionContext) -> Result<bool, 
CometError> {
+    let mut pulled = false;
+    for scan in exec_context.scans.iter_mut() {
+        pulled |= scan.get_next_batch()?;
+    }
+    for scan in exec_context.shuffle_scans.iter_mut() {
+        pulled |= scan.get_next_batch()?;
+    }
+    Ok(pulled)
+}
+
+/// Yields once, so the `block_on` thread sleeps until a waker registered by 
an earlier poll
+/// fires: a JVM-fed scan refilled by `pull_input_batches`, or native I/O that 
completed.
+async fn park_until_woken() {
+    let mut polled = false;
+    poll_fn(|_| {
+        if std::mem::replace(&mut polled, true) {
+            Poll::Ready(())
+        } else {
+            Poll::Pending
+        }
     })
+    .await
+}
+
+/// Drives `stream` to its next item. JVM-fed scans return `Pending` until 
`on_pending` refills
+/// them, so every pending poll runs it. `on_pending` returns whether it made 
a JNI call, and the
+/// loop parks only when it did not, because the stream is then waiting on 
native I/O.
+///
+/// After a JNI call the loop polls again without parking. The call can run 
another Comet plan on
+/// this thread, as when a native writer's input is itself native, and that 
plan's `block_on`
+/// shares this thread's parker. If the native I/O completes while the nested 
plan is parked, the
+/// nested park takes the wake-up, and a park here would wait for a wake that 
has already fired.
+async fn next_batch<S>(
+    stream: &mut S,
+    mut on_pending: impl FnMut() -> Result<bool, CometError>,
+) -> Result<Option<RecordBatch>, CometError>
+where
+    S: Stream<Item = DataFusionResult<RecordBatch>> + Unpin,
+{
+    loop {
+        match poll!(stream.next()) {
+            Poll::Ready(item) => return Ok(item.transpose()?),
+            Poll::Pending => {
+                // JNI call to pull batches from JVM into ScanExec operators.
+                // block_in_place lets tokio move other tasks off this worker
+                // while we wait for JVM data.
+                let pulled = tokio::task::block_in_place(&mut on_pending)?;
+                if !pulled {
+                    park_until_woken().await;
+                }
+            }
+        }
+    }
 }

Review Comment:
   > Thanks, this is a better design, and I've switched to it in 398b3f94e.
   
   I checked the new version locally. With the `flag.woken` check forced to 
false, `next_batch_polls_again_when_a_nested_block_on_took_the_wake_up` fails 
after 10.0 s with "a wake was lost, and only the timeout's timer woke the 
task", and at the head commit it passes.
   
   > On `acquireMemory` inside the poll: a nested plan can't run there, because 
tokio panics with "Cannot start a runtime from within a runtime" on a 
`block_on` outside `block_in_place`.
   
   That matches what I see. With `block_in_place(&mut on_pending)` replaced by 
a plain `on_pending()`, the same test panics with "Cannot start a runtime from 
within a runtime", so a nested plan inside the poll errors out rather than 
hanging.



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