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