This is an automated email from the ASF dual-hosted git repository.

baibaichen pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git


The following commit(s) were added to refs/heads/main by this push:
     new 7e309eec6a [GLUTEN-12569][CORE] Fix compatibility issues addressed in 
Spark 4.2 (#13020)
7e309eec6a is described below

commit 7e309eec6a5edfa2d3a4a37c80e238e4ee0b84c1
Author: manoj-ragupathy <[email protected]>
AuthorDate: Tue Sep 22 22:43:49 2026 -0700

    [GLUTEN-12569][CORE] Fix compatibility issues addressed in Spark 4.2 
(#13020)
    
    * [GLUTEN-12569][CORE] Avoid positional patterns broken by Spark 4.2
    
    Spark 4.2 added parameters to three case classes, which silently breaks
    positional extractor patterns:
    
    - CharType gained a collation parameter.
    - AppendDataExec / OverwriteByExpressionExec gained tableName and
      transaction.
    
    Match on type (and on the named field where a value is needed) instead,
    which is stable across all supported Spark versions.
    
    Co-authored-by: Copilot <[email protected]>
    Copilot-Session: a4b4451c-d099-4965-871d-12781176d082
    (cherry picked from commit b2a23c02af068c5aa0f7d37d77dd6563c877f046)
    
    * [GLUTEN-12569][CORE] Make GlutenQueryTest extend Spark's QueryTest
    
    Spark 4.2 hoisted checkAnswer, checkDataset, assertCached and friends out
    of QueryTest into a new QueryTestBase trait, which SharedSparkSession now
    mixes in. Because GlutenQueryTest declared its own copies of those
    members, every suite mixing both inherited conflicting definitions.
    
    Extending QueryTest (an abstract class up to 4.1, a trait in 4.2) makes
    Gluten's versions genuine overrides on every supported Spark version, and
    removes the duplicated-member conflict.
    
    Co-authored-by: Copilot <[email protected]>
    Copilot-Session: a4b4451c-d099-4965-871d-12781176d082
    (cherry picked from commit a0c9f51836f85ecb970603cc6dc63741e772f001)
    
    * [GLUTEN-12569][CORE] Fix Netty and datasketches fallout from Spark 4.2
    
    - Spark 4.2 bumps Netty to 4.2.13, which removed
      PlatformDependent.allocateDirectNoCleaner. Use ByteBuffer.allocateDirect,
      which works on every supported version.
    - Spark 4.2's spark-catalyst ships a patched copy of
      datasketches ResourceImpl, tripping the ban-duplicate-classes enforcer.
      Both artifacts are 'provided', so nothing extra is packaged; ignore it.
    
    Co-authored-by: Copilot <[email protected]>
    Copilot-Session: a4b4451c-d099-4965-871d-12781176d082
    (cherry picked from commit 5d4853dd5ab7298cb909f34e1f6fb3f1e771af38)
    
    * [GLUTEN-12569][CORE] Support SpecializedGetters.getBinaryView for Spark 
4.2
    
    Spark 4.2 adds `getBinaryView` to `SpecializedGetters`, which makes
    `PlaceholderRow` (and any other `InternalRowSparkCompatible` subclass) fail 
to
    compile as non-abstract:
    
        BatchCarrierRow.scala:110: error: class PlaceholderRow needs to be 
abstract.
        Missing implementation for member of trait SpecializedGetters:
          def getBinaryView(x$1: Int): org.apache.spark.unsafe.types.BinaryView
    
    Handled with the existing cross-version seam: 
`SpecializedGettersSparkCompatible`
    gains a `getBinaryView` stub and `InternalRowSparkCompatible` overrides it. 
The
    `Nothing` return type conforms to whichever concrete type the running Spark
    version declares, so one definition covers the whole 3.x/4.x matrix -- 
exactly
    how getVariant/getGeography/getGeometry are already handled.
    
    Verified with clean builds (stale target/ classes previously masked this):
    spark-3.4, spark-3.5, spark-4.0, spark-4.1 and spark-4.2 all BUILD SUCCESS.
    
    (cherry picked from commit 7210d79c11c7ebe11716e874c5e387a434c72823)
    
    * [GLUTEN-12569][CORE] Shim postDriverMetrics for Spark 4.2
    
    Spark 4.2 moved postDriverMetrics to SupportsCustomDriverMetrics and made
    the reported task metrics an explicit argument. Move doPostDriverMetrics
    out of the version-agnostic BatchScanExecTransformer into each
    BatchScanExecShim so the call site can differ per version.
    
    Co-authored-by: Copilot <[email protected]>
    Copilot-Session: 8ae7d6bd-f561-417a-a635-40fe24dc67a5
    
    * [GLUTEN-12569][CORE] Shim SampleExec seed and KeyGroupedPartitioning
    
    Two Spark 4.2 renames that leak into version-agnostic code:
    
    - SampleExec.seed became Option[Long], resolved lazily via resolvedSeed.
    - KeyGroupedPartitioning was renamed to KeyedPartitioning.
    
    Both are now hidden behind SparkShims (getSampleSeed,
    isKeyGroupedPartitioning) so gluten-substrait and backends-velox stay
    version-agnostic.
    
    Co-authored-by: Copilot <[email protected]>
    Copilot-Session: 8ae7d6bd-f561-417a-a635-40fe24dc67a5
    
    ---------
    
    Co-authored-by: MANOJ RAGUPATHY <[email protected]>
    Co-authored-by: Copilot <[email protected]>
    Copilot-Session: a4b4451c-d099-4965-871d-12781176d082
    Copilot-Session: 8ae7d6bd-f561-417a-a635-40fe24dc67a5
---
 .../backendsapi/velox/VeloxSparkPlanExecApi.scala  |  2 +-
 .../org/apache/gluten/fs/OnHeapFileSystemTest.java |  4 +--
 .../execution/BatchScanExecTransformer.scala       |  4 ---
 .../apache/gluten/expression/ConverterUtils.scala  |  2 +-
 .../columnar/offload/OffloadSingleNodeRules.scala  |  2 +-
 .../datasources/noop/GlutenNoopWriterRule.scala    |  6 ++--
 .../org/apache/spark/sql/GlutenQueryTest.scala     | 32 ++++++++++++----------
 package/pom.xml                                    |  4 +++
 .../execution/InternalRowSparkCompatible.scala     |  3 ++
 .../SpecializedGettersSparkCompatible.scala        | 10 ++++++-
 .../org/apache/gluten/sql/shims/SparkShims.scala   | 13 +++++++++
 .../gluten/sql/shims/spark34/Spark34Shims.scala    |  5 ++++
 .../datasources/v2/BatchScanExecShim.scala         |  5 ++++
 .../gluten/sql/shims/spark35/Spark35Shims.scala    |  5 ++++
 .../datasources/v2/BatchScanExecShim.scala         |  5 ++++
 .../gluten/sql/shims/spark40/Spark40Shims.scala    |  5 ++++
 .../datasources/v2/BatchScanExecShim.scala         |  5 ++++
 .../gluten/sql/shims/spark41/Spark41Shims.scala    |  5 ++++
 .../datasources/v2/BatchScanExecShim.scala         |  5 ++++
 19 files changed, 96 insertions(+), 26 deletions(-)

diff --git 
a/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
 
b/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
index 34eb926dfa..43f7ec6a48 100644
--- 
a/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
+++ 
b/backends-velox/src/main/scala/org/apache/gluten/backendsapi/velox/VeloxSparkPlanExecApi.scala
@@ -471,7 +471,7 @@ class VeloxSparkPlanExecApi extends SparkPlanExecApi with 
Logging {
             }
           }
         }
