[ 
https://issues.apache.org/jira/browse/SPARK-60041?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18124645#comment-18124645
 ] 

Yahya Kisana commented on SPARK-60041:
--------------------------------------

I can look into this. 

> 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
>            Priority: Major
>
> {{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