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