binwei yang created SPARK-59792:
-----------------------------------
Summary: Skip ColumnarToRow for Arrow-backed input to Python UDTFs
Key: SPARK-59792
URL: https://issues.apache.org/jira/browse/SPARK-59792
Project: Spark
Issue Type: Improvement
Components: SQL, PySpark
Affects Versions: 5.0.0
Reporter: binwei yang
SPARK-56350 lets ArrowEvalPythonExec accept Arrow-backed ColumnarBatch input,
so Arrow Python UDFs no longer pay a columnar -> row -> columnar round-trip
(ColumnarToRowExec, then ArrowWriter) when the child already produces Arrow
data.
Arrow-optimized Python UDTFs (ArrowEvalPythonUDTFExec, used for
SQL_ARROW_TABLE_UDF and SQL_ARROW_UDTF) still always go through
ColumnarToRowExec and re-encode their arguments row by row with ArrowWriter,
even when the child is an Arrow-backed DataSource V2 scan.
Proposal: apply the same approach to ArrowEvalPythonUDTFExec.
* When the child supports columnar output and all the UDTF arguments are
columns of the child, read the columnar child directly and serialize the
argument columns to the Python worker from their Arrow vectors (reusing
ColumnarArrowPythonInput from SPARK-56350).
* A UDTF returns any number of rows per input row (plus the rows of
terminate()), so its output is still joined row by row with the input rows;
only the child columns in the output need to be buffered.
* UDTF arguments that are not child columns (literals, expressions) keep the
existing row-based path.
* Guarded by a new conf,
spark.sql.execution.arrow.pythonUDTF.columnarInput.enabled (default true).
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]