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]

Reply via email to