Ben Hollis created SPARK-60066:
----------------------------------

             Summary: Nested UpdateFields expressions duplicate source 
expressions
                 Key: SPARK-60066
                 URL: https://issues.apache.org/jira/browse/SPARK-60066
             Project: Spark
          Issue Type: Improvement
          Components: SQL
    Affects Versions: 4.4.0
            Reporter: Ben Hollis


Spark represents `Column.withField` as the Catalyst expression `UpdateFields`. 
Before execution, Spark replaces that expression with a new struct containing 
the updated value and a read for every unchanged field. Each unchanged-field 
read contains the complete expression that produces the input struct.
 
When several updates are nested, the input to one replacement contains the 
replacement generated for the previous update. Spark repeats that entire 
earlier expression in every unchanged-field read. The logical expression can 
therefore grow much faster than the number of requested updates, increasing 
optimizer, expression-binding, and code-generation costs.
 
### Example
{code:java}
val updated = col("s")
.withField("nested.a", lit(1))
.withField("nested.b", lit(2))
.withField("nested.c", lit(3))


df.select(updated) {code}
 
Suppose `s` has the schema 
`struct<nested:struct<value:int>,d:int,e:int,f:int>`. The query changes three 
fields under `nested`; it does not change `d`, `e`, or `f`. Spark nevertheless 
generates reads for those unchanged fields at each replacement layer. Those 
reads repeat the expression generated for the earlier layer, so the Catalyst 
tree grows by far more than the three requested changes.



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