FelixYBW commented on PR #12276: URL: https://github.com/apache/gluten/pull/12276#issuecomment-5824667485
### Memory allocation and copies: Arrow UDF vs. Arrow UDTF | | `ColumnarArrowEvalPythonExec` (UDF) | `ColumnarArrowEvalPythonUDTFExec` (UDTF) | |---|---|---| | Python runner | Gluten's `ColumnarArrowPythonRunner` | Spark's `ArrowPythonUDTFRunner` (through `ArrowEvalPythonUDTFShim`) | | Allocator the Python output is read into | Gluten's `ArrowBufferAllocators.contextInstance()`, tracked by Gluten's off-heap memory management | Spark's `ArrowUtils.rootAllocator` child: off-heap (Netty direct memory) but not tracked by Spark or Gluten memory management | | Copies of the Python output | 1: socket → Gluten Arrow buffers | 2: socket → Spark Arrow buffers → Gluten Arrow buffers (`VectorAppender`, bulk per column) | | Output batch type | `ArrowJavaBatchType` | `ArrowJavaBatchType` | | Path into Velox | `OffloadArrowData` → `ArrowColumnarToVeloxColumnar` (C Data Interface + `importFromArrowAsOwner`, no copy) | same | **Why the UDTF copies the Python output into Gluten's allocator.** Spark's reader reuses one `VectorSchemaRoot`, so each batch's buffers are released when the next batch loads. When the stream ends, the reader calls `reader.close(false); allocator.close()` right away, not at task end. If the batches were wrapped and offloaded without a copy, a Velox operator still holding the imported buffers (hash build, aggregation, …) would make `allocator.close()` throw `Memory was leaked` as soon as the UDTF's output stream ends. The copy runs before the next `read()`, so nothing downstream refers to Spark's allocator, and the data is then counted in Gluten's off-heap budget. **Compared with the old UDTF output path.** Before, the output was `VanillaBatchType`, so a Velox consumer got `ColumnarToRow → RowToVeloxColumnar`: a per-value row round trip. Now it is `ColumnarArrowEvalPythonUDTF → OffloadArrowData → ArrowColumnarToVeloxColumnar`. **Possible follow-up.** A Gluten UDTF runner that reads the Python stream with Gluten's allocator, like `ColumnarArrowPythonRunner` does, would remove the extra copy. The cost is owning the UDTF worker protocol, which differs between Spark 3.5, 4.0 and 4.1. -- 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]
