andygrove opened a new pull request, #6372:
URL: https://github.com/apache/datafusion-comet/pull/6372

   ## Which issue does this PR close?
   
   Closes #6289.
   
   ## Rationale for this change
   
   Native takes ownership of an exported input `ArrowArrayStream` in 
`PhysicalPlanner::create_plan`, which runs on the plan's first `executePlan`. 
There, `AlignedArrowStreamReader::from_raw` moves the C struct out and leaves a 
released struct on the JVM side. Until then the JVM owns the stream. The 
task-completion listener that `CometArrowStream.stream` registers called only 
`ArrowArrayStream.close()`, which frees the struct's memory but does not run 
the release callback. The only JVM-side `release()` was a loop in 
`CometExecIterator`'s constructor. It ran only when `setShufflePartitionPusher` 
failed, after `createPlan` had already succeeded.
   
   So a stream that native never took kept its `ArrowReader` open. arrow-java's 
`exportArrayStream` creates a JNI global ref on the stream's private data, and 
only the release callback deletes it. That ref kept the reader, its source 
iterator and any first batch buffered by `reconcileStreamSchema` reachable for 
the life of the executor. This happens in four cases:
   
   - `createPlan` fails.
   - The iterator is never polled, for example when a task is killed before its 
first `hasNext`.
   - `CometExecRDD.compute` fails after `resolveInputObjects` has exported the 
streams.
   - `create_plan` fails after taking some of a plan's inputs but not all of 
them.
   
   ## What changes are included in this PR?
   
   - The completion listener in `CometArrowStream.stream` now calls 
`streamRef.release()` before `close()`. The call is safe on every stream. 
`FFI_ArrowArrayStream::from_raw` (arrow-rs 59.3.0) replaces the JVM struct with 
an empty one whose release callback is null. arrow-java 18.3.0's 
`JniWrapper.releaseArrayStream` does nothing when that callback is null, so a 
stream native took is not released twice.
   - The listener runs after `CometExecIterator.close()` has released the 
native plan. `TaskContextImpl` keeps completion listeners on a `Stack` and runs 
them in reverse registration order (in both 3.4.3 and 4.1.3), and a stream is 
always exported before the iterator that consumes it is built.
   - `CometExecIterator` no longer has its own release loop in the construction 
failure handler. The listener now releases every stream native did not take, so 
the loop was redundant, and it only covered a failure after `createPlan` had 
succeeded. I considered extending the loop to cover `createPlan` as well. That 
would keep two JVM-side places releasing the same stream for no practical gain: 
a construction failure fails the task, and the listener runs as part of that 
task's completion. `CometExecIterator.close()` never released input streams 
either, so on the JVM side the listener that exported a stream is now the only 
place that releases it.
   - The existing `CometNativeShuffleSuite` test `failed RSS callback 
registration closes every Arrow and shuffle input and preserves its error` used 
to check the Arrow reader inside the task, right after construction failed. It 
now checks the reader from a completion listener registered before the stream's 
listener, so the check runs after the release. The result comes back through an 
accumulator.
   
   ## How are these changes tested?
   
   I added three tests to `CometExecIteratorLifecycleSuite`. Each exports a 
stream with `CometArrowStream.stream` over an `ArrowReader` that counts its 
`closeReadSource()` calls, passes the stream as the only input of a 
`CometExecIterator`, calls `markTaskCompleted(None)`, and then asserts that the 
reader closed exactly once:
   
   - `an input stream native took is released once, when its plan is released` 
calls `hasNext` once on a limit plan. The reader closes when the plan is 
released and does not close again at task end.
   - `an input stream native never took is released when the task ends` uses 
the same plan and never calls `hasNext`.
   - `an input stream is released when the task ends after its plan fails to be 
created` uses plan bytes `[1, 2, 3]`, so construction throws 
`CometNativeException: Fail to deserialize to native operator ...`.
   
   Before the fix, the first test passes and the other two fail with `0 did not 
equal 1 the input reader closed 0 times`.
   
   With the fix, on the default Spark 4.1 profile on macOS:
   
   - All 16 tests in `CometExecIteratorLifecycleSuite` and the adjusted 
`CometNativeShuffleSuite` RSS test pass (17 tests).
   - `CometNativeShuffleSuite`, `CometNativeSuite`, `CometArrowStreamSuite` and 
`CometInMemoryCacheSuite` pass (141 tests). These export input streams through 
`CometArrowStream.stream` end to end, and the log has no leaked-memory or 
listener errors.
   


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