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


##########
native/core/src/execution/jni_api.rs:
##########
@@ -1081,6 +1088,79 @@ where
     .await
 }
 
+/// Runs a plan that has no JVM input on a Tokio task, which sends the plan's 
batches to the Spark
+/// task thread.
+struct BatchProducer {
+    batches: mpsc::Receiver<DataFusionResult<RecordBatch>>,
+    /// Owns the plan's stream, and with it every reservation the stream holds.
+    task: JoinHandle<()>,
+    runtime: Handle,
+}
+
+impl BatchProducer {
+    fn spawn(runtime: Handle, mut stream: SendableRecordBatchStream) -> Self {
+        // Channel capacity of 2 allows the producer to work one batch
+        // ahead while the consumer processes the current one via JNI,
+        // without buffering excessive memory. Increasing this would
+        // trade memory for latency hiding if JNI/FFI overhead dominates;
+        // decreasing to 1 would serialize production and consumption.
+        let (tx, batches) = mpsc::channel(2);
+        let task = runtime.spawn(async move {
+            let result = std::panic::AssertUnwindSafe(async {
+                while let Some(batch) = stream.next().await {
+                    if tx.send(batch).await.is_err() {
+                        break;
+                    }
+                }
+            })
+            .catch_unwind()
+            .await;
+
+            if let Err(panic) = result {
+                let msg = match panic.downcast_ref::<&str>() {
+                    Some(s) => s.to_string(),
+                    None => match panic.downcast_ref::<String>() {
+                        Some(s) => s.clone(),
+                        None => "unknown panic".to_string(),
+                    },
+                };
+                let _ = tx
+                    .send(Err(DataFusionError::Execution(format!(
+                        "native panic: {msg}"
+                    ))))
+                    .await;
+            }
+        });
+        Self {
+            batches,
+            task,
+            runtime,
+        }
+    }
+
+    /// Stops the task, and parks the calling thread until it has finished, by 
which point it has
+    /// dropped the plan's stream.
+    ///
+    /// Closing the channel stops a task that is waiting to send a batch, and 
aborting it stops one
+    /// that is waiting on the stream. A task in the middle of polling the 
stream stops when that
+    /// poll returns, so this waits at most for the work the stream does 
between two await points.
+    fn stop(self) -> CometResult<()> {
+        let Self {
+            batches,
+            task,
+            runtime,
+        } = self;
+        drop(batches);
+        task.abort();
+        match runtime.block_on(task) {

Review Comment:
   Reproduced: `stopping_a_batch_producer_does_not_need_a_free_worker` holds 
the only worker of a one-worker runtime in a synchronous wait for the stopped 
plan's memory, and on 18ddc8f09 `stop()` timed out waiting for a free worker.
   
   9d405f39e removes the dependence on a worker instead of making the acquire 
blocking-aware. The task locks the stream only while polling it, and `stop()` 
takes the stream and drops it on the calling thread, as `releasePlan` already 
does for a JVM-fed plan. That waits only for a poll the task is already in, 
which is already running on a worker. The task is still aborted, but no longer 
joined. The test passes now.
   
   Tasks that DataFusion operators spawn, such as the sort merge's per-run 
tasks, still need a worker to be cancelled. If every worker is blocked waiting 
on the memory they hold, the `PlanMemoryPool` wait gives up after one second, 
the task ends, and Spark's cleanup wakes the waiters, so that case falls back 
to the late release instead of deadlocking. Wrapping the JNI acquire in 
`block_in_place` would cover it too, but it hands the worker off to another 
thread on every acquire, so I'd rather do that separately if the fallback turns 
out to matter.
   



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