comphead commented on issue #3532: URL: https://github.com/apache/datafusion-comet/issues/3532#issuecomment-5898245257
I checked each claim above against `main` at `1628c52` and the pinned Arrow Java 18.3.0 and arrow-rs 59.3.0 sources. All of them hold in the source. I didn't re-run the measurements. **Evidence** - **The original report is fixed.** The scan input moves the whole stream with `FFI_ArrowArrayStream::from_raw` ([aligned_stream_reader.rs#L44](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/native/core/src/execution/operators/aligned_stream_reader.rs#L44)). The native-to-JVM import releases whatever it hasn't consumed, on failure and at EOF ([NativeUtil.scala#L192-L209](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala#L192-L209), [#L287-L291](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala#L287-L291)). - **Error path in native C2R.** `columnarToRowConvert` takes each column with `std::ptr::read` and returns on the first `?` ([jni_api.rs#L1916-L1924](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/native/core/src/execution/jni_api.rs#L1916-L1924)). Columns after the failing one are never read, and `convert` has no `try`/`finally` ([NativeColumnarToRowConverter.scala#L91-L97](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/NativeColumnarToRowConverter.scala#L91-L97)). `exportBatch` has the same gap if it throws part-way. - **Struct leak on every batch.** `exportBatchToAddresses` keeps only the addresses ([NativeUtil.scala#L102-L106](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala#L102-L106)), and only `close()` frees the struct buffers. `CometArrowAllocator` is a plain `RootAllocator` ([package.scala#L36](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/package.scala#L36)), which rounds 80 and 72 bytes up to 128 each. That matches the 19,200 bytes above: 25 batches × 3 columns × 256 bytes. - **Constant-column leak on every batch.** `exportCDataBuffers` retains each buffer it exports ([FieldVector.java#L85-L86](https://github.com/apache/arrow-java/blob/v18.3.0/vector/src/main/java/org/apache/arrow/vector/FieldVector.java#L85-L86)), so the export doesn't take ownership of the vector, although the `ConstantColumnVectors.materialize` doc says it does ([ArrowWriters.scala#L112-L113](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/spark/sql/comet/execution/arrow/ArrowWriters.scala#L112-L113)). `exportBatch` never closes it ([NativeUtil.scala#L154-L161](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/vector/NativeUtil.scala#L154-L161)). - **A JVM-side release would double free today.** Arrow Java's `release_exported` deletes `private_data` and clears only the struct it was called on ([jni_wrapper.cc#L185-L189](https://github.com/apache/arrow-java/blob/v18.3.0/c/src/main/cpp/jni_wrapper.cc#L185-L189)). After `ptr::read`, the JVM struct keeps its `release` pointer, and `releaseArray` calls any non-null `release` ([#L394-L395](https://github.com/apache/arrow-java/blob/v18.3.0/c/src/main/cpp/jni_wrapper.cc#L394-L395)). `FFI_ArrowArray::from_raw` leaves the source released ([ffi.rs#L234-L236](https://github.com/apache/arrow-rs/blob/59.3.0/arrow-data/src/ffi.rs#L234-L236)), which would make a JVM-side `finally` safe. **Impact** - All three leaks need `spark.comet.exec.columnarToRow.native.enabled=true`, which defaults to `false` ([CometConf.scala#L324-L333](https://github.com/apache/datafusion-comet/blob/1628c520096e12ba3fb13988d6ee235070e77e07/spark/src/main/scala/org/apache/comet/CometConf.scala#L324-L333)). - The struct leak is 256 bytes per column per batch. For a 10-column schema at the default 8192-row batch, that works out to about 300 MB per billion rows (arithmetic, not measured). - The constant-column leak needs a Spark `ConstantColumnVector` to reach native C2R, which is the only production caller of `exportBatch`. - The error path is hard to hit. `decode_string_arrays` repairs invalid UTF-8 rather than failing, and I found no other realistic trigger with Comet-produced vectors. - The memory is charged to the global `CometArrowAllocator`, not a per-task allocator. Task end doesn't free it, so it grows for the executor's lifetime. That means this isn't limited to shells and notebooks, as I suggested earlier. A failed task doesn't kill the executor, and two of these leaks happen on successful batches. -- 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]
