weiting-chen commented on code in PR #13126:
URL: https://github.com/apache/gluten/pull/13126#discussion_r4106168573


##########
shims/spark42/src/main/scala/org/apache/spark/sql/execution/datasources/v2/BatchScanExecShim.scala:
##########
@@ -0,0 +1,155 @@
+/*
+ * 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.execution.datasources.v2
+
+import org.apache.spark.SparkException
+import org.apache.spark.rdd.RDD
+import org.apache.spark.sql.catalyst.InternalRow
+import org.apache.spark.sql.catalyst.expressions._
+import org.apache.spark.sql.catalyst.plans.physical.KeyedPartitioning
+import org.apache.spark.sql.catalyst.util.InternalRowComparableWrapper
+import org.apache.spark.sql.connector.catalog.Table
+import org.apache.spark.sql.connector.catalog.functions.Reducer
+import org.apache.spark.sql.connector.expressions.aggregate.Aggregation
+import org.apache.spark.sql.connector.read.{HasPartitionKey, InputPartition, 
Scan}
+import org.apache.spark.sql.execution.datasources.v2.orc.OrcScan
+import org.apache.spark.sql.execution.datasources.v2.parquet.ParquetScan
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.vectorized.ColumnarBatch
+
+// Spark 4.2 removed `StoragePartitionJoinParams` and no longer accepts the 
SPJ parameters
+// (`joinKeyPositions`, `commonPartitionValues`, `reducers`, 
`applyPartialClustering`,
+// `replicatePartitions`) on the scan node -- that grouping/replication now 
happens in
+// `GroupPartitionsExec`. To keep the public constructor identical to the 
other Spark shims
+// (Gluten's own planner reads these vals), they are kept here as shim-local 
fields and are simply
+// not forwarded into the Spark superclass, which now only takes 
`keyGroupedPartitioning`.
+abstract class BatchScanExecShim(
+    output: Seq[AttributeReference],
+    @transient scan: Scan,
+    runtimeFilters: Seq[Expression],
+    keyGroupedPartitioning: Option[Seq[Expression]] = None,
+    ordering: Option[Seq[SortOrder]] = None,
+    @transient val table: Table,
+    val joinKeyPositions: Option[Seq[Int]] = None,
+    val commonPartitionValues: Option[Seq[(InternalRow, Int)]] = None,
+    val reducers: Option[Seq[Option[Reducer[_, _]]]] = None,
+    val applyPartialClustering: Boolean = false,
+    val replicatePartitions: Boolean = false)
+  extends AbstractBatchScanExec(
+    output,
+    scan,
+    runtimeFilters,
+    ordering,
+    table,
+    keyGroupedPartitioning
+  ) {
+
+  // Note: "metrics" is made transient to avoid sending driver-side metrics to 
tasks.
+  @transient override lazy val metrics: Map[String, SQLMetric] = Map()
+
+  lazy val metadataColumns: Seq[AttributeReference] = output.collect {
+    case FileSourceConstantMetadataAttribute(attr) => attr
+    case FileSourceGeneratedMetadataAttribute(attr, _) => attr
+  }
+
+  def hasUnsupportedColumns: Boolean = {
+    // TODO, fallback if user define same name column due to we can't right now
+    // detect which column is metadata column which is user defined column.
+    val metadataColumnsNames = metadataColumns.map(_.name)
+    output
+      .filterNot(metadataColumns.toSet)
+      .exists(v => metadataColumnsNames.contains(v.name))
+  }
+
+  // Spark 4.2 moved `postDriverMetrics` to SupportsCustomDriverMetrics and 
made the reported
+  // task metrics an explicit argument (see BatchScanExec in Spark 4.2).
+  def doPostDriverMetrics(): Unit = {
+    postDriverMetrics(scan.reportDriverMetrics())
+  }
+
+  override def doExecuteColumnar(): RDD[ColumnarBatch] = {
+    throw new UnsupportedOperationException("Need to implement this method")
+  }
+
+  @transient protected lazy val filteredPartitions: Seq[Seq[InputPartition]] = 
{
+    val originalPartitioning = outputPartitioning
+
+    val filtered = PushDownUtils.pushRuntimeFilters(scan, runtimeFilters, 
table, output)
+    // call toBatch again to get filtered partitions if any runtime filter was 
pushed
+    val newPartitions =
+      if (filtered) scan.toBatch.planInputPartitions().toSeq else 
inputPartitions
+
+    originalPartitioning match {
+      case k: KeyedPartitioning =>
+        if (newPartitions.exists(!_.isInstanceOf[HasPartitionKey])) {
+          throw new SparkException(
+            "Data source must have preserved the original partitioning " +
+              "during runtime filtering: not all partitions implement 
HasPartitionKey after " +
+              "filtering")
+        }
+
+        if (filtered) {
+          // Validate that runtime filtering only removed partition keys, 
never introduced new ones.
+          val newPartitionKeys = newPartitions
+            .map(
+              partition =>
+                InternalRowComparableWrapper(
+                  partition.asInstanceOf[HasPartitionKey].partitionKey(),
+                  k.expressions))
+            .toSet
+          val oldPartitionKeys = k.partitionKeys.toSet
+          // We require the new number of partition keys to be equal or less 
than the old number.
+          if (oldPartitionKeys.size < newPartitionKeys.size) {
+            throw new SparkException(
+              "During runtime filtering, data source must either report " +
+                "the same number of partition values, or a subset of partition 
values from the " +
+                s"original. Before: ${oldPartitionKeys.size} partition values. 
" +
+                s"After: ${newPartitionKeys.size} partition values")
+          }
+          if (!newPartitionKeys.forall(oldPartitionKeys.contains)) {
+            throw new SparkException(
+              "During runtime filtering, data source must not report new " +
+                "partition values that are not present in the original 
partitioning.")
+          }
+        }
+
+        // Group the splits that share the same partition key into a single 
group and sort the
+        // groups by partition key in ascending order. This reproduces the 
key-grouped layout that
+        // Spark 4.1's `BatchScanExec`/`KeyGroupedPartitionedScan` used to 
produce and that Gluten's
+        // planner (`SparkShims.orderPartitions`) still expects. In Spark 4.2 
this grouping is
+        // otherwise deferred to `GroupPartitionsExec`.
+        newPartitions
+          .map(part => (part.asInstanceOf[HasPartitionKey].partitionKey(), 
part))
+          .groupBy { case (key, _) => InternalRowComparableWrapper(key, 
k.expressions) }

Review Comment:
   **Preserve per-split keyed partition slots before enabling native keyed 
scans**
   
   **Target location:** `BatchScanExecShim.scala:130-140`, together with 
`Spark42Shims.orderPartitions:296-309`.
   
   **Problem:** This is a conditional follow-up, not a demonstrated failure of 
the currently supported Spark42 connector configuration. The native adapter 
collapses duplicate-key splits while keeping Spark42's ungrouped 
`KeyedPartitioning` metadata. For input keys `[A,A,B]`, it advertises three 
partitions but creates only two native RDD partitions.
   
   **Evidence:**
   ```scala
   .groupBy { case (key, _) => InternalRowComparableWrapper(key, k.expressions) 
}
   .toSeq
   .sortBy(_._1)(k.keyOrdering)
   .map { case (_, keyedParts) => keyedParts.map(_._2) }
   ```
   The new `orderPartitions` also iterates `p.toGrouped.partitionKeys`, i.e. 
distinct keys. Spark42's [scan 
partitioning](https://github.com/apache/spark/blob/32f7299601108917fb01920a54e084595b7b3bf8/sql/core/src/main/scala/org/apache/spark/sql/execution/datasources/v2/DataSourceV2ScanExecBase.scala#L91-L103)
 retains duplicate keys. `GroupPartitionsExec` derives parent indices `[0,1]` 
and `[2]` from that metadata, so a two-partition native parent cannot satisfy 
index 2. Runtime filtering likewise needs original per-key multiplicity, not 
one empty group per distinct key.
   
   The correct padding in `AbstractBatchScanExec.inputRDD` does not cover this 
path: native `BatchScanExecTransformer.finalPartitions` consumes the protected 
shim groups and passes them through whole-stage native partition construction.
   
   **Suggested Fix:** Before enabling native keyed connectors on Spark42, 
retain one slot per original sorted split, validating filtered per-key counts 
and padding removed splits; let Spark's `GroupPartitionsExec` perform grouping. 
The required layout is:
   ```text
   original [A1,A2,B1]       -> [[A1],[A2],[B1]]
   filter removes A2        -> [[A1],[],[B1]]
   filter removes A1 and A2 -> [[],[],[B1]]
   ```
   Alternatively, explicitly fall back for native keyed scans until 
implemented. Cover these cases under an actually offloaded scan with 
partition-count and result assertions.
   
   Scope qualification: ordinary native `FileScan` partitions do not report 
these keys, and generic keyed DSv2 scans fall back. The optional Iceberg 
transformer can consume this path, but the PR declares its Spark42 dependency 
unavailable and does not enable it. Thus this should not be described as an 
already reproduced supported-runtime crash.



##########
shims/pom.xml:
##########
@@ -95,6 +95,12 @@
         <module>spark41</module>
       </modules>
     </profile>
+    <profile>
+      <id>spark-4.2</id>
+      <modules>
+        <module>spark42</module>
+      </modules>

Review Comment:
   **Register the new shim module in both developer test helpers**
   
   **Target location:** `shims/pom.xml:99-104`; corresponding maps in 
`dev/bloop-test.sh:74-79` and `dev/run-scala-test.sh:161-166`.
   
   **Problem:** The module is registered in Maven but neither developer helper 
recognizes it. Bloop's fallback converts `shims/spark42` to `shims-spark42`, 
not the declared `spark-sql-columnar-shims-spark42` project. The direct Scala 
runner reports `Unknown gluten module in classpath` when its dependency 
classpath contains the Spark42 shim JAR from the Maven repository.
   
   **Evidence:**
   ```xml
   <profile>
     <id>spark-4.2</id>
     <modules>
       <module>spark42</module>
     </modules>
   </profile>
   ```
   Both helper maps stop at Spark41. Already-resolved reactor `target/` 
directories bypass the direct runner's JAR lookup, so this is not a claim that 
every possible invocation fails; direct Maven invocation remains available.
   
   **Suggested Fix:** Add the matching entries:
   ```bash
   # dev/bloop-test.sh MODULE_MAP
   ["shims/spark42"]="spark-sql-columnar-shims-spark42"
   
   # dev/run-scala-test.sh MODULE_MAP
   ["spark-sql-columnar-shims-spark42"]="shims/spark42:java"
   ```
   After fixing the compile issue, exercise the new 
`Spark42LocalTableScanStreamSuite` through these helpers. This only wires the 
shim introduced here; no deferred `gluten-ut/spark42` module is needed.



##########
shims/spark42/src/main/scala/org/apache/gluten/sql/shims/spark42/Spark42Shims.scala:
##########
@@ -0,0 +1,495 @@
+/*
+ * 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.gluten.sql.shims.spark42
+
+import org.apache.gluten.execution.PartitionedFileUtilShim
+import org.apache.gluten.expression.{ExpressionNames, Sig}
+import org.apache.gluten.sql.shims.SparkShims
+
+import org.apache.spark._
+import org.apache.spark.sql.{AnalysisException, SparkSession}
+import org.apache.spark.sql.catalyst.{ExtendedAnalysisException, InternalRow}
+import org.apache.spark.sql.catalyst.analysis.DecimalPrecisionTypeCoercion
+import org.apache.spark.sql.catalyst.expressions._
+import org.apache.spark.sql.catalyst.expressions.aggregate._
+import org.apache.spark.sql.catalyst.plans.{JoinType, LeftSingle}
+import org.apache.spark.sql.catalyst.plans.QueryPlan
+import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
+import org.apache.spark.sql.catalyst.plans.physical.{KeyedPartitioning, 
Partitioning}
+import org.apache.spark.sql.catalyst.types.DataTypeUtils
+import org.apache.spark.sql.catalyst.util.{CollationFactory, 
InternalRowComparableWrapper, MapData}
+import org.apache.spark.sql.catalyst.util.RebaseDateTime.RebaseSpec
+import org.apache.spark.sql.connector.read.{HasPartitionKey, InputPartition, 
Scan}
+import org.apache.spark.sql.connector.read.streaming.SparkDataStream
+import org.apache.spark.sql.execution._
+import org.apache.spark.sql.execution.datasources._
+import org.apache.spark.sql.execution.datasources.parquet.{ParquetFileFormat, 
ParquetFilters}
+import org.apache.spark.sql.execution.datasources.v2.{BatchScanExec, 
DataSourceV2ScanExecBase}
+import org.apache.spark.sql.execution.exchange.{BroadcastExchangeLike, 
ShuffleExchangeLike}
+import org.apache.spark.sql.execution.window.{Final, Partial, _}
+import org.apache.spark.sql.internal.{LegacyBehaviorPolicy, SQLConf}
+import org.apache.spark.sql.types._
+import org.apache.spark.storage.{GlutenShuffleBlockFetcherIterator, 
GlutenShuffleBlockFetcherIteratorBase, ShuffleBlockFetcherIteratorParams}
+
+import org.apache.hadoop.fs.{FileStatus, Path}
+import org.apache.parquet.hadoop.metadata.{CompressionCodecName, 
ParquetMetadata}
+import org.apache.parquet.hadoop.metadata.FileMetaData.EncryptionType
+import org.apache.parquet.schema.{GroupType, LogicalTypeAnnotation, 
MessageType}
+
+import java.util.{Map => JMap}
+
+import scala.jdk.CollectionConverters._
+
+class Spark42Shims extends SparkShims {
+
+  override def getSampleSeed(plan: SampleExec): Long = plan.resolvedSeed
+
+  override def isKeyGroupedPartitioning(partitioning: Partitioning): Boolean =
+    partitioning.isInstanceOf[KeyedPartitioning]
+
+  override def getLocalTableScanStream(plan: LocalTableScanExec): 
Option[SparkDataStream] =
+    plan.stream
+
+  override def scalarExpressionMappings: Seq[Sig] = {
+    Seq(
+      Sig[Empty2Null](ExpressionNames.EMPTY2NULL),
+      Sig[Mask](ExpressionNames.MASK),
+      Sig[ArrayInsert](ExpressionNames.ARRAY_INSERT),
+      
Sig[CheckOverflowInTableInsert](ExpressionNames.CHECK_OVERFLOW_IN_TABLE_INSERT),
+      Sig[ArrayAppend](ExpressionNames.ARRAY_APPEND),
+      Sig[UrlEncode](ExpressionNames.URL_ENCODE),
+      Sig[KnownNotContainsNull](ExpressionNames.KNOWN_NOT_CONTAINS_NULL),
+      Sig[UrlDecode](ExpressionNames.URL_DECODE),
+      Sig[ToPrettyString](ExpressionNames.TO_PRETTY_STRING),
+      Sig[RandStr](ExpressionNames.RANDSTR),
+      Sig[RegExpInStr](ExpressionNames.REGEXP_INSTR),
+      Sig[DayName](ExpressionNames.DAY_NAME),
+      Sig[MonthName](ExpressionNames.MONTH_NAME)
+    )
+  }
+
+  override def aggregateExpressionMappings: Seq[Sig] = {
+    Seq(
+      Sig[RegrSlope](ExpressionNames.REGR_SLOPE),
+      Sig[RegrIntercept](ExpressionNames.REGR_INTERCEPT),
+      Sig[RegrSXY](ExpressionNames.REGR_SXY),
+      Sig[RegrReplacement](ExpressionNames.REGR_REPLACEMENT),
+      Sig[BitmapConstructAgg](ExpressionNames.BITMAP_CONSTRUCT_AGG)
+    )
+  }
+
+  override def runtimeReplaceableExpressionMappings: Seq[Sig] = {
+    Seq(
+      Sig[ArrayCompact](ExpressionNames.ARRAY_COMPACT),
+      Sig[ArrayPrepend](ExpressionNames.ARRAY_PREPEND),
+      Sig[EqualNull](ExpressionNames.EQUAL_NULL),
+      Sig[Get](ExpressionNames.GET),
+      Sig[Luhncheck](ExpressionNames.LUHN_CHECK)
+    )
+  }
+
+  override def isNullIntolerant(expr: Expression): Boolean = 
expr.nullIntolerant
+
+  override def filesGroupedToBuckets(
+      selectedPartitions: Array[PartitionDirectory]): Map[Int, 
Array[PartitionedFile]] = {
+    selectedPartitions
+      .flatMap(p => p.files.map(f => 
PartitionedFileUtilShim.getPartitionedFile(f, p.values)))
+      .groupBy {
+        f =>
+          BucketingUtils
+            .getBucketId(f.toPath.getName)
+            .getOrElse(throw invalidBucketFile(f.urlEncodedPath))
+      }
+  }
+
+  // https://issues.apache.org/jira/browse/SPARK-40400
+  private def invalidBucketFile(path: String): Throwable = {
+    new SparkException(
+      errorClass = "INVALID_BUCKET_FILE",
+      messageParameters = Map("path" -> path),
+      cause = null)
+  }
+
+  override def isWindowGroupLimitExec(plan: SparkPlan): Boolean = plan match {
+    case _: WindowGroupLimitExec => true
+    case _ => false
+  }
+
+  override def isEmptyRelationExec(plan: SparkPlan): Boolean = plan match {
+    case _: EmptyRelationExec => true
+    case _ => false
+  }
+
+  override def getWindowGroupLimitExecShim(plan: SparkPlan): 
WindowGroupLimitExecShim = {
+    val windowGroupLimitPlan = plan.asInstanceOf[WindowGroupLimitExec]
+    val mode = windowGroupLimitPlan.mode match {
+      case Partial => GlutenPartial
+      case Final => GlutenFinal
+    }
+    WindowGroupLimitExecShim(
+      windowGroupLimitPlan.partitionSpec,
+      windowGroupLimitPlan.orderSpec,
+      windowGroupLimitPlan.rankLikeFunction,
+      windowGroupLimitPlan.limit,
+      mode,
+      windowGroupLimitPlan.child
+    )
+  }
+
+  override def getWindowGroupLimitExec(
+      windowGroupLimitExecShim: WindowGroupLimitExecShim): SparkPlan = {
+    val mode = windowGroupLimitExecShim.mode match {
+      case GlutenPartial => Partial
+      case GlutenFinal => Final
+    }
+    WindowGroupLimitExec(
+      windowGroupLimitExecShim.partitionSpec,
+      windowGroupLimitExecShim.orderSpec,
+      windowGroupLimitExecShim.rankLikeFunction,
+      windowGroupLimitExecShim.limit,
+      mode,
+      windowGroupLimitExecShim.child
+    )
+  }
+
+  override def setJobDescriptionOrTagForBroadcastExchange(
+      sc: SparkContext,
+      broadcastExchange: BroadcastExchangeLike): Unit = {
+    // Setup a job tag here so later it may get cancelled by tag if necessary.
+    sc.addJobTag(broadcastExchange.jobTag)
+    sc.setInterruptOnCancel(true)
+  }
+
+  override def cancelJobGroupForBroadcastExchange(
+      sc: SparkContext,
+      broadcastExchange: BroadcastExchangeLike): Unit = {
+    sc.cancelJobsWithTag(broadcastExchange.jobTag)
+  }
+
+  override def getShuffleAdvisoryPartitionSize(shuffle: ShuffleExchangeLike): 
Option[Long] =
+    shuffle.advisoryPartitionSize
+
+  def getFileStatus(partition: PartitionDirectory): Seq[(FileStatus, 
Map[String, Any])] =
+    partition.files.map(f => (f.fileStatus, f.metadata))
+
+  def isFileSplittable(
+      relation: HadoopFsRelation,
+      filePath: Path,
+      sparkSchema: StructType): Boolean = {
+    relation.fileFormat
+      .isSplitable(relation.sparkSession, relation.options, filePath)
+  }
+
+  def isRowIndexMetadataColumn(name: String): Boolean =
+    name == ParquetFileFormat.ROW_INDEX_TEMPORARY_COLUMN_NAME ||
+      name.equalsIgnoreCase("__delta_internal_is_row_deleted")
+
+  def findRowIndexColumnIndexInSchema(sparkSchema: StructType): Int = {
+    sparkSchema.fields.zipWithIndex.find {
+      case (field: StructField, _: Int) =>
+        field.name == ParquetFileFormat.ROW_INDEX_TEMPORARY_COLUMN_NAME
+    } match {
+      case Some((field: StructField, idx: Int)) =>
+        if (field.dataType != LongType && field.dataType != IntegerType) {
+          throw new RuntimeException(
+            s"${ParquetFileFormat.ROW_INDEX_TEMPORARY_COLUMN_NAME} " +
+              "must be of LongType or IntegerType")
+        }
+        idx
+      case _ => -1
+    }
+  }
+
+  def splitFiles(
+      sparkSession: SparkSession,
+      file: FileStatus,
+      filePath: Path,
+      isSplitable: Boolean,
+      maxSplitBytes: Long,
+      partitionValues: InternalRow,
+      metadata: Map[String, Any] = Map.empty): Seq[PartitionedFile] = {
+    PartitionedFileUtilShim.splitFiles(
+      sparkSession,
+      FileStatusWithMetadata(file, metadata),
+      isSplitable,
+      maxSplitBytes,
+      partitionValues)
+  }
+
+  def structFromAttributes(attrs: Seq[Attribute]): StructType = {
+    DataTypeUtils.fromAttributes(attrs)
+  }
+
+  def attributesFromStruct(structType: StructType): Seq[Attribute] = {
+    DataTypeUtils.toAttributes(structType)
+  }
+
+  def getAnalysisExceptionPlan(ae: AnalysisException): Option[LogicalPlan] = {
+    ae match {
+      case eae: ExtendedAnalysisException =>
+        eae.plan
+      case _ =>
+        None
+    }
+  }
+  override def getCommonPartitionValues(
+      batchScan: BatchScanExec): Option[Seq[(InternalRow, Int)]] = {
+    // Spark 4.2 removed `StoragePartitionJoinParams` (and 
`BatchScanExec.spjParams`), so the
+    // "common partition values" that a partially-clustered 
storage-partitioned join used to expose
+    // on the scan node are no longer available here -- Spark 4.2 computes and 
applies them in
+    // `EnsureRequirements`/`GroupPartitionsExec` instead. There is no 
equivalent accessor on the
+    // 4.2 `BatchScanExec`, so we conservatively return `None`, which simply 
disables the
+    // partially-clustered-distribution refinement in Gluten's own scan 
planner (DEGRADED: see the
+    // note in `orderPartitions`). This does not affect the base 
(fully-clustered) SPJ path.
+    None
+  }
+
+  // please ref BatchScanExec::inputRDD
+  override def orderPartitions(
+      batchScan: DataSourceV2ScanExecBase,
+      scan: Scan,
+      keyGroupedPartitioning: Option[Seq[Expression]],
+      filteredPartitions: Seq[Seq[InputPartition]],
+      outputPartitioning: Partitioning,
+      commonPartitionValues: Option[Seq[(InternalRow, Int)]],
+      applyPartialClustering: Boolean,
+      replicatePartitions: Boolean,
+      joinKeyPositions: Option[Seq[Int]] = None): Seq[Seq[InputPartition]] = {
+    scan match {
+      case _ if keyGroupedPartitioning.isDefined =>
+        outputPartitioning match {
+          case p: KeyedPartitioning =>
+            val partExpressions = keyGroupedPartitioning.get
+
+            // DEGRADED (Spark 4.2 port): Spark 4.2 removed 
`KeyGroupedPartitioning` and
+            // `StoragePartitionJoinParams`, and moved the 
storage-partitioned-join refinements that
+            // used to run here into 
`EnsureRequirements`/`GroupPartitionsExec`:
+            //   - subset-of-join-keys projection (`joinKeyPositions`),
+            //   - compatible partition-expression reduction (`reducers`),
+            //   - partially-clustered replication (`commonPartitionValues` /
+            //     `applyPartialClustering` / `replicatePartitions`).
+            // Gluten never populates `joinKeyPositions`/`reducers`, and 
`getCommonPartitionValues`
+            // returns `None` on 4.2, so `commonPartitionValues` is always 
empty here. These
+            // parameters therefore have no 4.2 equivalent that can be 
reproduced on the scan node
+            // and are intentionally NOT applied; only the base key-grouped 
ordering is reproduced.
+            // The base (fully-clustered) SPJ path is unaffected.
+            val groupedPartitions = filteredPartitions.map {
+              splits =>
+                assert(splits.nonEmpty && 
splits.head.isInstanceOf[HasPartitionKey])
+                (splits.head.asInstanceOf[HasPartitionKey].partitionKey(), 
splits)
+            }
+
+            val partitionMapping = groupedPartitions.map {
+              case (partValue, splits) =>
+                InternalRowComparableWrapper(partValue, partExpressions) -> 
splits
+            }.toMap
+
+            // Use the unique, sorted partition keys as the canonical 
partition order (Spark 4.2's
+            // `KeyedPartitioning.toGrouped` returns distinct keys sorted 
ascending), filling absent
+            // keys with empty split groups so both sides of a 
storage-partitioned join stay
+            // aligned. This mirrors the old 
`KeyGroupedPartitioning.uniquePartitionValues` path.
+            p.toGrouped.partitionKeys.map {
+              keyWrapper =>
+                // Use empty partition for those partition values that are not 
present
+                partitionMapping.getOrElse(keyWrapper, Seq.empty)
+            }
+
+          case _ => filteredPartitions
+        }
+      case _ =>
+        filteredPartitions
+    }
+  }
+
+  override def createParquetFilters(
+      conf: SQLConf,
+      schema: MessageType,
+      caseSensitive: Option[Boolean] = None): ParquetFilters = {
+    new ParquetFilters(
+      schema,
+      conf.parquetFilterPushDownDate,
+      conf.parquetFilterPushDownTimestamp,
+      conf.parquetFilterPushDownDecimal,
+      conf.parquetFilterPushDownStringPredicate,
+      conf.parquetFilterPushDownInFilterThreshold,
+      caseSensitive.getOrElse(conf.caseSensitiveAnalysis),
+      RebaseSpec(LegacyBehaviorPolicy.CORRECTED)
+    )
+  }
+
+  override def withOperatorIdMap[T](idMap: java.util.Map[QueryPlan[_], 
Int])(body: => T): T = {
+    val prevIdMap = QueryPlan.localIdMap.get()
+    try {
+      QueryPlan.localIdMap.set(idMap)
+      body
+    } finally {
+      QueryPlan.localIdMap.set(prevIdMap)
+    }
+  }
+
+  override def getOperatorId(plan: QueryPlan[_]): Option[Int] = {
+    Option(QueryPlan.localIdMap.get().get(plan))
+  }
+
+  override def setOperatorId(plan: QueryPlan[_], opId: Int): Unit = {
+    val map = QueryPlan.localIdMap.get()
+    assert(!map.containsKey(plan))
+    map.put(plan, opId)
+  }
+
+  override def unsetOperatorId(plan: QueryPlan[_]): Unit = {
+    QueryPlan.localIdMap.get().remove(plan)
+  }
+
+  override def isParquetFileEncrypted(footer: ParquetMetadata): Boolean = {
+    footer.getFileMetaData.getEncryptionType match {
+      // UNENCRYPTED file has a plaintext footer and no file encryption,
+      // We can leverage file metadata for this check and return unencrypted.
+      case EncryptionType.UNENCRYPTED =>
+        false
+      // PLAINTEXT_FOOTER has a plaintext footer however the file is encrypted.
+      // In such cases, read the footer and use the metadata for encryption 
check.
+      case EncryptionType.PLAINTEXT_FOOTER =>
+        true
+      case _ =>
+        false
+    }
+  }
+
+  override def shouldFallbackForParquetVariantAnnotation(footer: 
ParquetMetadata): Boolean = {
+    if (SQLConf.get.getConf(SQLConf.PARQUET_IGNORE_VARIANT_ANNOTATION)) {
+      false
+    } else {
+      containsVariantAnnotation(footer.getFileMetaData.getSchema)
+    }
+  }
+
+  private def containsVariantAnnotation(groupType: GroupType): Boolean = {
+    groupType.getFields.asScala.exists {
+      field =>
+        Option(field.getLogicalTypeAnnotation)
+          
.exists(_.isInstanceOf[LogicalTypeAnnotation.VariantLogicalTypeAnnotation]) ||
+        (!field.isPrimitive && containsVariantAnnotation(field.asGroupType()))
+    }
+  }
+
+  override def getOtherConstantMetadataColumnValues(file: PartitionedFile): 
JMap[String, Object] =
+    file.otherConstantMetadataColumnValues.asJava.asInstanceOf[JMap[String, 
Object]]
+
+  override def extractExpressionTimestampAddUnit(exp: Expression): 
Option[Seq[String]] = {
+    exp match {
+      // Velox does not support quantity larger than Int.MaxValue.
+      case TimestampAdd(_, LongLiteral(quantity), _, _) if quantity > 
Integer.MAX_VALUE =>
+        Option.empty
+      case timestampAdd: TimestampAdd =>
+        Option.apply(Seq(timestampAdd.unit, 
timestampAdd.timeZoneId.getOrElse("")))
+      case _ => Option.empty
+    }
+  }
+
+  override def widerDecimalType(d1: DecimalType, d2: DecimalType): DecimalType 
= {

Review Comment:
   **Remove the stale override so the new Spark42 profile can compile**
   
   **Target location:** `Spark42Shims.scala:404-406`.
   
   **Problem:** The current source-head `SparkShims` trait no longer declares 
`widerDecimalType`, and `Spark42Shims` has no other superclass contract 
supplying it. This declaration therefore causes `method widerDecimalType 
overrides nothing` when the new profile is compiled. This is a source-head API 
mismatch, not merely an old CI result or target-merge conflict.
   
   **Evidence:**
   ```scala
   override def widerDecimalType(d1: DecimalType, d2: DecimalType): DecimalType 
= {
     DecimalPrecisionTypeCoercion.widerDecimalType(d1, d2)
   }
   ```
   Compare the [current common 
trait](https://github.com/apache/gluten/blob/10054d77c9e020e263f168bad70ce6e6efc72da5/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala):
 the method is absent there and in the existing Spark41 shim. Spark's 
`DecimalPrecisionTypeCoercion` method exists, but that does not make this an 
override.
   
   **Suggested Fix:** Remove this obsolete method and its unused import:
   ```diff
   -  override def widerDecimalType(d1: DecimalType, d2: DecimalType): 
DecimalType = {
   -    DecimalPrecisionTypeCoercion.widerDecimalType(d1, d2)
   -  }
   ```
   Then validate the actual Spark42 profile, for example:
   ```bash
   ./build/mvn -Pspark-4.2 -Pscala-2.13 -Pjava-17 -Pbackends-velox   -pl 
shims/spark42 -am test-compile -DskipTests
   ```
   The green checks currently select older Spark profiles: the new `shims42` 
workflow output has no downstream consumer. A focused Spark42 compile/shim-test 
CI job would catch this without introducing the deferred full UT module. Also 
include `spark-4.2` in `dev/format-scala-code.sh`'s explicit profile list so 
the new module participates in that check.
   
   This finding is based on the explicit current-source override contract; I 
have not executed a local Scala compilation.



##########
shims/spark42/src/main/scala/org/apache/spark/sql/execution/python/BasePythonRunnerShim.scala:
##########
@@ -0,0 +1,65 @@
+/*
+ * 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.execution.python
+
+import org.apache.spark.SparkEnv
+import org.apache.spark.TaskContext
+import org.apache.spark.api.python.{BasePythonRunner, ChainedPythonFunctions, 
PythonWorker}
+import org.apache.spark.sql.execution.metric.SQLMetric
+import org.apache.spark.sql.execution.python.EvalPythonExec.ArgumentMetadata
+import org.apache.spark.sql.vectorized.ColumnarBatch
+
+import java.io.DataOutputStream
+
+abstract class BasePythonRunnerShim(
+    funcs: Seq[(ChainedPythonFunctions, Long)],
+    evalType: Int,
+    argMetas: Array[Array[(Int, Option[String])]],
+    pythonMetrics: Map[String, SQLMetric])
+  extends BasePythonRunner[ColumnarBatch, ColumnarBatch](

Review Comment:
   **Adapt the Arrow Python command framing, not only `writeUDFs`**
   
   **Target location:** `BasePythonRunnerShim.scala:33-38,47-53`.
   
   **Problem:** This new shim connects the existing columnar Arrow writer to 
Spark42's changed worker protocol without adapting the command prefix. Eligible 
nonempty Pandas/Arrow UDF queries select this runner by default 
(`spark.gluten.sql.columnar.arrowUdf=true`); a malformed worker command fails 
execution rather than falling back to Spark.
   
   **Evidence:**
   ```scala
   extends BasePythonRunner[ColumnarBatch, ColumnarBatch](
     funcs.map(_._1),
     evalType,
     argMetas.map(_.map(_._1)),
     None,
     pythonMetrics)
   ```
   Spark42's 
[BasePythonRunner](https://github.com/apache/spark/blob/32f7299601108917fb01920a54e084595b7b3bf8/core/src/main/scala/org/apache/spark/api/python/PythonRunner.scala#L519-L522)
 already emits `evalType -> runnerConf -> evalConf -> writeCommand`. Neither 
this shim nor Gluten's runner overrides the two configuration maps, so both are 
empty.
   
   The shared [Gluten 
writer](https://github.com/apache/gluten/blob/10054d77c9e020e263f168bad70ce6e6efc72da5/backends-velox/src/main/scala/org/apache/spark/api/python/ColumnarArrowEvalPythonExec.scala#L150-L160)
 then still writes `conf.size`, the three configuration pairs, and the old raw 
schema for eval type 101 before the UDF definitions. Spark42's 
[worker](https://github.com/apache/spark/blob/32f7299601108917fb01920a54e084595b7b3bf8/python/pyspark/worker.py#L2493-L2497)
 reads that extra `3` as the number of UDFs and configuration-string bytes as 
argument metadata. Inheriting the new superclass does not repair the extra 
prefix.
   
   **Suggested Fix:** Introduce version-specific framing hooks. The Spark42 
adapter must receive the runner configuration/schema, expose them through the 
new configuration methods, and emit only UDF definitions from `writeCommand`. 
Its contract should match:
   ```scala
   override protected def runnerConf: Map[String, String] =
     super.runnerConf ++ pythonRunnerConf
   
   override protected def evalConf: Map[String, String] =
     if (evalType == 101) super.evalConf + ("input_type" -> schema.json)
     else super.evalConf
   ```
   Keep the legacy prefix only on older Spark versions. For this profile-only 
PR, a safe alternative is to reject Spark42 columnar Arrow-Python offload 
during validation and retain vanilla Spark's `ArrowEvalPythonExec`. Do not just 
disable the existing `arrowUdf` flag: its false branch selects 
`EvalPythonExecTransformer`, not directly vanilla Spark.
   
   A protocol-level check using the pinned worker's argument-parser AST 
reproduces the extra-prefix misdecode; this is not an end-to-end Spark run. 
Add/enable Spark42 scalar Pandas UDF and Arrow-101 tests with executed-plan 
assertions, including a nondefault timezone and worker reuse.



##########
pom.xml:
##########
@@ -1445,6 +1445,94 @@
         </plugins>
       </build>
     </profile>
+    <profile>
+      <id>spark-4.2</id>
+      <properties>
+        <sparkbundle.version>4.2</sparkbundle.version>

Review Comment:
   **Preserve Spark42 bundles across subsequent clean builds**
   
   **Target location:** New `pom.xml:1451` bundle version; corresponding clean 
exclusions are in `package/pom.xml:230-247`.
   
   **Problem:** Building Spark42 and then cleaning/building another Spark 
version in the same checkout removes the Spark42 bundle, unlike the existing 
versioned bundles. This is a nonblocking packaging consistency issue; a single 
Spark42 build is unaffected.
   
   **Evidence:**
   ```xml
   <sparkbundle.version>4.2</sparkbundle.version>
   ```
   The package phase uses this value in `...-bundle-spark4.2_...jar`. Its clean 
plugin sets `excludeDefaultDirectories=true` and selectively cleans `target`, 
preserving only `*spark3.4*`, `*spark3.5*`, `*spark4.0*`, and `*spark4.1*`. 
There is no matching Spark42 exclusion.
   
   **Suggested Fix:** Add alongside the existing package clean exclusions:
   ```xml
   <exclude>*spark4.2*</exclude>
   ```
   A focused check is to build Spark42, run a subsequent Spark41 clean/package 
in the same checkout, and confirm both versioned bundles remain.



##########
.github/workflows/util/install-spark-resources.sh:
##########
@@ -122,6 +122,12 @@ if [[ "${BASH_SOURCE[0]}" == "${0}" ]]; then
       install_spark "4.1.1" "3" "2.12"
       mv /opt/shims/spark41/spark_home/assembly/target/scala-2.12 
/opt/shims/spark41/spark_home/assembly/target/scala-2.13
       ;;
+  4.2)
+      # Spark-4.x, scala 2.12 // using 2.12 as a hack as 4.2 does not have a 
2.13 suffix
+      cd ${INSTALL_DIR} && \
+      install_spark "4.2.0" "3" "2.12"
+      mv /opt/shims/spark42/spark_home/assembly/target/scala-2.12 
/opt/shims/spark42/spark_home/assembly/target/scala-2.13

Review Comment:
   **Honor the selected installation directory in the new Spark42 branch**
   
   **Target location:** `install-spark-resources.sh:125-129`.
   
   **Problem:** The script advertises an optional installation directory, and 
`install_spark` installs beneath `INSTALL_DIR`. With a custom root such as 
`/tmp/spark-resources`, the new branch nevertheless renames a directory under 
`/opt`. It either fails at the final `mv` or targets an unrelated existing 
installation. Default `/opt` use is unaffected; this is a nonblocking 
custom-path defect.
   
   **Evidence:**
   ```bash
   cd ${INSTALL_DIR} && \
   install_spark "4.2.0" "3" "2.12"
   mv /opt/shims/spark42/spark_home/assembly/target/scala-2.12 
/opt/shims/spark42/spark_home/assembly/target/scala-2.13
   ```
   
   **Suggested Fix:**
   ```bash
   mv "${INSTALL_DIR}/shims/spark42/spark_home/assembly/target/scala-2.12" \
      "${INSTALL_DIR}/shims/spark42/spark_home/assembly/target/scala-2.13"
   ```
   I checked the exact branch using side-effect-free command stubs: 
custom-directory routing fails for the current operands and passes with the 
substitution. No actual downloads or filesystem moves were performed.



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