[ 
https://issues.apache.org/jira/browse/SPARK-59995?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Peter Toth updated SPARK-59995:
-------------------------------
    Summary: Drop sort orders that hold a partition transform from a V2 scan's 
output ordering  (was: GroupPartitionsExec's k-way merge fails when the 
ordering has a partition transform)

> Drop sort orders that hold a partition transform from a V2 scan's output 
> ordering
> ---------------------------------------------------------------------------------
>
>                 Key: SPARK-59995
>                 URL: https://issues.apache.org/jira/browse/SPARK-59995
>             Project: Spark
>          Issue Type: Bug
>          Components: SQL
>    Affects Versions: 4.2.0
>            Reporter: Peter Toth
>            Priority: Major
>              Labels: pull-request-available
>
> With {{spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled}} 
> on, {{GroupPartitionsExec}} can coalesce partitions with a k-way merge over 
> its child's ordering (SPARK-55715). The merge compares rows with a 
> code-generated ordering. When that ordering contains a partition transform, 
> such as {{years(arrive_time)}}, the task fails:
> {code}
> [INTERNAL_ERROR] Cannot generate code for expression: 
> transformexpression(org.apache.spark.sql.connector.catalog.functions.YearsFunction$@...,
>  input[2, timestamp, true], None) SQLSTATE: XX000
> {code}
> * {{TransformExpression.eval}} calls the bound {{ScalarFunction}}, but 
> {{TransformExpression.doGenCode}} always throws.
> * The merge's ordering ({{LazyCodeGenOrdering}} in {{GroupPartitionsExec}}) 
> calls {{GenerateOrdering}} directly. Unlike {{RowOrdering.create}}, it has no 
> interpreted fallback.
> The ordering can contain a transform in two ways. Both are measured on master:
> * The source reports it via {{SupportsReportOrdering}}. Example: a table 
> partitioned by {{identity(id)}} reports {{[id, name, years(arrive_time)]}}. A 
> sort-merge join on {{(id, name)}} fails with:
> ** {{spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled=true}}
> ** {{spark.sql.requireAllClusterKeysForCoPartition=false}}
> * The scan derives it from its partition keys (SPARK-56241). Example: a table 
> partitioned by {{(identity(id), years(arrive_time))}} derives {{[id, 
> years(arrive_time)]}}. A sort-merge join on {{id}} fails with:
> ** {{spark.sql.sources.v2.bucketing.preserveOrderingOnCoalesce.enabled=true}}
> ** 
> {{spark.sql.sources.v2.bucketing.allowKeysSubsetOfPartitionKeys.enabled=true}}
> ** {{spark.sql.sources.v2.bucketing.partitionKeyOrdering.enabled=true}} (the 
> default on master)
> ** 
> {{spark.sql.sources.v2.bucketing.preserveKeyOrderingOnCoalesce.enabled=false}}
> In both examples two partitions with {{id}} = 1 are coalesced. The merge's 
> ordering satisfies the join's required ordering, so {{EnsureRequirements}} 
> picks the merge instead of a sort. In the second example, 
> {{preserveKeyOrderingOnCoalesce}} on keeps the ordering on {{id}} without the 
> merge, so the query does not fail.
> Found while working on SPARK-59981, which derives the ordering in more cases.



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