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


##########
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:
   Since your approval, 5d2423231 also makes the acquire blocking-aware, as you 
first suggested. `SparkMemory` now calls Spark's `acquireMemory` inside 
`block_in_place`, so a worker blocked there hands its other tasks to another 
thread, and aborted DataFusion tasks holding the memory it waits for can still 
be cancelled. `a_blocked_acquire_leaves_its_worker_running_other_tasks` covers 
it: on a one-worker runtime, an acquire that waits for another task to run 
times out without the change and completes with it.
   
   A micro-benchmark with an acquire that takes 1 µs measured about 0.1 to 0.3 
µs added per call for one task, and up to about 2.4 µs with eight tasks 
acquiring at once. The Comet suites I ran took no longer. Stopping the producer 
still drops the stream on the calling thread, so releasing a plan does not 
depend on the scheduler either way.
   



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