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]
