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]

Reply via email to