-      case _: KeyGroupedPartitioning =>
+      case p if SparkShimLoader.getSparkShims.isKeyGroupedPartitioning(p) =>
         FallbackTags.add(
           shuffle,
           ValidationResult.failed(
diff --git 
a/backends-velox/src/test/java/org/apache/gluten/fs/OnHeapFileSystemTest.java 
b/backends-velox/src/test/java/org/apache/gluten/fs/OnHeapFileSystemTest.java
index eaa4835de1..b467a67937 100644
--- 
a/backends-velox/src/test/java/org/apache/gluten/fs/OnHeapFileSystemTest.java
+++ 
b/backends-velox/src/test/java/org/apache/gluten/fs/OnHeapFileSystemTest.java
@@ -35,7 +35,7 @@ public class OnHeapFileSystemTest {
     JniFilesystem.WriteFile writeFile = fs.openFileForWrite(path);
     try {
       byte[] bytes = text.getBytes(StandardCharsets.UTF_8);
-      ByteBuffer buf = PlatformDependent.allocateDirectNoCleaner(bytes.length);
+      ByteBuffer buf = ByteBuffer.allocateDirect(bytes.length);
       buf.put(bytes);
       writeFile.append(bytes.length, 
PlatformDependent.directBufferAddress(buf));
       writeFile.flush();
@@ -47,7 +47,7 @@ public class OnHeapFileSystemTest {
 
     JniFilesystem.ReadFile readFile = fs.openFileForRead(path);
     Assert.assertEquals(fileSize, readFile.size());
-    ByteBuffer buf = PlatformDependent.allocateDirectNoCleaner((int) fileSize);
+    ByteBuffer buf = ByteBuffer.allocateDirect((int) fileSize);
     readFile.pread(0, fileSize, PlatformDependent.directBufferAddress(buf));
     byte[] out = new byte[(int) fileSize];
     buf.get(out);
diff --git 
a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala
 
b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala
index 6d2854f30b..a50b252043 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/gluten/execution/BatchScanExecTransformer.scala
@@ -113,10 +113,6 @@ abstract class BatchScanExecTransformerBase(
     BackendsApiManager.getMetricsApiInstance.genBatchScanTransformerMetrics(
       sparkContext) ++ customMetrics
 
-  def doPostDriverMetrics(): Unit = {
-    postDriverMetrics()
-  }
-
   override def scanFilters: Seq[Expression] = scan match {
     case fileScan: FileScan => fileScan.dataFilters
     case _ =>
diff --git 
a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ConverterUtils.scala
 
b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ConverterUtils.scala
index ca83ccbd5b..4c2528b588 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ConverterUtils.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ConverterUtils.scala
@@ -444,7 +444,7 @@ object ConverterUtils extends Logging {
         sigName = sigName.concat(getTypeSigName(valueType))
         sigName = sigName.concat(">")
         sigName
-      case CharType(_) =>
+      case _: CharType =>
         "fchar"
       case NullType =>
         "nothing"
diff --git 
a/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/offload/OffloadSingleNodeRules.scala
 
b/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/offload/OffloadSingleNodeRules.scala
index 61520c6e42..26fe2da4df 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/offload/OffloadSingleNodeRules.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/gluten/extension/columnar/offload/OffloadSingleNodeRules.scala
@@ -319,7 +319,7 @@ object OffloadOthers {
             plan.lowerBound,
             plan.upperBound,
             plan.withReplacement,
-            plan.seed,
+            SparkShimLoader.getSparkShims.getSampleSeed(plan),
             child)
         case plan: RDDScanExec if 
RDDScanTransformer.isSupportRDDScanExec(plan) =>
           RDDScanTransformer.getRDDScanTransform(plan)
diff --git 
a/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/datasources/noop/GlutenNoopWriterRule.scala
 
b/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/datasources/noop/GlutenNoopWriterRule.scala
index bedf006510..db3e6b66da 100644
--- 
a/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/datasources/noop/GlutenNoopWriterRule.scala
+++ 
b/gluten-substrait/src/main/scala/org/apache/spark/sql/execution/datasources/noop/GlutenNoopWriterRule.scala
@@ -33,9 +33,11 @@ import 
org.apache.spark.sql.execution.datasources.v2.{AppendDataExec, OverwriteB
  */
 case class GlutenNoopWriterRule(session: SparkSession) extends Rule[SparkPlan] 
{
   override def apply(p: SparkPlan): SparkPlan = p match {
-    case rc @ AppendDataExec(_, _, NoopWrite) =>
+    // Matched by field rather than positionally: Spark 4.2 added 
`tableName`/`transaction`
+    // parameters to these node types, which breaks positional extraction.
+    case rc: AppendDataExec if rc.write == NoopWrite =>
       injectFakeRowAdaptor(rc, rc.child)
-    case rc @ OverwriteByExpressionExec(_, _, NoopWrite) =>
+    case rc: OverwriteByExpressionExec if rc.write == NoopWrite =>
       injectFakeRowAdaptor(rc, rc.child)
     case _ => p
   }
diff --git 
a/gluten-substrait/src/test/scala/org/apache/spark/sql/GlutenQueryTest.scala 
b/gluten-substrait/src/test/scala/org/apache/spark/sql/GlutenQueryTest.scala
index eb9819998c..1b541b3267 100644
--- a/gluten-substrait/src/test/scala/org/apache/spark/sql/GlutenQueryTest.scala
+++ b/gluten-substrait/src/test/scala/org/apache/spark/sql/GlutenQueryTest.scala
@@ -26,7 +26,6 @@ import org.apache.gluten.sql.shims.SparkShimLoader
 
 import org.apache.spark.{SPARK_VERSION_SHORT, SparkConf}
 import org.apache.spark.sql.catalyst.expressions.Attribute
-import org.apache.spark.sql.catalyst.plans._
 import org.apache.spark.sql.catalyst.plans.logical._
 import org.apache.spark.sql.catalyst.util._
 import org.apache.spark.sql.classic.ClassicConversions._
@@ -45,7 +44,7 @@ import scala.collection.JavaConverters._
 import scala.reflect.ClassTag
 import scala.reflect.runtime.universe
 
-abstract class GlutenQueryTest extends PlanTest with AdaptiveSparkPlanHelper {
+abstract class GlutenQueryTest extends QueryTest with AdaptiveSparkPlanHelper {
 
   // TODO: remove this if we can suppress unused import error.
   locally {
@@ -125,7 +124,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
   }
 
   /** Runs the plan and makes sure the answer contains all of the keywords. */
-  def checkKeywordsExist(df: DataFrame, keywords: String*): Unit = {
+  override def checkKeywordsExist(df: DataFrame, keywords: String*): Unit = {
     val outputs = df.collect().map(_.mkString).mkString
     for (key <- keywords) {
       assert(outputs.contains(key), s"Failed for $df ($key doesn't exist in 
result)")
@@ -133,7 +132,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
   }
 
   /** Runs the plan and makes sure the answer does NOT contain any of the 
keywords. */
-  def checkKeywordsNotExist(df: DataFrame, keywords: String*): Unit = {
+  override def checkKeywordsNotExist(df: DataFrame, keywords: String*): Unit = 
{
     val outputs = df.collect().map(_.mkString).mkString
     for (key <- keywords) {
       assert(!outputs.contains(key), s"Failed for $df ($key existed in the 
result)")
@@ -144,7 +143,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
    * Evaluates a dataset to make sure that the result of calling collect 
matches the given expected
    * answer.
    */
-  protected def checkDataset[T](ds: => Dataset[T], expectedAnswer: T*): Unit = 
{
+  override protected def checkDataset[T](ds: => Dataset[T], expectedAnswer: 
T*): Unit = {
     val result = getResult(ds)
 
     if (!GlutenQueryTest.compare(result.toSeq, expectedAnswer)) {
@@ -161,7 +160,9 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
    * Evaluates a dataset to make sure that the result of calling collect 
matches the given expected
    * answer, after sort.
    */
-  protected def checkDatasetUnorderly[T: Ordering](ds: => Dataset[T], 
expectedAnswer: T*): Unit = {
+  override protected def checkDatasetUnorderly[T: Ordering](
+      ds: => Dataset[T],
+      expectedAnswer: T*): Unit = {
     val result = getResult(ds)
 
     if (!GlutenQueryTest.compare(result.toSeq.sorted, expectedAnswer.sorted)) {
@@ -218,7 +219,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
    * @param expectedAnswer
    *   the expected result in a [[Seq]] of [[Row]]s.
    */
-  protected def checkAnswer(df: => DataFrame, expectedAnswer: Seq[Row]): Unit 
= {
+  override protected def checkAnswer(df: => DataFrame, expectedAnswer: 
Seq[Row]): Unit = {
     val analyzedDF =
       try df
       catch {
@@ -241,11 +242,11 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
     GlutenQueryTest.checkAnswer(analyzedDF, expectedAnswer)
   }
 
-  protected def checkAnswer(df: => DataFrame, expectedAnswer: Row): Unit = {
+  override protected def checkAnswer(df: => DataFrame, expectedAnswer: Row): 
Unit = {
     checkAnswer(df, Seq(expectedAnswer))
   }
 
-  protected def checkAnswer(df: => DataFrame, expectedAnswer: DataFrame): Unit 
= {
+  override protected def checkAnswer(df: => DataFrame, expectedAnswer: 
DataFrame): Unit = {
     checkAnswer(df, expectedAnswer.collect())
   }
 
@@ -259,7 +260,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
    * @param absTol
    *   the absolute tolerance between actual and expected answers.
    */
-  protected def checkAggregatesWithTol(
+  override protected def checkAggregatesWithTol(
       dataFrame: DataFrame,
       expectedAnswer: Seq[Row],
       absTol: Double): Unit = {
@@ -275,7 +276,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
     }
   }
 
-  protected def checkAggregatesWithTol(
+  override protected def checkAggregatesWithTol(
       dataFrame: DataFrame,
       expectedAnswer: Row,
       absTol: Double): Unit = {
@@ -283,7 +284,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
   }
 
   /** Asserts that a given [[Dataset]] will be executed using the given number 
of cached results. */
-  def assertCached(query: Dataset[_], numCachedTables: Int = 1): Unit = {
+  override def assertCached(query: Dataset[_], numCachedTables: Int = 1): Unit 
= {
     val planWithCaching = query.queryExecution.withCachedData
     val cachedData = planWithCaching.collect { case cached: InMemoryRelation 
=> cached }
 
@@ -297,7 +298,10 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
    * Asserts that a given [[Dataset]] will be executed using the cache with 
the given name and
    * storage level.
    */
-  def assertCached(query: Dataset[_], cachedName: String, storageLevel: 
StorageLevel): Unit = {
+  override def assertCached(
+      query: Dataset[_],
+      cachedName: String,
+      storageLevel: StorageLevel): Unit = {
     val planWithCaching = query.queryExecution.withCachedData
     val matched = planWithCaching
       .collectFirst {
@@ -315,7 +319,7 @@ abstract class GlutenQueryTest extends PlanTest with 
AdaptiveSparkPlanHelper {
   }
 
   /** Asserts that a given [[Dataset]] does not have missing inputs in all the 
analyzed plans. */
-  def assertEmptyMissingInput(query: Dataset[_]): Unit = {
+  override def assertEmptyMissingInput(query: Dataset[_]): Unit = {
     assert(
       query.queryExecution.analyzed.missingInput.isEmpty,
       s"The analyzed logical plan has missing 
inputs:\n${query.queryExecution.analyzed}")
diff --git a/package/pom.xml b/package/pom.xml
index 6e3ee88c9c..1ad1c06b5d 100644
--- a/package/pom.xml
+++ b/package/pom.xml
@@ -292,6 +292,10 @@
                     <ignoreClass>org.apache.hadoop.hive.*</ignoreClass>
                     <!-- hawtjni -->
                     
<ignoreClass>org.fusesource.hawtjni.runtime.Library</ignoreClass>
+                    <!-- Spark 4.2's spark-catalyst ships a patched copy of 
this datasketches
+                         class, which collides with datasketches-memory. Both 
artifacts are
+                         `provided`, so neither is packaged into the Gluten 
jar. -->
+                    
<ignoreClass>org.apache.datasketches.memory.internal.ResourceImpl</ignoreClass>
                     <!-- The overridden class list by Gluten. Carefully add 
entries to this list only when you knew exactly what is going to happen -->
                     
<ignoreClass>org.apache.spark.sql.hive.execution.HiveFileFormat</ignoreClass>
                     
<ignoreClass>org.apache.spark.sql.hive.execution.HiveFileFormat$$$$anon$1</ignoreClass>
diff --git 
a/shims/common/src/main/scala/org/apache/gluten/execution/InternalRowSparkCompatible.scala
 
b/shims/common/src/main/scala/org/apache/gluten/execution/InternalRowSparkCompatible.scala
index f1b5c7708d..e4b76e9b38 100644
--- 
a/shims/common/src/main/scala/org/apache/gluten/execution/InternalRowSparkCompatible.scala
+++ 
b/shims/common/src/main/scala/org/apache/gluten/execution/InternalRowSparkCompatible.scala
@@ -31,4 +31,7 @@ abstract class InternalRowSparkCompatible
 
   override def getGeometry(ordinal: Int): Nothing =
     throw new UnsupportedOperationException()
+
+  override def getBinaryView(ordinal: Int): Nothing =
+    throw new UnsupportedOperationException()
 }
diff --git 
a/shims/common/src/main/scala/org/apache/gluten/expression/SpecializedGettersSparkCompatible.scala
 
b/shims/common/src/main/scala/org/apache/gluten/expression/SpecializedGettersSparkCompatible.scala
index a957affed5..fcaad118a2 100644
--- 
a/shims/common/src/main/scala/org/apache/gluten/expression/SpecializedGettersSparkCompatible.scala
+++ 
b/shims/common/src/main/scala/org/apache/gluten/expression/SpecializedGettersSparkCompatible.scala
@@ -21,7 +21,11 @@ package org.apache.gluten.expression
  * compatible with Spark 3.x and 4.x at the same time.
  *
  * Provides stub implementations for methods that exist in Spark 4.x but not 
in Spark 3.x, including
- * getVariant, getGeography, and getGeometry.
+ * getVariant, getGeography, getGeometry and getBinaryView.
+ *
+ * The return type is `Nothing` on purpose: it conforms to whichever concrete 
type the running Spark
+ * version declares (and these types do not all exist in every version), so a 
single definition
+ * works across the whole 3.x/4.x matrix.
  */
 trait SpecializedGettersSparkCompatible {
   def getVariant(ordinal: Int): Nothing = {
@@ -33,4 +37,8 @@ trait SpecializedGettersSparkCompatible {
 
   def getGeometry(ordinal: Int): Nothing =
     throw new UnsupportedOperationException()
+
+  // Added by SpecializedGetters in Spark 4.2.
+  def getBinaryView(ordinal: Int): Nothing =
+    throw new UnsupportedOperationException()
 }
diff --git 
a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala 
b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
index 18fba65f5c..8e41170efc 100644
--- a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
+++ b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
@@ -89,6 +89,19 @@ trait SparkShims {
    */
   def isEmptyRelationExec(plan: SparkPlan): Boolean = false
 
+  /**
+   * The seed of a [[SampleExec]]. Spark 4.2 (SPARK-53564) made the seed 
optional and resolves it
+   * lazily via `resolvedSeed`, so the accessor is shimmed per version.
+   */
+  def getSampleSeed(plan: SampleExec): Long
+
+  /**
+   * Whether the partitioning is Spark's storage-partitioned-join 
partitioning. Spark 4.2
+   * (SPARK-53401) renamed `KeyGroupedPartitioning` to `KeyedPartitioning`, so 
the type test is
+   * shimmed per version.
+   */
+  def isKeyGroupedPartitioning(partitioning: Partitioning): Boolean
+
   def getWindowGroupLimitExecShim(plan: SparkPlan): WindowGroupLimitExecShim = 
null
 
   def getWindowGroupLimitExec(windowGroupLimitExecShim: 
WindowGroupLimitExecShim): SparkPlan = null
diff --git 
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
 
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
index cfd06b6e68..baf3884248 100644
--- 
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
+++ 
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
@@ -49,6 +49,11 @@ import org.apache.parquet.schema.MessageType
 
 class Spark34Shims extends SparkShims {
 
+  override def getSampleSeed(plan: SampleExec): Long = plan.seed
+
+  override def isKeyGroupedPartitioning(partitioning: Partitioning): Boolean =
+    partitioning.isInstanceOf[KeyGroupedPartitioning]
+
   override def scalarExpressionMappings: Seq[Sig] = {
     Seq(
       Sig[Empty2Null](ExpressionNames.EMPTY2NULL),
diff --git 
a/shims/spark34/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
 
b/shims/spark34/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
index 81446cf23a..986401f8b6 100644
--- 
a/shims/spark34/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
+++ 
b/shims/spark34/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
@@ -68,6 +68,11 @@ abstract class BatchScanExecShim(
       .exists(v => metadataColumnsNames.contains(v.name))
   }
 
+  // Spark 4.2 changed the signature of `postDriverMetrics`, so the call is 
shimmed per version.
+  def doPostDriverMetrics(): Unit = {
+    postDriverMetrics()
+  }
+
   override def doExecuteColumnar(): RDD[ColumnarBatch] = {
     throw new UnsupportedOperationException("Need to implement this method")
   }
diff --git 
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
 
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
index 29c5195fbb..b986fd2ce6 100644
--- 
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
+++ 
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
@@ -53,6 +53,11 @@ import scala.collection.JavaConverters._
 
 class Spark35Shims extends SparkShims {
 
+  override def getSampleSeed(plan: SampleExec): Long = plan.seed
+
+  override def isKeyGroupedPartitioning(partitioning: Partitioning): Boolean =
+    partitioning.isInstanceOf[KeyGroupedPartitioning]
+
   override def scalarExpressionMappings: Seq[Sig] = {
     Seq(
       Sig[Empty2Null](ExpressionNames.EMPTY2NULL),
diff --git 
a/shims/spark35/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
 
b/shims/spark35/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
index 0767f5e9f4..7196af63ed 100644
--- 
a/shims/spark35/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
+++ 
b/shims/spark35/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
@@ -70,6 +70,11 @@ abstract class BatchScanExecShim(
       .exists(v => metadataColumnsNames.contains(v.name))
   }
 
+  // Spark 4.2 changed the signature of `postDriverMetrics`, so the call is 
shimmed per version.
+  def doPostDriverMetrics(): Unit = {
+    postDriverMetrics()
+  }
+
   override def doExecuteColumnar(): RDD[ColumnarBatch] = {
     throw new UnsupportedOperationException("Need to implement this method")
   }
diff --git 
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
 
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
index bf9369c6fa..c189367411 100644
--- 
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
+++ 
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
@@ -56,6 +56,11 @@ import scala.jdk.CollectionConverters._
 
 class Spark40Shims extends SparkShims {
 
+  override def getSampleSeed(plan: SampleExec): Long = plan.seed
+
+  override def isKeyGroupedPartitioning(partitioning: Partitioning): Boolean =
+    partitioning.isInstanceOf[KeyGroupedPartitioning]
+
   override def getLocalTableScanStream(plan: LocalTableScanExec): 
Option[SparkDataStream] =
     plan.stream
 
diff --git 
a/shims/spark40/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
 
b/shims/spark40/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
index c00ebe11c0..8a40a8130a 100644
--- 
a/shims/spark40/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
+++ 
b/shims/spark40/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
@@ -76,6 +76,11 @@ abstract class BatchScanExecShim(
       .exists(v => metadataColumnsNames.contains(v.name))
   }
 
+  // Spark 4.2 changed the signature of `postDriverMetrics`, so the call is 
shimmed per version.
+  def doPostDriverMetrics(): Unit = {
+    postDriverMetrics()
+  }
+
   override def doExecuteColumnar(): RDD[ColumnarBatch] = {
     throw new UnsupportedOperationException("Need to implement this method")
   }
diff --git 
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
 
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
index 46ce0ab640..dc0f3e5ae5 100644
--- 
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
+++ 
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
@@ -55,6 +55,11 @@ import scala.jdk.CollectionConverters._
 
 class Spark41Shims extends SparkShims {
 
+  override def getSampleSeed(plan: SampleExec): Long = plan.seed
+
+  override def isKeyGroupedPartitioning(partitioning: Partitioning): Boolean =
+    partitioning.isInstanceOf[KeyGroupedPartitioning]
+
   override def getLocalTableScanStream(plan: LocalTableScanExec): 
Option[SparkDataStream] =
     plan.stream
 
diff --git 
a/shims/spark41/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
 
b/shims/spark41/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
index 604a8c0ee2..d7ab949f47 100644
--- 
a/shims/spark41/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
+++ 
b/shims/spark41/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala
@@ -77,6 +77,11 @@ abstract class BatchScanExecShim(
       .exists(v => metadataColumnsNames.contains(v.name))
   }
 
+  // Spark 4.2 changed the signature of `postDriverMetrics`, so the call is 
shimmed per version.
+  def doPostDriverMetrics(): Unit = {
+    postDriverMetrics()
+  }
+
   override def doExecuteColumnar(): RDD[ColumnarBatch] = {
     throw new UnsupportedOperationException("Need to implement this method")
   }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to