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