andygrove opened a new issue, #6289:
URL: https://github.com/apache/datafusion-comet/issues/6289

   ### Describe the bug
   
   A JVM input `ArrowArrayStream` that native never takes is never released. 
Its `ArrowReader` is never closed, and the JNI global reference that Arrow 
Java's C Data layer holds on the stream's private data is never deleted. That 
reference pins the reader and whatever the reader holds, including the source 
iterator and, for `CometArrowStream.wrapColumnarBatchRDD`, the first input 
batch that `reconcileStreamSchema` buffered. It stays pinned for the life of 
the executor. (`exportArrayStream` in arrow-java's `jni_wrapper.cc` creates the 
global ref, and only the release callback deletes it.)
   
   Native takes ownership of an input stream in `planner.create_plan` (the 
`OpStruct::Scan` arm), which runs on the first `executePlan` call. There, 
`AlignedArrowStreamReader::from_raw` moves the struct out and leaves the JVM 
copy released. Until then the JVM still owns it. On the JVM side, the only 
`release()` of an input stream is in `CometExecIterator`'s catch around 
`setShufflePartitionPusher`. The `createPlan` call itself sits outside that 
`try`. The task-completion listener that `CometArrowStream.stream` registers 
calls `streamRef.close()` and `allocator.close()`. `ArrowArrayStream.close()` 
frees the struct's memory and does not run the release callback.
   
   So an exported stream leaks whenever the task ends before native takes it:
   
   - **`createPlan` fails.** Verified: plan bytes that fail to deserialize 
leave the reader unclosed after task completion.
   - **The iterator is never polled.** Verified: `CometExecIterator` is built, 
`hasNext` is never called, and the task completes. In a real job this happens 
when the task is killed or cancelled between `CometExecRDD.compute` and the 
first `hasNext`.
   - **`CometExecRDD.compute` fails after `resolveInputObjects` has exported 
the streams.** Examples are `PlanDataInjector`, or `resolveInputObjects` itself 
failing on a later slot. That includes the first-batch read 
`reconcileStreamSchema` does for that slot, which leaves the earlier slots' 
streams exported. This one is inferred from the code.
   - **The first `executePlan` fails inside `create_plan` after it has taken 
some input streams but not all.** For example, a join plans its left child, and 
so takes the left stream, before the right. The taken readers are dropped and 
released, and the rest leak. Also inferred.
   
   ### Steps to reproduce
   
   A unit test in the `org.apache.spark` package, modeled on 
`CometExecIteratorLifecycleSuite`. It exports a stream with 
`CometArrowStream.stream` over an `ArrowReader` that records 
`closeReadSource()`, and passes it as the only input of a `CometExecIterator`. 
Then it calls `taskContext.markTaskCompleted(None)`:
   
   - with `CometExecUtils.getLimitNativePlan` as the plan and one `hasNext` 
call, the reader is closed;
   - with the same plan and no `hasNext` call, the reader is never closed;
   - with `protobufQueryPlan = Array[Byte](1, 2, 3)`, construction throws 
`CometNativeException: Fail to deserialize to native operator ...` and the 
reader is never closed.
   
   Reproduced on `main` at `31b38196d` (Spark 4.1).
   
   ### Expected behavior
   
   Every exported input stream is released exactly once, whether or not native 
ever took it.
   
   ### Additional context
   
   The smallest fix is probably in the `CometArrowStream.stream` listener: call 
`streamRef.release()` before `close()`. That listener runs after 
`CometExecIterator.close()` has released the native plan, because listeners 
fire in reverse registration order. The call is idempotent. `from_raw` leaves a 
null release pointer in the JVM struct, and arrow-java's 
`JniWrapper.releaseArrayStream` does nothing when that pointer is null, so a 
stream native took is not released twice. `CometExecIterator`'s own release 
loop could then go, or it could move to cover the `createPlan` call as well.
   


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