Peter Toth created SPARK-60041:
----------------------------------
Summary: Resolve nested DSv2 references without an Alias when
converting partition keys and sort orders
Key: SPARK-60041
URL: https://issues.apache.org/jira/browse/SPARK-60041
Project: Spark
Issue Type: Improvement
Components: SQL
Affects Versions: 4.2.0
Reporter: Peter Toth
{{V2ExpressionUtils.toCatalystOpt}} resolves a {{FieldReference}}, and the
column of an {{IdentityTransform}} or {{BucketTransform}}, with
{{resolveRef[NamedExpression]}}. For a nested field such as {{s.x}},
{{LogicalPlan.resolve}} returns {{Alias(GetStructField(...))}} with a new
{{ExprId}} on every call. The alias makes sense in a projection list, but not
inside a sort order or a partition key. Semantic comparison keeps the alias's
{{ExprId}}, so the converted expression never matches the same field converted
again, or a plain {{GetStructField}} from a query.
SPARK-59948 strips the alias from a scan's reported ordering in
{{V2ScanPartitioningAndOrdering}}. The same alias still reaches two other
places:
* _Reported partitioning._ A scan that reports {{bucket(4, s.x)}} gets the key
{{transformexpression(..., s#7.x AS x#11, Some(4))}}.
** {{PlanMerger.combineRequiredKeyGroupedPartitioning}} declines to merge two
scans of the table. Measured: two scalar subqueries over the table keep 2
scans. Stripping the alias merges them into 1.
** {{KeyedPartitioning.supportsExpressions}} accepts a {{GetStructField}} chain
but not an {{Alias}}, so a storage-partitioned join never applies to a nested
partition key.
* _Write requirements._ {{DistributionAndOrderingUtils.prepareQuery}} converts
a table's required clustering and ordering through the same code. Measured for
the ordering, with SPARK-59948 applied: {{INSERT INTO dst SELECT * FROM src}},
where {{src}} reports and {{dst}} requires an ordering on {{s.x}}, keeps {{Sort
[s#15.x AS x#26L ASC NULLS FIRST]}} above a scan that reports {{s.x}}. With a
top-level column there is no sort.
Stripping the alias where {{toCatalystOpt}} resolves these references would fix
all three, and would make the strip in SPARK-59948 unnecessary. The other
callers of {{resolveRef}}, e.g. {{ResolveChangelogTable}} and
{{RewriteRowLevelCommand}}, need a named expression and stay as they are.
Allowing a storage-partitioned join on a nested partition key this way needs
its own tests.
Query results are correct. The cost is a lost scan merge, a lost
storage-partitioned join, and a redundant sort on write.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]