weiting-chen commented on code in PR #13163:
URL: https://github.com/apache/gluten/pull/13163#discussion_r4150187823
##########
gluten-ut/pom.xml:
##########
@@ -238,5 +238,11 @@
<module>spark41</module>
</modules>
</profile>
+ <profile>
+ <id>spark-4.2</id>
+ <modules>
+ <module>spark42</module>
Review Comment:
**Use the annotations-specific Jackson version before enabling Spark42 UT**
**Target Location:** `gluten-ut/pom.xml:242-245`; the dependency to update
is at lines 142-145.
**Problem:** This activates a Spark42 UT reactor whose parent explicitly
resolves `jackson-annotations` with `${fasterxml.version}` (2.21.2), instead of
`${fasterxml.annotations.version}` (2.21). The declaration is inherited, not
newly added here, but it prevents the new module from building.
**Evidence:**
```xml
<id>spark-4.2</id>
<modules>
<module>spark42</module>
</modules>
```
The [new group3
job](https://github.com/apache/gluten/actions/runs/36683918472/job/109792884514)
fails before `gluten-ut-spark42` with `Could not find artifact
com.fasterxml.jackson.core:jackson-annotations:jar:2.21.2`. Extended and
slow-hive jobs fail identically. Maven Central has the annotations 2.21
artifact, not 2.21.2; explicit dependency versions override root dependency
management.
**Suggested Fix:** Use the already-defined annotations property in the UT
parent's dependency, then rerun the Spark42 UT lanes:
```xml
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson-annotations</artifactId>
<version>${fasterxml.annotations.version}</version>
</dependency>
```
##########
gluten-ut/spark42/src/test/scala/org/apache/spark/sql/connector/GlutenKeyGroupedPartitioningSuite.scala:
##########
@@ -0,0 +1,2088 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.spark.sql.connector
+
+import org.apache.gluten.config.GlutenConfig
+import org.apache.gluten.execution.SortMergeJoinExecTransformer
+
+import org.apache.spark.SparkConf
+import org.apache.spark.sql.{DataFrame, GlutenSQLTestsBaseTrait, Row}
+import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning
+import org.apache.spark.sql.connector.catalog.{Column, Identifier,
InMemoryTableCatalog}
+import org.apache.spark.sql.connector.distributions.Distributions
+import org.apache.spark.sql.connector.expressions.Expressions.{bucket, days,
identity, years}
+import org.apache.spark.sql.connector.expressions.Transform
+import org.apache.spark.sql.execution.{ColumnarShuffleExchangeExec, SparkPlan}
+import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
+import org.apache.spark.sql.execution.exchange.{ShuffleExchangeExec,
ShuffleExchangeLike}
+import org.apache.spark.sql.execution.joins.SortMergeJoinExec
+import org.apache.spark.sql.functions.{col, max}
+import org.apache.spark.sql.internal.SQLConf
+import org.apache.spark.sql.types._
+
+import java.util.Collections
+
+class GlutenKeyGroupedPartitioningSuite
+ extends KeyGroupedPartitioningSuite
+ with GlutenSQLTestsBaseTrait {
+ override def sparkConf: SparkConf = {
+ // Native SQL configs
+ super.sparkConf
+ .set(GlutenConfig.COLUMNAR_FORCE_SHUFFLED_HASH_JOIN_ENABLED.key, "false")
+ .set("spark.sql.adaptive.enabled", "false")
+ .set("spark.sql.shuffle.partitions", "5")
+ }
+
+ private val emptyProps: java.util.Map[String, String] = {
+ Collections.emptyMap[String, String]
+ }
+
+ private val columns: Array[Column] = Array(
+ Column.create("id", IntegerType),
+ Column.create("data", StringType),
+ Column.create("ts", TimestampType))
+
+ private val columns2: Array[Column] = Array(
+ Column.create("store_id", IntegerType),
+ Column.create("dept_id", IntegerType),
+ Column.create("data", StringType))
+
+ private def createTable(
+ table: String,
+ columns: Array[Column],
+ partitions: Array[Transform],
+ catalog: InMemoryTableCatalog = catalog): Unit = {
+ catalog.createTable(
+ Identifier.of(Array("ns"), table),
+ columns,
+ partitions,
+ emptyProps,
+ Distributions.unspecified(),
+ Array.empty,
+ None,
+ None,
+ numRowsPerSplit = 1)
+ }
+
+ private def collectColumnarShuffleExchangeExec(
+ plan: SparkPlan): Seq[ColumnarShuffleExchangeExec] = {
+ // here we skip collecting shuffle operators that are not associated with
SMJ
+ collect(plan) {
+ case s: SortMergeJoinExecTransformer => s
+ case s: SortMergeJoinExec => s
+ }.flatMap(smj => collect(smj) { case s: ColumnarShuffleExchangeExec => s })
+ }
+
+ override protected def collectShuffles(plan: SparkPlan):
Seq[ShuffleExchangeLike] = {
+ // here we skip collecting shuffle operators that are not associated with
SMJ
+ collect(plan) {
+ case s: SortMergeJoinExec => s
+ case s: SortMergeJoinExecTransformer => s
+ }.flatMap(
+ smj =>
+ collect(smj) {
+ case s: ShuffleExchangeExec => s
+ case s: ColumnarShuffleExchangeExec => s
+ })
+ }
+
+ override protected def collectAllShuffles(plan: SparkPlan):
Seq[ColumnarShuffleExchangeExec] = {
+ collect(plan) { case s: ColumnarShuffleExchangeExec => s }
Review Comment:
**Keep vanilla exchanges visible to inherited shuffle assertions**
**Target Location:** `GlutenKeyGroupedPartitioningSuite.scala:103-105`.
**Problem:** This is now an override of Spark42's protected
`collectAllShuffles`, so inherited tests dispatch to it too. Returning only
columnar exchanges hides vanilla fallback exchanges from those tests. For
example, the enabled upstream `SPARK-55535: Order by on partitions keys` uses
this helper to assert that no shuffle remains. An unwanted vanilla shuffle
would now pass unnoticed. This is a nonblocking test-coverage issue, not
evidence of a production regression.
**Evidence:**
```scala
override protected def collectAllShuffles(plan: SparkPlan):
Seq[ColumnarShuffleExchangeExec] = {
collect(plan) { case s: ColumnarShuffleExchangeExec => s }
}
```
The [inherited call
sites](https://github.com/apache/spark/blob/32f7299601108917fb01920a54e084595b7b3bf8/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala#L3435-L3436)
rely on the broader contract.
**Suggested Fix:** Preserve that contract and split out the intentionally
columnar-only checks:
```scala
override protected def collectAllShuffles(plan: SparkPlan):
Seq[ShuffleExchangeLike] = {
collect(plan) {
case s: ShuffleExchangeExec => s
case s: ColumnarShuffleExchangeExec => s
}
}
private def collectAllColumnarShuffles(
plan: SparkPlan): Seq[ColumnarShuffleExchangeExec] = {
collect(plan) { case s: ColumnarShuffleExchangeExec => s }
}
```
Update the native-only checks at lines 1116-1119 and 2030-2032 to use
`collectAllColumnarShuffles`. Do not substitute the existing join-scoped
`collectColumnarShuffleExchangeExec`: it would miss the standalone sort
exchange expected by the latter test.
##########
.github/workflows/velox_backend_x86.yml:
##########
@@ -1508,3 +1508,132 @@ jobs:
**/target/*.log
**/gluten-ut/**/hs_err_*.log
**/gluten-ut/**/core.*
+
+ spark-test-spark42:
+ needs: [detect-changes, build-native-lib-centos-8]
+ if: >-
+ needs.detect-changes.outputs.java == 'true' ||
+ needs.detect-changes.outputs.shims42 == 'true' ||
+ needs.detect-changes.outputs.cpp == 'true'
+ runs-on: ubuntu-22.04
+ strategy:
+ fail-fast: false
+ matrix:
+ # Split tests into 3 groups to run in parallel and cut the ~2h
wall-clock time.
+ # group1 – streaming tests
+ # group2 – execution / catalyst / errors / extension tests
+ # group3 – top-level sql, connector, sources, hive and remaining
tests
+ group: [1, 2, 3]
+ env:
+ SPARK_TESTING: true
+ container: apache/gluten:centos-9-jdk17
+ steps:
+ - uses: actions/checkout@v7
+ - name: Download All Artifacts
+ uses: actions/download-artifact@v8
+ with:
+ name: velox-native-lib-centos-8-${{github.sha}}
+ path: ./cpp/build/releases/
+ - name: Prepare
+ run: |
+ dnf install -y python3.11 python3.11-pip python3.11-devel && \
+ ls -la /usr/bin/python3.11 && \
+ alternatives --install /usr/bin/python3 python3 /usr/bin/python3.11
1 && \
+ alternatives --set python3 /usr/bin/python3.11 && \
+ pip3 install setuptools==77.0.3 && \
+ pip3 install pyspark==3.5.5 cython && \
+ pip3 install pandas==2.2.3 pyarrow==20.0.0
+ - name: Build and Run unit test for Spark 4.2.0 with scala-2.13 (other
tests, group ${{ matrix.group }})
+ run: |
+ cd $GITHUB_WORKSPACE/
+ export SPARK_SCALA_VERSION=2.13
+ yum install -y java-17-openjdk-devel
+ export JAVA_HOME=/usr/lib/jvm/java-17-openjdk
+ export PATH=$JAVA_HOME/bin:$PATH
+ java -version
+
TAGS_EXCLUDE="org.apache.spark.tags.ExtendedSQLTest,org.apache.spark.tags.SlowHiveTest,org.apache.gluten.tags.UDFTest,org.apache.gluten.tags.EnhancedFeaturesTest,org.apache.gluten.tags.CudfTest,org.apache.gluten.tags.SkipTest"
+ if [ "${{ matrix.group }}" = "1" ]; then
+ # Group 1: streaming + gluten utils + top-level spark tests (~55
classes)
+ $MVN_CMD clean test -Pspark-4.2 -Pscala-2.13 -Pjava-17
-Pbackends-velox -Pspark-ut \
+ -DargLine="-Dspark.test.home=/opt/shims/spark42/spark_home/"
-DtagsToExclude="$TAGS_EXCLUDE" \
+
-DwildcardSuites="org.apache.spark.sql.streaming,org.apache.spark.GlutenSortShuffleSuite,org.apache.gluten"
Review Comment:
**Include the copied RPC suite in a Spark42 test group**
**Target Location:** `.github/workflows/velox_backend_x86.yml:1559`.
**Problem:** None of the new wildcard groups selects
`org.apache.spark.rpc.GlutenDriverEndpointSuite`. It is compiled from
`spark42/src/test/backends-velox`, but is a plain, untagged `SparkFunSuite`, so
the extended/slow-hive jobs do not pick it up either. This is a nonblocking
coverage omission; Spark41 has the same omission.
**Evidence:**
```sh
-DwildcardSuites="org.apache.spark.sql.streaming,org.apache.spark.GlutenSortShuffleSuite,org.apache.gluten"
```
The newly copied
`gluten-ut/spark42/src/test/backends-velox/org/apache/spark/rpc/GlutenDriverEndpointSuite.scala:24-25`
defines the concrete suite and ordinary test; none of the remaining group
prefixes includes its package.
**Suggested Fix:** Add its package to group1:
```sh
-DwildcardSuites="org.apache.spark.sql.streaming,org.apache.spark.GlutenSortShuffleSuite,org.apache.spark.rpc,org.apache.gluten"
```
No `VeloxTestSettings` registration is needed for this plain `SparkFunSuite`.
##########
.github/workflows/velox_backend_x86.yml:
##########
@@ -1508,3 +1508,132 @@ jobs:
**/target/*.log
**/gluten-ut/**/hs_err_*.log
**/gluten-ut/**/core.*
+
+ spark-test-spark42:
+ needs: [detect-changes, build-native-lib-centos-8]
+ if: >-
+ needs.detect-changes.outputs.java == 'true' ||
+ needs.detect-changes.outputs.shims42 == 'true' ||
+ needs.detect-changes.outputs.cpp == 'true'
+ runs-on: ubuntu-22.04
+ strategy:
+ fail-fast: false
+ matrix:
+ # Split tests into 3 groups to run in parallel and cut the ~2h
wall-clock time.
+ # group1 – streaming tests
+ # group2 – execution / catalyst / errors / extension tests
+ # group3 – top-level sql, connector, sources, hive and remaining
tests
+ group: [1, 2, 3]
+ env:
+ SPARK_TESTING: true
+ container: apache/gluten:centos-9-jdk17
+ steps:
+ - uses: actions/checkout@v7
+ - name: Download All Artifacts
+ uses: actions/download-artifact@v8
+ with:
+ name: velox-native-lib-centos-8-${{github.sha}}
+ path: ./cpp/build/releases/
+ - name: Prepare
+ run: |
+ dnf install -y python3.11 python3.11-pip python3.11-devel && \
+ ls -la /usr/bin/python3.11 && \
+ alternatives --install /usr/bin/python3 python3 /usr/bin/python3.11
1 && \
+ alternatives --set python3 /usr/bin/python3.11 && \
+ pip3 install setuptools==77.0.3 && \
+ pip3 install pyspark==3.5.5 cython && \
+ pip3 install pandas==2.2.3 pyarrow==20.0.0
+ - name: Build and Run unit test for Spark 4.2.0 with scala-2.13 (other
tests, group ${{ matrix.group }})
+ run: |
+ cd $GITHUB_WORKSPACE/
+ export SPARK_SCALA_VERSION=2.13
+ yum install -y java-17-openjdk-devel
+ export JAVA_HOME=/usr/lib/jvm/java-17-openjdk
+ export PATH=$JAVA_HOME/bin:$PATH
+ java -version
+
TAGS_EXCLUDE="org.apache.spark.tags.ExtendedSQLTest,org.apache.spark.tags.SlowHiveTest,org.apache.gluten.tags.UDFTest,org.apache.gluten.tags.EnhancedFeaturesTest,org.apache.gluten.tags.CudfTest,org.apache.gluten.tags.SkipTest"
+ if [ "${{ matrix.group }}" = "1" ]; then
+ # Group 1: streaming + gluten utils + top-level spark tests (~55
classes)
+ $MVN_CMD clean test -Pspark-4.2 -Pscala-2.13 -Pjava-17
-Pbackends-velox -Pspark-ut \
+ -DargLine="-Dspark.test.home=/opt/shims/spark42/spark_home/"
-DtagsToExclude="$TAGS_EXCLUDE" \
+
-DwildcardSuites="org.apache.spark.sql.streaming,org.apache.spark.GlutenSortShuffleSuite,org.apache.gluten"
+ elif [ "${{ matrix.group }}" = "2" ]; then
+ # Group 2: execution + catalyst + errors + extension tests (~140
classes)
+ $MVN_CMD clean test -Pspark-4.2 -Pscala-2.13 -Pjava-17
-Pbackends-velox -Pspark-ut \
+ -DargLine="-Dspark.test.home=/opt/shims/spark42/spark_home/"
-DtagsToExclude="$TAGS_EXCLUDE" \
+
-DwildcardSuites="org.apache.spark.sql.execution,org.apache.spark.sql.catalyst,org.apache.spark.sql.errors,org.apache.spark.sql.extension"
Review Comment:
**Adapt the inherited task-CPU error assertion for this Spark42 lane**
**Target Location:** `.github/workflows/velox_backend_x86.yml:1562-1564`;
the assertion to update is
`gluten-substrait/src/test/scala/org/apache/spark/sql/execution/GlutenAutoAdjustStageResourceProfileSuite.scala:126`.
**Problem:** The new group2 runs an existing test whose expected wording is
incompatible with Spark42. This is an inherited test-oracle issue exposed by
the new lane, not a newly introduced resource-profile bug, but it currently
prevents this lane from reaching the Spark42 UT module.
**Evidence:**
```sh
-DwildcardSuites="org.apache.spark.sql.execution,org.apache.spark.sql.catalyst,org.apache.spark.sql.errors,org.apache.spark.sql.extension"
```
The [group2
failure](https://github.com/apache/gluten/actions/runs/36683918472/job/109792884598)
reports that `updateResourceSetting rejects a non-positive task cpus` received:
```text
[INVALID_CONF_VALUE.REQUIREMENT] The value '0' in the config
"spark.task.cpus" is invalid.
Number of cores to allocate for each task should be positive. SQLSTATE: 22022
```
The existing assertion requires the contiguous substring `"spark.task.cpus
should be positive"`. The exception still correctly rejects the invalid
configuration.
**Suggested Fix:** Retain the `IllegalArgumentException` interception and
assert the key and requirement independently, so both old and structured
Spark42 wording satisfy the same semantic check:
```scala
assert(e.getMessage.contains("spark.task.cpus"))
assert(e.getMessage.contains("positive"))
```
Please include or coordinate this compatibility fix before relying on the
new group2 gate.
##########
gluten-ut/spark42/src/test/scala/org/apache/spark/sql/connector/GlutenKeyGroupedPartitioningSuite.scala:
##########
@@ -0,0 +1,2088 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.spark.sql.connector
+
+import org.apache.gluten.config.GlutenConfig
+import org.apache.gluten.execution.SortMergeJoinExecTransformer
+
+import org.apache.spark.SparkConf
+import org.apache.spark.sql.{DataFrame, GlutenSQLTestsBaseTrait, Row}
+import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning
+import org.apache.spark.sql.connector.catalog.{Column, Identifier,
InMemoryTableCatalog}
+import org.apache.spark.sql.connector.distributions.Distributions
+import org.apache.spark.sql.connector.expressions.Expressions.{bucket, days,
identity, years}
+import org.apache.spark.sql.connector.expressions.Transform
+import org.apache.spark.sql.execution.{ColumnarShuffleExchangeExec, SparkPlan}
+import org.apache.spark.sql.execution.datasources.v2.BatchScanExec
+import org.apache.spark.sql.execution.exchange.{ShuffleExchangeExec,
ShuffleExchangeLike}
+import org.apache.spark.sql.execution.joins.SortMergeJoinExec
+import org.apache.spark.sql.functions.{col, max}
+import org.apache.spark.sql.internal.SQLConf
+import org.apache.spark.sql.types._
+
+import java.util.Collections
+
+class GlutenKeyGroupedPartitioningSuite
+ extends KeyGroupedPartitioningSuite
+ with GlutenSQLTestsBaseTrait {
+ override def sparkConf: SparkConf = {
+ // Native SQL configs
+ super.sparkConf
+ .set(GlutenConfig.COLUMNAR_FORCE_SHUFFLED_HASH_JOIN_ENABLED.key, "false")
+ .set("spark.sql.adaptive.enabled", "false")
+ .set("spark.sql.shuffle.partitions", "5")
+ }
+
+ private val emptyProps: java.util.Map[String, String] = {
+ Collections.emptyMap[String, String]
+ }
+
+ private val columns: Array[Column] = Array(
+ Column.create("id", IntegerType),
+ Column.create("data", StringType),
+ Column.create("ts", TimestampType))
+
+ private val columns2: Array[Column] = Array(
+ Column.create("store_id", IntegerType),
+ Column.create("dept_id", IntegerType),
+ Column.create("data", StringType))
+
+ private def createTable(
+ table: String,
+ columns: Array[Column],
+ partitions: Array[Transform],
+ catalog: InMemoryTableCatalog = catalog): Unit = {
+ catalog.createTable(
+ Identifier.of(Array("ns"), table),
+ columns,
+ partitions,
+ emptyProps,
+ Distributions.unspecified(),
+ Array.empty,
+ None,
+ None,
+ numRowsPerSplit = 1)
+ }
+
+ private def collectColumnarShuffleExchangeExec(
+ plan: SparkPlan): Seq[ColumnarShuffleExchangeExec] = {
+ // here we skip collecting shuffle operators that are not associated with
SMJ
+ collect(plan) {
+ case s: SortMergeJoinExecTransformer => s
+ case s: SortMergeJoinExec => s
+ }.flatMap(smj => collect(smj) { case s: ColumnarShuffleExchangeExec => s })
+ }
+
+ override protected def collectShuffles(plan: SparkPlan):
Seq[ShuffleExchangeLike] = {
+ // here we skip collecting shuffle operators that are not associated with
SMJ
+ collect(plan) {
+ case s: SortMergeJoinExec => s
+ case s: SortMergeJoinExecTransformer => s
+ }.flatMap(
+ smj =>
+ collect(smj) {
+ case s: ShuffleExchangeExec => s
+ case s: ColumnarShuffleExchangeExec => s
+ })
+ }
+
+ override protected def collectAllShuffles(plan: SparkPlan):
Seq[ColumnarShuffleExchangeExec] = {
+ collect(plan) { case s: ColumnarShuffleExchangeExec => s }
+ }
+
+ private def collectVanillaShuffles(plan: SparkPlan):
Seq[ShuffleExchangeExec] = {
+ collect(plan) { case s: ShuffleExchangeExec => s }
+ }
+
+ private def collectScans(plan: SparkPlan): Seq[BatchScanExec] = {
+ collect(plan) { case s: BatchScanExec => s }
+ }
+
+ private def selectWithMergeJoinHint(t1: String, t2: String): String = {
+ s"SELECT /*+ MERGE($t1, $t2) */ "
+ }
+
+ private def createJoinTestDF(
+ keys: Seq[(String, String)],
+ extraColumns: Seq[String] = Nil,
+ joinType: String = ""): DataFrame = {
+ val extraColList = if (extraColumns.isEmpty) "" else
extraColumns.mkString(", ", ", ", "")
+ sql(s"""
+ |${selectWithMergeJoinHint("i", "p")}
+ |id, name, i.price as purchase_price, p.price as sale_price
$extraColList
+ |FROM testcat.ns.$items i $joinType JOIN testcat.ns.$purchases p
+ |ON ${keys.map(k => s"i.${k._1} = p.${k._2}").mkString(" AND ")}
+ |ORDER BY id, purchase_price, sale_price $extraColList
+ |""".stripMargin)
+ }
+
+ private val customers: String = "customers"
+ private val customersColumns: Array[Column] = Array(
+ Column.create("customer_name", StringType),
+ Column.create("customer_age", IntegerType),
+ Column.create("customer_id", LongType))
+
+ private val orders: String = "orders"
+ private val ordersColumns: Array[Column] =
+ Array(Column.create("order_amount", DoubleType),
Column.create("customer_id", LongType))
+
+ private def testWithCustomersAndOrders(
+ customers_partitions: Array[Transform],
+ orders_partitions: Array[Transform],
+ expectedNumOfShuffleExecs: Int): Unit = {
+ createTable(customers, customersColumns, customers_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$customers VALUES " +
+ s"('aaa', 10, 1), ('bbb', 20, 2), ('ccc', 30, 3)")
+
+ createTable(orders, ordersColumns, orders_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$orders VALUES " +
+ s"(100.0, 1), (200.0, 1), (150.0, 2), (250.0, 2), (350.0, 2), (400.50,
3)")
+
+ val df = sql(
+ "SELECT customer_name, customer_age, order_amount " +
+ s"FROM testcat.ns.$customers c JOIN testcat.ns.$orders o " +
+ "ON c.customer_id = o.customer_id ORDER BY c.customer_id,
order_amount")
+
+ val shuffles =
collectColumnarShuffleExchangeExec(df.queryExecution.executedPlan)
+ assert(shuffles.length == expectedNumOfShuffleExecs)
+
+ checkAnswer(
+ df,
+ Seq(
+ Row("aaa", 10, 100.0),
+ Row("aaa", 10, 200.0),
+ Row("bbb", 20, 150.0),
+ Row("bbb", 20, 250.0),
+ Row("bbb", 20, 350.0),
+ Row("ccc", 30, 400.50)))
+ }
+
+ testGluten("partitioned join: only one side reports partitioning") {
+ val customers_partitions = Array(bucket(4, "customer_id"))
+ val orders_partitions = Array(bucket(2, "customer_id"))
+
+ testWithCustomersAndOrders(customers_partitions, orders_partitions, 2)
+ }
+ testGluten("partitioned join: exact distribution (same number of buckets)
from both sides") {
+ val customers_partitions = Array(bucket(4, "customer_id"))
+ val orders_partitions = Array(bucket(4, "customer_id"))
+
+ testWithCustomersAndOrders(customers_partitions, orders_partitions, 0)
+ }
+
+ private val items: String = "items"
+ private val itemsColumns: Array[Column] = Array(
+ Column.create("id", LongType),
+ Column.create("name", StringType),
+ Column.create("price", FloatType),
+ Column.create("arrive_time", TimestampType))
+ private val purchases: String = "purchases"
+ private val purchasesColumns: Array[Column] = Array(
+ Column.create("item_id", LongType),
+ Column.create("price", FloatType),
+ Column.create("time", TimestampType))
+
+ testGluten(
+ "SPARK-41413: partitioned join: partition values" +
+ " from one side are subset of those from the other side") {
+ val items_partitions = Array(bucket(4, "id"))
+ createTable(items, itemsColumns, items_partitions)
+
+ sql(
+ s"INSERT INTO testcat.ns.$items VALUES " +
+ "(1, 'aa', 40.0, cast('2020-01-01' as timestamp)), " +
+ "(3, 'bb', 10.0, cast('2020-01-01' as timestamp)), " +
+ "(4, 'cc', 15.5, cast('2020-02-01' as timestamp))")
+
+ val purchases_partitions = Array(bucket(4, "item_id"))
+ createTable(purchases, purchasesColumns, purchases_partitions)
+
+ sql(
+ s"INSERT INTO testcat.ns.$purchases VALUES " +
+ "(1, 42.0, cast('2020-01-01' as timestamp)), " +
+ "(3, 19.5, cast('2020-02-01' as timestamp))")
+
+ Seq(true, false).foreach {
+ pushDownValues =>
+ withSQLConf(SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key ->
pushDownValues.toString) {
+ val df = sql(
+ "SELECT id, name, i.price as purchase_price, p.price as sale_price
" +
+ s"FROM testcat.ns.$items i JOIN testcat.ns.$purchases p " +
+ "ON i.id = p.item_id ORDER BY id, purchase_price, sale_price")
+
+ val shuffles =
collectColumnarShuffleExchangeExec(df.queryExecution.executedPlan)
+ if (pushDownValues) {
+ assert(shuffles.isEmpty, "should not add shuffle when partition
values mismatch")
+ } else {
+ assert(
+ shuffles.nonEmpty,
+ "should add shuffle when partition values mismatch, and " +
+ "pushing down partition values is not enabled")
+ }
+
+ checkAnswer(df, Seq(Row(1, "aa", 40.0, 42.0), Row(3, "bb", 10.0,
19.5)))
+ }
+ }
+ }
+
+ testGluten("SPARK-41413: partitioned join: partition values from both sides
overlaps") {
+ val items_partitions = Array(identity("id"))
+ createTable(items, itemsColumns, items_partitions)
+
+ sql(
+ s"INSERT INTO testcat.ns.$items VALUES " +
+ "(1, 'aa', 40.0, cast('2020-01-01' as timestamp)), " +
+ "(2, 'bb', 10.0, cast('2020-01-01' as timestamp)), " +
+ "(3, 'cc', 15.5, cast('2020-02-01' as timestamp))")
+
+ val purchases_partitions = Array(identity("item_id"))
+ createTable(purchases, purchasesColumns, purchases_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$purchases VALUES " +
+ "(1, 42.0, cast('2020-01-01' as timestamp)), " +
+ "(2, 19.5, cast('2020-02-01' as timestamp)), " +
+ "(4, 30.0, cast('2020-02-01' as timestamp))")
+
+ Seq(true, false).foreach {
+ pushDownValues =>
+ withSQLConf(SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key ->
pushDownValues.toString) {
+ val df = sql(
+ "SELECT id, name, i.price as purchase_price, p.price as sale_price
" +
+ s"FROM testcat.ns.$items i JOIN testcat.ns.$purchases p " +
+ "ON i.id = p.item_id ORDER BY id, purchase_price, sale_price")
+
+ val shuffles =
collectColumnarShuffleExchangeExec(df.queryExecution.executedPlan)
+ if (pushDownValues) {
+ assert(shuffles.isEmpty, "should not add shuffle when partition
values mismatch")
+ } else {
+ assert(
+ shuffles.nonEmpty,
+ "should add shuffle when partition values mismatch, and " +
+ "pushing down partition values is not enabled")
+ }
+
+ checkAnswer(df, Seq(Row(1, "aa", 40.0, 42.0), Row(2, "bb", 10.0,
19.5)))
+ }
+ }
+ }
+
+ testGluten("SPARK-41413: partitioned join: non-overlapping partition values
from both sides") {
+ val items_partitions = Array(identity("id"))
+ createTable(items, itemsColumns, items_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$items VALUES " +
+ "(1, 'aa', 40.0, cast('2020-01-01' as timestamp)), " +
+ "(2, 'bb', 10.0, cast('2020-01-01' as timestamp)), " +
+ "(3, 'cc', 15.5, cast('2020-02-01' as timestamp))")
+
+ val purchases_partitions = Array(identity("item_id"))
+ createTable(purchases, purchasesColumns, purchases_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$purchases VALUES " +
+ "(4, 42.0, cast('2020-01-01' as timestamp)), " +
+ "(5, 19.5, cast('2020-02-01' as timestamp)), " +
+ "(6, 30.0, cast('2020-02-01' as timestamp))")
+
+ Seq(true, false).foreach {
+ pushDownValues =>
+ withSQLConf(SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key ->
pushDownValues.toString) {
+ val df = sql(
+ "SELECT id, name, i.price as purchase_price, p.price as sale_price
" +
+ s"FROM testcat.ns.$items i JOIN testcat.ns.$purchases p " +
+ "ON i.id = p.item_id ORDER BY id, purchase_price, sale_price")
+
+ val shuffles =
collectColumnarShuffleExchangeExec(df.queryExecution.executedPlan)
+ if (pushDownValues) {
+ assert(shuffles.isEmpty, "should not add shuffle when partition
values mismatch")
+ } else {
+ assert(
+ shuffles.nonEmpty,
+ "should add shuffle when partition values mismatch, and " +
+ "pushing down partition values is not enabled")
+ }
+
+ checkAnswer(df, Seq.empty)
+ }
+ }
+ }
+
+ testGluten(
+ "SPARK-42038: partially clustered:" +
+ " with same partition keys and one side fully clustered") {
+ val items_partitions = Array(identity("id"))
+ createTable(items, itemsColumns, items_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$items VALUES " +
+ s"(1, 'aa', 40.0, cast('2020-01-01' as timestamp)), " +
+ s"(2, 'bb', 10.0, cast('2020-01-01' as timestamp)), " +
+ s"(3, 'cc', 15.5, cast('2020-02-01' as timestamp))")
+
+ val purchases_partitions = Array(identity("item_id"))
+ createTable(purchases, purchasesColumns, purchases_partitions)
+ sql(
+ s"INSERT INTO testcat.ns.$purchases VALUES " +
+ s"(1, 45.0, cast('2020-01-01' as timestamp)), " +
+ s"(1, 50.0, cast('2020-01-02' as timestamp)), " +
+ s"(2, 15.0, cast('2020-01-02' as timestamp)), " +
+ s"(2, 20.0, cast('2020-01-03' as timestamp)), " +
+ s"(3, 20.0, cast('2020-02-01' as timestamp))")
+
+ Seq(true, false).foreach {
+ pushDownValues =>
+ Seq(("true", 5), ("false", 3)).foreach {
+ case (enable, expected) =>
+ withSQLConf(
+ SQLConf.V2_BUCKETING_PUSH_PART_VALUES_ENABLED.key ->
pushDownValues.toString,
+
SQLConf.V2_BUCKETING_PARTIALLY_CLUSTERED_DISTRIBUTION_ENABLED.key -> enable
+ ) {
+ val df = sql(
+ "SELECT id, name, i.price as purchase_price, p.price as
sale_price " +
+ s"FROM testcat.ns.$items i JOIN testcat.ns.$purchases p " +
+ "ON i.id = p.item_id ORDER BY id, purchase_price,
sale_price")
+
+ val shuffles =
collectColumnarShuffleExchangeExec(df.queryExecution.executedPlan)
+ assert(shuffles.isEmpty, "should not contain any shuffle")
+ if (pushDownValues) {
+ val scans = collectScans(df.queryExecution.executedPlan)
+ assert(scans.forall(_.inputRDD.partitions.length == expected))
+ }
Review Comment:
**Assert Spark42 grouping output, not the raw scan split count**
**Target Location:** `GlutenKeyGroupedPartitioningSuite.scala:361-364` and
the same migrated oracle at lines 1601-1611.
**Problem:** Spark42 moved SPJ grouping/replication out of `BatchScanExec`
into `GroupPartitionsExec`. This fixture sets `numRowsPerSplit = 1`, so the
first test's three item rows and five purchase rows produce raw scan counts
`[3, 5]`. Neither `forall(_ == 5)` nor `forall(_ == 3)` can pass. The later
SPARK-47094 fixture likewise has raw counts `[3, 3]`, not its asserted `[2,
2]`/`[3, 2]`. The rewritten `Gluten - ...` tests remain enabled; exclusions for
the unprefixed upstream names do not exclude them.
**Evidence:**
```scala
val scans = collectScans(df.queryExecution.executedPlan)
assert(scans.forall(_.inputRDD.partitions.length == expected))
```
The [pinned Spark42
suite](https://github.com/apache/spark/blob/32f7299601108917fb01920a54e084595b7b3bf8/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala#L830-L835)
checks grouping-node output instead. These are source-contract failures, not
failures observed in the current CI run: dependency resolution stops group3
earlier.
**Suggested Fix:** For the first assertion, inspect grouping output and
require a nonempty collection:
```scala
val groups = collectAllGroupPartitions(df.queryExecution.executedPlan)
assert(groups.nonEmpty)
assert(groups.forall(_.outputPartitioning.numPartitions == expected))
```
`collectAllGroupPartitions` is inherited from the pinned parent and
traverses the whole plan. Alternatively, adapt the join-scoped collector for
both `SortMergeJoinExec` and `SortMergeJoinExecTransformer`; the inherited
join-scoped helper only recognizes vanilla joins.
For SPARK-47094, retain the answer/shuffle checks and follow the [Spark42
branch-specific
oracle](https://github.com/apache/spark/blob/32f7299601108917fb01920a54e084595b7b3bf8/sql/core/src/test/scala/org/apache/spark/sql/connector/KeyGroupedPartitioningSuite.scala#L2386-L2396):
```scala
val partitions = collectAllGroupPartitions(df.queryExecution.executedPlan)
.map(_.outputPartitioning.numPartitions)
(allowPushDown, partiallyClustered) match {
case (true, false) => assert(partitions == Seq(2, 2))
case _ => assert(partitions.isEmpty)
}
```
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]