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


##########
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:
   [P1] Keep producer cancellation runnable before joining it. With 
`COMET_WORKER_THREADS=1` and two Spark tasks, task A can stop early while its 
native producer still holds off-heap reservations. Task B can occupy the sole 
Tokio worker inside Spark's `ExecutionMemoryPool.acquireMemory`, waiting for A 
to release memory. `task.abort()` only schedules cancellation, so this join 
waits for a worker that cannot become available until A releases its 
reservations. A also cannot return to Spark's final task-memory cleanup. Both 
tasks therefore hang, and the later one-second timeout is never reached. The 
base teardown returned and allowed Spark cleanup to unblock B. Please make 
potentially blocking JNI memory acquisition blocking-aware, or otherwise 
guarantee cancellation can execute while preserving the memory-release ordering.
   
   Evidence: Compiled a disposable harness with byte-for-byte current 
`BatchProducer` code and the pinned Tokio 1.53.1 dependencies. On a one-worker 
runtime, producer A yields one batch then remains pending while retaining 
memory. Another task occupies the worker with a synchronous wait for that 
memory. Calling `stop()` from a separate thread remained blocked through the 
300 ms observation window and completed only after externally releasing the 
wait. The base receiver-drop/detach control returned immediately. Wrapping 
acquisition in `tokio::task::block_in_place` also allowed the head 
implementation to finish. Reproduction: `timeout 30 
/tmp/comet-6261-current-review/lifecycle-tests 
review_stop_with_saturated_workers_and_memory_wait --nocapture`. Production 
source confirms `JniMemoryManager::acquire` calls Spark synchronously without a 
blocking-aware wrapper, Spark's acquisition uses `lock.wait()`, and executor 
memory cleanup follows task completion.



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