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]

Reply via email to