mbutrovich commented on code in PR #6219:
URL: https://github.com/apache/datafusion-comet/pull/6219#discussion_r4105893189
##########
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:
This comment came over from the old loop, but I don't think it describes
what `block_in_place` does here anymore. `next_batch` runs inside
`Runtime::block_on` on the Spark task thread, which isn't a tokio worker (the
`development.md` change in this PR says the same). On that thread
`block_in_place` has no worker core to hand off. It only stops the coop budget
and calls `exit_runtime` ([tokio 1.53.1
`worker.rs`](https://github.com/tokio-rs/tokio/blob/75fef53d0a8590c2d1dbb63672aa7b7d1ef51155/tokio/src/runtime/scheduler/multi_thread/worker.rs#L418-L429)
and
[L497-L505](https://github.com/tokio-rs/tokio/blob/75fef53d0a8590c2d1dbb63672aa7b7d1ef51155/tokio/src/runtime/scheduler/multi_thread/worker.rs#L497-L505)).
Leaving the runtime context is what lets a nested Comet plan's `block_on` run
inside `on_pending`. When I replaced the call with `on_pending()`, the nested
`block_on` test panicked with "Cannot start a runtime from within a runtime".
Since the flag design depends on that nested
case, could the comment say why the call is there? `on_pending` also isn't
only a JNI pull anymore, since it runs the metrics update too. Something like:
```suggestion
// `on_pending` calls into the JVM, which can run another Comet plan
on this thread.
// `block_in_place` exits the runtime context so that plan's
`block_on` doesn't panic.
```
--
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]