timsaucer opened a new issue, #1705:
URL: https://github.com/apache/datafusion-python/issues/1705

   **Is your feature request related to a problem or challenge? Please describe 
what you are trying to do.**
   
   In a distributed setup the scheduler plans a query and hands stages to 
executors; only the executors ever call a Python UDF. Decoding an inlined 
Python UDF unpickles the function, which requires a Python interpreter and 
every module the function closes over to be importable. Doing that on the 
scheduler costs work nobody needs and forces the scheduler image to carry 
Python and the full dependency set of user code it will never run. Raised in 
https://github.com/apache/datafusion-python/pull/1678#pullrequestreview-5100366976.
   
   **Describe the solution you'd like**
   
   An opaque decode path: a `ScalarUDFImpl` that holds the still-pickled blob 
rather than a live Python object, and a codec that produces it. A scheduler 
installs that codec, decodes a plan into something it can inspect, route, and 
re-encode, and never touches cloudpickle. The executor installs the ordinary 
codec and unpickles as it does now.
   
   The wire format already allows this. An inlined UDF payload is `DFPYUDF` 
followed by a version byte and the cloudpickle blob 
(`crates/core/src/codec.rs`), so an opaque holder can carry those bytes 
verbatim and no format change is needed.
   
   The part that needs design is re-encoding. A scheduler that forwards a stage 
has to emit the blob byte-identically, so the executor sees exactly what the 
client wrote. That also raises what such a UDF should report for the things 
DataFusion asks of a `ScalarUDFImpl` during planning — name, signature, and 
return type are all recoverable from the payload without unpickling, since they 
are stored alongside the function, but `invoke` has to be an error rather than 
a surprise.
   
   **Describe alternatives you've considered**
   
   Encoding Python UDFs by name only and registering them on every node. 
Already supported and appropriate when the function is available everywhere; it 
does not cover the case inlining exists for, which is a function the receiving 
process does not have.
   
   Having the scheduler unpickle and immediately drop the object. Keeps the 
code simple, and still requires Python plus all user dependencies on the 
scheduler, which is the actual cost being avoided.
   
   **Additional context**
   
   Follow-up from #1678, which made extension codecs compose so a setup like 
this can install a scheduler-side codec alongside others. Likely also depends 
on #1703, gating `pyo3/extension-module`, if the consumer is a Rust crate 
rather than a Python process.
   


-- 
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