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]