[ 
https://issues.apache.org/jira/browse/SPARK-59996?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18124672#comment-18124672
 ] 

Charan Rathore commented on SPARK-59996:
----------------------------------------

I'd like to work on nested support here. Would a schema-driven size estimator 
for string and binary leaves, running before Java conversion, be a good 
starting point? I'd cover arrays, maps, structs and nulls, and benchmark the 
per-row overhead. I see this depends on #59122; should the design build on that 
PR while it is under review?

> Support nested values in the pickle Python UDF row-size guard
> -------------------------------------------------------------
>
>                 Key: SPARK-59996
>                 URL: https://issues.apache.org/jira/browse/SPARK-59996
>             Project: Spark
>          Issue Type: Task
>          Components: PySpark
>    Affects Versions: 4.3.0
>            Reporter: Tengfei Huang
>            Priority: Major
>
> h3. Background
> SPARK-59824 adds an opt-in, pre-conversion row-size guard for pickle Python 
> UDF inputs. The initial implementation estimates top-level {{StringType}} and 
> {{BinaryType}} arguments after projection and before Java conversion and 
> pickle serialization.
> Nested values inside {{{}ArrayType{}}}, {{{}MapType{}}}, and {{StructType}} 
> are intentionally excluded and documented as unsupported.
> h3. Problem
> Large string or binary values inside nested columns contribute nothing to the 
> current estimate. A nested-heavy row can therefore pass the guard and still 
> cause executor OOM during conversion or pickling.
> h3. Why this is deferred
> Nested support requires a recursive, size-only traversal before conversion. 
> This needs separate design and performance evaluation to:
>  * Avoid materializing converted Java objects.
>  * Avoid excessive per-row traversal overhead.
>  * Handle arrays, maps, structs, nulls, deeply nested values, and 
> user-defined types correctly.
>  * Preserve zero overhead when the guard is disabled or the projected schema 
> has no supported leaves.
> The existing nested-aware {{PickledSizeAccumulator}} path runs during 
> conversion, so it cannot directly provide the required pre-conversion 
> protection.
> h3. Proposed work
>  * Design a schema-driven estimator for nested string and binary leaves.
>  * Run the estimator after UDF argument projection and before conversion and 
> pickling.
>  * Evaluate whether logic from {{EvaluatePython.toJava}} and 
> {{PickledSizeAccumulator}} can be factored or shared.
>  * Add coverage for arrays, maps, structs, null values, computed arguments, 
> limit boundaries, and disabled behavior.
>  * Measure the per-row performance overhead.
> h3. Related
>  * SPARK-59824
>  * [https://github.com/apache/spark/pull/59122]
>  * [https://github.com/apache/spark/pull/59122#pullrequestreview-5368972749]



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to