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]
