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


##########
native/core/src/execution/jni_api.rs:
##########
@@ -1021,6 +1022,68 @@ fn pull_input_batches(exec_context: &mut 
ExecutionContext) -> Result<(), CometEr
     })
 }
 
+/// Forwards a wake-up to the `block_on` task and records that it happened, so 
`next_batch` can
+/// tell whether the stream was woken even if something else took the wake-up 
from the thread's
+/// parker.
+struct WakeFlag {
+    woken: AtomicBool,
+    parent: Waker,
+}
+
+impl Wake for WakeFlag {
+    fn wake(self: Arc<Self>) {
+        self.wake_by_ref()
+    }
+
+    fn wake_by_ref(self: &Arc<Self>) {
+        self.woken.store(true, Ordering::Release);
+        self.parent.wake_by_ref();
+    }
+}
+
+/// Drives `stream` to its next item. JVM-fed scans return `Pending` until 
`on_pending` refills
+/// them, so every pending poll runs it, and the refill wakes the stream.
+///
+/// Each poll gives the stream a `WakeFlag` waker. If the flag is still clear 
when `on_pending`
+/// returns, the stream is waiting on native I/O and `block_on` parks until it 
completes.
+/// Otherwise `next_batch` wakes `block_on` itself so that it polls again at 
once. It can't rely on
+/// the wake-up that set the flag, because `on_pending` can run another Comet 
plan on this thread,
+/// as when a native writer's input is itself native. That plan's `block_on` 
shares this thread's
+/// parker, which holds a single wake-up, and its park can take this one.
+///
+/// It polls again by yielding to `block_on` rather than looping here, so that 
each poll starts
+/// with a fresh coop budget. A stream that has spent its budget wakes itself 
and returns
+/// `Pending`, and a loop here would only get past that because 
`block_in_place` happens to leave
+/// this thread's budget unconstrained.
+async fn next_batch<S>(
+    stream: &mut S,
+    mut on_pending: impl FnMut() -> Result<(), CometError>,
+) -> Result<Option<RecordBatch>, CometError>
+where
+    S: Stream<Item = DataFusionResult<RecordBatch>> + Unpin,
+{
+    poll_fn(|cx| {
+        let flag = Arc::new(WakeFlag {
+            woken: AtomicBool::new(false),
+            parent: cx.waker().clone(),
+        });
+        let waker = Waker::from(Arc::clone(&flag));
+        if let Poll::Ready(item) = stream.poll_next_unpin(&mut 
Context::from_waker(&waker)) {
+            return Poll::Ready(Ok(item.transpose()?));
+        }
+        // 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.

Review Comment:
   Thanks, that comment was left over from the old loop. I applied your 
suggestion in b83dce380.



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