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]