Repository navigation
Replies: 1 comment
|
+1 on the direction. The conversion cost is real: DataProcessor._set_output_tuple finalizes every output row individually, OutputManager.tuple_to_batch feeds partitioners one tuple at a time, and tuple_to_frame rebuilds the Arrow table from tuples. For a batch or table UDF that already produced a DataFrame, that round trip is pure overhead. Please check #8914 (per-port output, I implemented a prototype on a branch). You might want to take into account
Happy to coordinate so that the two branches don't have conflicts in design, merge conflicts can be resolved :) |
0 replies
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Uh oh!
There was an error while loading. Please reload this page.
Feature Summary
Allow Python UDFs to explicitly receive and return Arrow batches, avoiding the intermediate Python Tuple conversion while keeping existing UDF APIs unchanged.
PyAmber currently expands incoming Arrow tables into Tuples and converts UDF table output back through Tuples before rebuilding Arrow. This adds substantial overhead for operators that already work on entire batches.
An exploratory local conversion benchmark on upstream commit ec3a9dd, using PyArrow 23.0.1 and 10,000 rows with 10 string columns of 64-character values, measured:
This is approximately 77 times faster for the measured conversion path. It indicates potential savings from avoiding row conversion, not a measured 77 times improvement in workflow execution. The proposed engine path has not been implemented. The full Arrow Flight benchmark was blocked locally by a JOOQ schema mismatch, so an end-to-end benchmark is still needed.
Proposed Solution or Design
Introduce an explicit opt-in API, for example:
The engine would deliver Arrow batches directly and accept Arrow output without expanding every row into a Tuple. Syntax alone would not remove the current input and output conversions.
Start with ArrowBatchOperator. A separate ArrowTableOperator could later support whole-port input at completion. Validate the proposal with reproducible conversion and full engine benchmarks, plus tests for control handling and output equivalence where the APIs share semantics.
Affected Area
Workflow Engine (Amber)
Originally raised in #8476. Continuing the proposal here for discussion.
All reactions