FelixYBW commented on issue #13140: URL: https://github.com/apache/gluten/issues/13140#issuecomment-5905274748
# The final plan: https://github.com/FelixYBW/spark/commits/SPARK-57468-convention-4.2 https://github.com/FelixYBW/gluten/commits/spark-convention-arrow-udf-4.2/ <img width="573" height="915" alt="Image" src="https://github.com/user-attachments/assets/8adc7998-38c9-4bdd-bd01-093846de5b81" /> ## 1. LoadArrowData (Velox → Java Arrow) Velox exports to Arrow ([`exportToArrow`](cpp/core/jni/JniCommon.cc) in `VeloxColumnarBatch.cc`), and Java imports it through the Arrow C Data Interface. - **Flat fixed-width columns** (int, long, double, date, etc.) — shared without copying. - **Dictionary- or constant-encoded Velox vectors** — flattened first via `ensureFlattened()`, which copies. - **String columns** — copied. Gluten's bridge options don't enable Arrow string views, so Velox's string layout is rebuilt as an Arrow value buffer plus offsets. --- ## 2. ArrowEvalPython No conversion copy. **Reuse Spark's ArrowEvalPythonExec** - Spark takes ownership of the input vectors by transfer (`ArrowBatchOwner`), which doesn't copy. - The UDF input columns are written from those buffers straight into the Arrow stream to the Python worker. That write is inherent to sending data to Python, not a conversion. - The result comes back as Arrow and is read into vectors allocated by Spark. --- ## 3. OffloadArrowData The copy is **only for UDF result columns**. Spark's output batch holds plain `ArrowColumnVector`s, so `ColumnarBatches.offload` first calls `adopt()` on each column. It moves a vector without copying if it has the same root allocator as Gluten's, and otherwise copies it with `VectorAppender`: - **Pass-through columns — no copy.** They were allocated by `LoadArrowData` in Gluten's per-task allocator, and Spark kept them there when it took ownership, so they are moved. - **UDF result columns — copied.** Spark's runner allocates them from `ArrowUtils.rootAllocator`. Gluten's per-task allocator is its own `RootAllocator` ([`ArrowBufferAllocators.java:87`](jvm/src/main/java/org/apache/gluten/memory/arrow/alloc/ArrowBufferAllocators.java)), and Arrow cannot move buffers between two roots. - After that, handing the batch to native code through the C Data Interface doesn't copy. --- ## 4. ArrowColumnarToVeloxColumnar Arrow-to-Velox conversion via `importFromArrowAsOwner`. - **Fixed-width buffers** — wrapped without copying. - **String columns** — get a new Velox string-view buffer, one entry per row. Strings of 12 bytes or fewer are copied into it; longer ones point into the Arrow data. --- # The test ```python import pyarrow as pa from pyspark.sql.functions import arrow_udf # Create a simple arrow UDF with @arrow_udf @arrow_udf("long") def multiply_arrow_func(a, b): return pa.compute.multiply(a, b) df=spark.read.parquet("file:///opt/spark/database/store_sales/") df.select(multiply_arrow_func(F.col("ss_promo_sk"),F.col("ss_ticket_number")).alias("udfout")).agg(F.sum("udfout")).show() ``` -- 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]
