diff --git a/.github/workflows/build-amd64-releases.yml b/.github/workflows/build-amd64-releases.yml index 9209a221a..82c775f0a 100644 --- a/.github/workflows/build-amd64-releases.yml +++ b/.github/workflows/build-amd64-releases.yml @@ -37,7 +37,7 @@ jobs: runs-on: ${{ matrix.runner }} strategy: matrix: - sparkver: [spark-3.0, spark-3.1, spark-3.2, spark-3.3, spark-3.4, spark-3.5, spark-4.0, spark-4.1] + sparkver: [spark-3.0, spark-3.1, spark-3.2, spark-3.3, spark-3.4, spark-3.5, spark-4.0, spark-4.1, spark-4.2] scalaver: [ 2.12, 2.13 ] javaver: [ 8, 21 ] auronver: [9.0.0-SNAPSHOT] @@ -69,10 +69,14 @@ jobs: scalaver: '2.12' - sparkver: spark-4.1 scalaver: '2.12' + - sparkver: spark-4.2 + scalaver: '2.12' - sparkver: spark-4.0 javaver: '8' - sparkver: spark-4.1 javaver: '8' + - sparkver: spark-4.2 + javaver: '8' # Spark 4.x JDK 21 builds are handled by explicit rockylinux8 include entries below. - sparkver: spark-4.0 scalaver: '2.13' @@ -80,6 +84,9 @@ jobs: - sparkver: spark-4.1 scalaver: '2.13' javaver: '21' + - sparkver: spark-4.2 + scalaver: '2.13' + javaver: '21' include: # Spark 4.x uses rockylinux8 image (with JDK 21) - sparkver: spark-4.0 @@ -94,6 +101,12 @@ jobs: auronver: 9.0.0-SNAPSHOT runner: ubuntu-24.04 image: rockylinux8 + - sparkver: spark-4.2 + scalaver: '2.13' + javaver: '21' + auronver: 9.0.0-SNAPSHOT + runner: ubuntu-24.04 + image: rockylinux8 steps: - uses: actions/checkout@v7 diff --git a/.github/workflows/build-arm-releases.yml b/.github/workflows/build-arm-releases.yml index 3eed52fae..8243eea10 100644 --- a/.github/workflows/build-arm-releases.yml +++ b/.github/workflows/build-arm-releases.yml @@ -37,7 +37,7 @@ jobs: runs-on: ${{ matrix.runner }} strategy: matrix: - sparkver: [spark-3.0, spark-3.1, spark-3.2, spark-3.3, spark-3.4, spark-3.5, spark-4.0, spark-4.1] + sparkver: [spark-3.0, spark-3.1, spark-3.2, spark-3.3, spark-3.4, spark-3.5, spark-4.0, spark-4.1, spark-4.2] auronver: [9.0.0-SNAPSHOT] scalaver: [2.12, 2.13] javaver: [ 8, 21 ] @@ -69,10 +69,14 @@ jobs: scalaver: '2.12' - sparkver: spark-4.1 scalaver: '2.12' + - sparkver: spark-4.2 + scalaver: '2.12' - sparkver: spark-4.0 javaver: '8' - sparkver: spark-4.1 javaver: '8' + - sparkver: spark-4.2 + javaver: '8' steps: - uses: actions/checkout@v7 diff --git a/.github/workflows/build-macos-releases.yml b/.github/workflows/build-macos-releases.yml index 62a54379a..e9cbb334c 100644 --- a/.github/workflows/build-macos-releases.yml +++ b/.github/workflows/build-macos-releases.yml @@ -37,7 +37,7 @@ jobs: runs-on: ${{ matrix.runner }} strategy: matrix: - sparkver: [spark-3.0, spark-3.1, spark-3.2, spark-3.3, spark-3.4, spark-3.5, spark-4.0, spark-4.1] + sparkver: [spark-3.0, spark-3.1, spark-3.2, spark-3.3, spark-3.4, spark-3.5, spark-4.0, spark-4.1, spark-4.2] auronver: [9.0.0-SNAPSHOT] scalaver: [2.12, 2.13] javaver: [8, 21] @@ -69,10 +69,14 @@ jobs: scalaver: '2.12' - sparkver: spark-4.1 scalaver: '2.12' + - sparkver: spark-4.2 + scalaver: '2.12' - sparkver: spark-4.0 javaver: '8' - sparkver: spark-4.1 javaver: '8' + - sparkver: spark-4.2 + javaver: '8' steps: - uses: actions/checkout@v7 diff --git a/.github/workflows/tpcds-reusable.yml b/.github/workflows/tpcds-reusable.yml index ee176d0d1..01ae0bde0 100644 --- a/.github/workflows/tpcds-reusable.yml +++ b/.github/workflows/tpcds-reusable.yml @@ -226,7 +226,7 @@ jobs: if: steps.cache-spark-bin.outputs.cache-hit != 'true' run: | SPARK_PATH="spark/spark-${{ steps.get-dependency-version.outputs.sparkversion }}" - if [[ ${{ inputs.scalaver }} = "2.13" && "${{ inputs.sparkver }}" != "spark-4.0" && "${{ inputs.sparkver }}" != "spark-4.1" ]]; then + if [[ ${{ inputs.scalaver }} = "2.13" && "${{ inputs.sparkver }}" != "spark-4.0" && "${{ inputs.sparkver }}" != "spark-4.1" && "${{ inputs.sparkver }}" != "spark-4.2" ]]; then SPARK_FILE="spark-${{ steps.get-dependency-version.outputs.sparkversion }}-bin-${{ inputs.hadoop-profile }}-scala${{ inputs.scalaver }}.tgz" else SPARK_FILE="spark-${{ steps.get-dependency-version.outputs.sparkversion }}-bin-${{ inputs.hadoop-profile }}.tgz" diff --git a/.github/workflows/tpcds.yml b/.github/workflows/tpcds.yml index d2711dee5..537b4b996 100644 --- a/.github/workflows/tpcds.yml +++ b/.github/workflows/tpcds.yml @@ -112,3 +112,13 @@ jobs: scalaver: '2.13' hadoop-profile: 'hadoop3' sparktests: 'true' + + test-spark-42-jdk21-scala-2-13: + name: Test spark-4.2 JDK21 Scala-2.13 + uses: ./.github/workflows/tpcds-reusable.yml + with: + sparkver: spark-4.2 + javaver: '21' + scalaver: '2.13' + hadoop-profile: 'hadoop3' + sparktests: 'true' diff --git a/auron-build.sh b/auron-build.sh index 523ecf1a1..9f329aa20 100755 --- a/auron-build.sh +++ b/auron-build.sh @@ -30,7 +30,7 @@ # Define constants for supported component versions # ----------------------------------------------------------------------------- SUPPORTED_OS_IMAGES=("centos7" "ubuntu24" "rockylinux8" "debian11" "azurelinux3") -SUPPORTED_SPARK_VERSIONS=("3.0" "3.1" "3.2" "3.3" "3.4" "3.5" "4.0" "4.1") +SUPPORTED_SPARK_VERSIONS=("3.0" "3.1" "3.2" "3.3" "3.4" "3.5" "4.0" "4.1" "4.2") SUPPORTED_SCALA_VERSIONS=("2.12" "2.13") SUPPORTED_CELEBORN_VERSIONS=("0.5" "0.6") # Currently only one supported version, but kept plural for consistency diff --git a/auron-spark-tests/common/src/test/scala/org/apache/spark/sql/SparkExpressionTestsBase.scala b/auron-spark-tests/common/src/test/scala/org/apache/spark/sql/SparkExpressionTestsBase.scala index bb13a6362..5c61ca55f 100644 --- a/auron-spark-tests/common/src/test/scala/org/apache/spark/sql/SparkExpressionTestsBase.scala +++ b/auron-spark-tests/common/src/test/scala/org/apache/spark/sql/SparkExpressionTestsBase.scala @@ -386,7 +386,7 @@ trait SparkExpressionTestsBase Column(expression) } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") private def columnFromExpression(expression: Expression): Column = { new Column(org.apache.spark.sql.classic.ExpressionColumnNode(expression)) } @@ -398,7 +398,7 @@ trait SparkExpressionTestsBase _spark.internalCreateDataFrame(rows, schema) } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") private def internalCreateDataFrame( rows: org.apache.spark.rdd.RDD[InternalRow], schema: StructType): org.apache.spark.sql.classic.DataFrame = { diff --git a/auron-spark-ui/src/main/scala/org/apache/spark/sql/execution/ui/AuronAllExecutionsPage.scala b/auron-spark-ui/src/main/scala/org/apache/spark/sql/execution/ui/AuronAllExecutionsPage.scala index 840abca67..3523d4adc 100644 --- a/auron-spark-ui/src/main/scala/org/apache/spark/sql/execution/ui/AuronAllExecutionsPage.scala +++ b/auron-spark-ui/src/main/scala/org/apache/spark/sql/execution/ui/AuronAllExecutionsPage.scala @@ -32,7 +32,7 @@ private[ui] class AuronAllExecutionsPage(parent: AuronSQLTab) extends WebUIPage( UIUtils.headerSparkPage(request, "Auron", buildInfoSummary(sqlStore.buildInfo()), parent) } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override def render(request: jakarta.servlet.http.HttpServletRequest): Seq[Node] = { UIUtils.headerSparkPage(request, "Auron", buildInfoSummary(sqlStore.buildInfo()), parent) } diff --git a/pom.xml b/pom.xml index 30f27a413..0fbc1b8fa 100644 --- a/pom.xml +++ b/pom.xml @@ -1154,6 +1154,64 @@ + + spark-4.2 + + spark-4.2 + 3.2.9 + 4.2.0 + 1.17.0 + 4.2 + 4.2.13.Final + + + + + org.apache.maven.plugins + maven-enforcer-plugin + ${maven-enforcer-plugin.version} + + + spark42-enforce-java-scala-version + + enforce + + + + + + [17,) + Spark 4.2 requires JDK 17 or higher. Current: ${java.version} + + + scalaLongVersion + 2\.13\.\d+ + Spark 4.2 requires Scala 2.13.x. Current: ${scalaLongVersion} + + + + + + spark42-disallow-iceberg + + enforce + + + + + icebergEnabled + false + Spark 4.2 is not supported for Iceberg. Use Spark 3.4-4.0 or disable icebergEnabled. + + + + + + + + + + jdk-8 diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/InterceptedValidateSparkPlan.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/InterceptedValidateSparkPlan.scala index b61dba4e1..bd89f1c44 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/InterceptedValidateSparkPlan.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/InterceptedValidateSparkPlan.scala @@ -25,7 +25,7 @@ import org.apache.auron.sparkver object InterceptedValidateSparkPlan extends Logging { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") def validate(plan: SparkPlan): Unit = { import org.apache.spark.sql.execution.adaptive.BroadcastQueryStageExec import org.apache.spark.sql.execution.auron.plan.NativeRenameColumnsBase @@ -79,7 +79,7 @@ object InterceptedValidateSparkPlan extends Logging { throw new UnsupportedOperationException("validate is not supported in spark 3.0.3 or 3.1.3") } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def errorOnInvalidBroadcastQueryStage(plan: SparkPlan): Unit = { import org.apache.spark.sql.execution.adaptive.InvalidAQEPlanException throw InvalidAQEPlanException("Invalid broadcast query stage", plan) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/ShimsImpl.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/ShimsImpl.scala index 99a9e088b..bb7ffc86f 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/ShimsImpl.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/auron/ShimsImpl.scala @@ -136,8 +136,10 @@ class ShimsImpl extends Shims with Logging { override def shimVersion: String = "spark-4.0" @sparkver("4.1") override def shimVersion: String = "spark-4.1" + @sparkver("4.2") + override def shimVersion: String = "spark-4.2" - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def initExtension(): Unit = { ValidateSparkPlanInjector.inject() @@ -341,16 +343,22 @@ class ShimsImpl extends Shims with Logging { runtimeFilters: Seq[Expression]): BatchScanExec = exec.copy(exec.output, exec.scan, runtimeFilters, exec.ordering, exec.table, exec.spjParams) - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("4.2") + override def copyBatchScanExecWithRuntimeFilters( + exec: BatchScanExec, + runtimeFilters: Seq[Expression]): BatchScanExec = + exec.copy(exec.output, exec.scan, runtimeFilters, exec.ordering, exec.table) + + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def effectiveLimit(rawLimit: Int): Int = if (rawLimit == -1) Int.MaxValue else rawLimit - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getLimitAndOffset(plan: GlobalLimitExec): (Int, Int) = { (effectiveLimit(plan.limit), plan.offset) } - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getLimitAndOffset(plan: TakeOrderedAndProjectExec): (Int, Int) = { (effectiveLimit(plan.limit), plan.offset) } @@ -364,7 +372,7 @@ class ShimsImpl extends Shims with Logging { override def createNativeLocalLimitExec(limit: Int, child: SparkPlan): NativeLocalLimitBase = NativeLocalLimitExec(limit, child) - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getLimitAndOffset(plan: CollectLimitExec): (Int, Int) = { (effectiveLimit(plan.limit), plan.offset) } @@ -525,7 +533,7 @@ class ShimsImpl extends Shims with Logging { length: Long, numRecords: Long): FileSegment = new FileSegment(file, offset, length) - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def commit( dep: ShuffleDependency[_, _, _], shuffleBlockResolver: IndexShuffleBlockResolver, @@ -732,7 +740,7 @@ class ShimsImpl extends Shims with Logging { expr.asInstanceOf[AggregateExpression].filter } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def isAQEShuffleRead(exec: SparkPlan): Boolean = { import org.apache.spark.sql.execution.adaptive.AQEShuffleReadExec exec.isInstanceOf[AQEShuffleReadExec] @@ -744,7 +752,7 @@ class ShimsImpl extends Shims with Logging { exec.isInstanceOf[CustomShuffleReaderExec] } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def executeNativeAQEShuffleReader(exec: SparkPlan): NativeRDD = { import org.apache.spark.sql.execution.adaptive.AQEShuffleReadExec import org.apache.spark.sql.execution.CoalescedMapperPartitionSpec @@ -1044,7 +1052,7 @@ class ShimsImpl extends Shims with Logging { } } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getSqlContext(sparkPlan: SparkPlan): SQLContext = sparkPlan.session.sqlContext @@ -1066,7 +1074,7 @@ class ShimsImpl extends Shims with Logging { size: Long): PartitionedFile = PartitionedFile(partitionValues, filePath, offset, size) - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getPartitionedFile( partitionValues: InternalRow, filePath: String, @@ -1077,7 +1085,7 @@ class ShimsImpl extends Shims with Logging { PartitionedFile(partitionValues, SparkPath.fromPath(new Path(filePath)), offset, size) } - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getMinPartitionNum(sparkSession: SparkSession): Int = sparkSession.sessionState.conf.filesMinPartitionNum .getOrElse(sparkSession.sparkContext.defaultParallelism) @@ -1100,13 +1108,13 @@ class ShimsImpl extends Shims with Logging { } @nowarn("cat=unused") // Some params temporarily unused - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def convertPromotePrecision( e: Expression, isPruningExpr: Boolean, fallback: Expression => pb.PhysicalExprNode): Option[pb.PhysicalExprNode] = None - @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def convertBloomFilterAgg( agg: AggregateFunction): Option[pb.PhysicalAggExprNode.Builder] = { import org.apache.spark.sql.catalyst.expressions.aggregate.BloomFilterAggregate @@ -1137,7 +1145,7 @@ class ShimsImpl extends Shims with Logging { private def convertBloomFilterAgg( agg: AggregateFunction): Option[pb.PhysicalAggExprNode.Builder] = None - @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") private def convertBloomFilterMightContain( e: Expression, isPruningExpr: Boolean, @@ -1172,7 +1180,7 @@ class ShimsImpl extends Shims with Logging { exec.initialPlan } - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getAdaptiveInputPlan(exec: AdaptiveSparkPlanExec): SparkPlan = { exec.inputPlan } @@ -1202,7 +1210,7 @@ class ShimsImpl extends Shims with Logging { }) } - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getJoinBuildSide(exec: SparkPlan): JoinBuildSide = { import org.apache.spark.sql.catalyst.optimizer.BuildLeft convertJoinBuildSide( @@ -1213,19 +1221,19 @@ class ShimsImpl extends Shims with Logging { }) } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getIsSkewJoinFromSHJ(exec: ShuffledHashJoinExec): Boolean = exec.isSkewJoin @sparkver("3.0 / 3.1") override def getIsSkewJoinFromSHJ(exec: ShuffledHashJoinExec): Boolean = false - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getShuffleOrigin(exec: ShuffleExchangeExec): Option[Any] = Some(exec.shuffleOrigin) @sparkver("3.0") override def getShuffleOrigin(exec: ShuffleExchangeExec): Option[Any] = None - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def isNullAwareAntiJoin(exec: BroadcastHashJoinExec): Boolean = exec.isNullAwareAntiJoin @@ -1236,7 +1244,7 @@ class ShimsImpl extends Shims with Logging { case class ForceNativeExecutionWrapper(override val child: SparkPlan) extends ForceNativeExecutionWrapperBase(child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) @@ -1251,6 +1259,6 @@ case class NativeExprWrapper( override val nullable: Boolean) extends NativeExprWrapperBase(nativeExpr, dataType, nullable) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def withNewChildrenInternal(newChildren: IndexedSeq[Expression]): Expression = copy() } diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/ConvertToNativeExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/ConvertToNativeExec.scala index ead37c9fa..dad518ea4 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/ConvertToNativeExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/ConvertToNativeExec.scala @@ -22,7 +22,7 @@ import org.apache.auron.sparkver case class ConvertToNativeExec(override val child: SparkPlan) extends ConvertToNativeBase(child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeAggExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeAggExec.scala index c72569d5b..3c90819c2 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeAggExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeAggExec.scala @@ -44,22 +44,22 @@ case class NativeAggExec( child) with BaseAggregateExec { - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override val requiredChildDistributionExpressions: Option[Seq[Expression]] = theRequiredChildDistributionExpressions - @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override val initialInputBufferOffset: Int = theInitialInputBufferOffset - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def isStreaming: Boolean = false - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def numShufflePartitions: Option[Int] = None override def resultExpressions: Seq[NamedExpression] = outputAttributes - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeExec.scala index 1f0434cef..655975313 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeExec.scala @@ -43,7 +43,7 @@ case class NativeBroadcastExchangeExec(mode: BroadcastMode, override val child: relationFuturePromise.future } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitExec.scala index 476db05eb..452b1dafc 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeCollectLimitExec.scala @@ -23,7 +23,7 @@ import org.apache.auron.sparkver case class NativeCollectLimitExec(limit: Int, offset: Int, override val child: SparkPlan) extends NativeCollectLimitBase(limit, offset, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeExpandExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeExpandExec.scala index 4b3d221f7..1aba78912 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeExpandExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeExpandExec.scala @@ -28,7 +28,7 @@ case class NativeExpandExec( override val child: SparkPlan) extends NativeExpandBase(projections, output, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeFilterExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeFilterExec.scala index d5335101b..32d6fbbce 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeFilterExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeFilterExec.scala @@ -24,7 +24,7 @@ import org.apache.auron.sparkver case class NativeFilterExec(condition: Expression, override val child: SparkPlan) extends NativeFilterBase(condition, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateExec.scala index 04ad22dbc..26cf0087c 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGenerateExec.scala @@ -30,7 +30,7 @@ case class NativeGenerateExec( override val child: SparkPlan) extends NativeGenerateBase(generator, requiredChildOutput, outer, generatorOutput, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGlobalLimitExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGlobalLimitExec.scala index 4a1407ca7..0412ee6bd 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGlobalLimitExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeGlobalLimitExec.scala @@ -23,7 +23,7 @@ import org.apache.auron.sparkver case class NativeGlobalLimitExec(limit: Int, offset: Int, override val child: SparkPlan) extends NativeGlobalLimitBase(limit, offset, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeLocalLimitExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeLocalLimitExec.scala index bb84b1805..a62a0e724 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeLocalLimitExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeLocalLimitExec.scala @@ -23,7 +23,7 @@ import org.apache.auron.sparkver case class NativeLocalLimitExec(limit: Int, override val child: SparkPlan) extends NativeLocalLimitBase(limit, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcInsertIntoHiveTableExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcInsertIntoHiveTableExec.scala index e364f2434..a2cca528c 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcInsertIntoHiveTableExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcInsertIntoHiveTableExec.scala @@ -69,7 +69,7 @@ case class NativeOrcInsertIntoHiveTableExec( metrics) } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override protected def getInsertIntoHiveTableCommand( table: CatalogTable, partition: Map[String, Option[String]], @@ -88,7 +88,7 @@ case class NativeOrcInsertIntoHiveTableExec( metrics) } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) @@ -294,7 +294,7 @@ case class NativeOrcInsertIntoHiveTableExec( } } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") class AuronInsertIntoHiveTable41( table: CatalogTable, partition: Map[String, Option[String]], diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcSinkExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcSinkExec.scala index 183e4271a..407896150 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcSinkExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeOrcSinkExec.scala @@ -31,7 +31,7 @@ case class NativeOrcSinkExec( override val metrics: Map[String, SQLMetric]) extends NativeOrcSinkBase(sparkSession, table, partition, child, metrics) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetInsertIntoHiveTableExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetInsertIntoHiveTableExec.scala index d7311d447..8871802ef 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetInsertIntoHiveTableExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetInsertIntoHiveTableExec.scala @@ -69,7 +69,7 @@ case class NativeParquetInsertIntoHiveTableExec( metrics) } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override protected def getInsertIntoHiveTableCommand( table: CatalogTable, partition: Map[String, Option[String]], @@ -88,7 +88,7 @@ case class NativeParquetInsertIntoHiveTableExec( metrics) } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) @@ -295,7 +295,7 @@ case class NativeParquetInsertIntoHiveTableExec( } } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") class AuronInsertIntoHiveTable41( table: CatalogTable, partition: Map[String, Option[String]], diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetSinkExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetSinkExec.scala index 78056cbd2..0d997dfb9 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetSinkExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeParquetSinkExec.scala @@ -31,7 +31,7 @@ case class NativeParquetSinkExec( override val metrics: Map[String, SQLMetric]) extends NativeParquetSinkBase(sparkSession, table, partition, child, metrics) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativePartialTakeOrderedExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativePartialTakeOrderedExec.scala index ec1563c80..5106c4472 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativePartialTakeOrderedExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativePartialTakeOrderedExec.scala @@ -29,7 +29,7 @@ case class NativePartialTakeOrderedExec( override val metrics: Map[String, SQLMetric]) extends NativePartialTakeOrderedBase(limit, sortOrder, child, metrics) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeProjectExecProvider.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeProjectExecProvider.scala index e341dffd7..7140d2631 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeProjectExecProvider.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeProjectExecProvider.scala @@ -22,7 +22,7 @@ import org.apache.spark.sql.execution.SparkPlan import org.apache.auron.sparkver case object NativeProjectExecProvider { - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") def provide(projectList: Seq[NamedExpression], child: SparkPlan): NativeProjectBase = { import org.apache.spark.sql.execution.OrderPreservingUnaryExecNode import org.apache.spark.sql.execution.PartitioningPreservingUnaryExecNode diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeRenameColumnsExecProvider.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeRenameColumnsExecProvider.scala index 6e62ba14f..ff6b506af 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeRenameColumnsExecProvider.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeRenameColumnsExecProvider.scala @@ -21,7 +21,7 @@ import org.apache.spark.sql.execution.SparkPlan import org.apache.auron.sparkver case object NativeRenameColumnsExecProvider { - @sparkver("3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.4 / 3.5 / 4.0 / 4.1 / 4.2") def provide(child: SparkPlan, renamedColumnNames: Seq[String]): NativeRenameColumnsBase = { import org.apache.spark.sql.catalyst.expressions.NamedExpression import org.apache.spark.sql.catalyst.expressions.SortOrder diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeShuffleExchangeExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeShuffleExchangeExec.scala index 577973d49..090a22d79 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeShuffleExchangeExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeShuffleExchangeExec.scala @@ -132,7 +132,7 @@ case class NativeShuffleExchangeExec( internalWrite(rdd, dep, mapId, context, partition) } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override def write( inputs: Iterator[_], dep: ShuffleDependency[_, _, _], @@ -195,7 +195,7 @@ case class NativeShuffleExchangeExec( // for databricks testing val causedBroadcastJoinBuildOOM = false - @sparkver("3.5 / 4.0 / 4.1") + @sparkver("3.5 / 4.0 / 4.1 / 4.2") override def advisoryPartitionSize: Option[Long] = None // If users specify the num partitions via APIs like `repartition`, we shouldn't change it. @@ -204,13 +204,13 @@ case class NativeShuffleExchangeExec( override def canChangeNumPartitions: Boolean = outputPartitioning != SinglePartition - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def shuffleOrigin: org.apache.spark.sql.execution.exchange.ShuffleOrigin = { import org.apache.spark.sql.execution.exchange.ShuffleOrigin; _shuffleOrigin.get.asInstanceOf[ShuffleOrigin] } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) @@ -218,7 +218,7 @@ case class NativeShuffleExchangeExec( override def withNewChildren(newChildren: Seq[SparkPlan]): SparkPlan = copy(child = newChildren.head) - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override def shuffleId: Int = { shuffleDependency.shuffleId } diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeSortExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeSortExec.scala index 6d47afb83..03f5fadd2 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeSortExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeSortExec.scala @@ -27,7 +27,7 @@ case class NativeSortExec( override val child: SparkPlan) extends NativeSortBase(sortOrder, global, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeTakeOrderedExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeTakeOrderedExec.scala index 2c939c4f8..f044dece3 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeTakeOrderedExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeTakeOrderedExec.scala @@ -28,7 +28,7 @@ case class NativeTakeOrderedExec( override val child: SparkPlan) extends NativeTakeOrderedBase(limit, offset, sortOrder, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeUnionExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeUnionExec.scala index f4d4ad44d..6b19b20b2 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeUnionExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeUnionExec.scala @@ -26,7 +26,7 @@ case class NativeUnionExec( override val output: Seq[Attribute]) extends NativeUnionBase(children, output) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildrenInternal(newChildren: IndexedSeq[SparkPlan]): SparkPlan = copy(children = newChildren) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeWindowExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeWindowExec.scala index 028c41a29..e74cf12e9 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeWindowExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeWindowExec.scala @@ -31,7 +31,7 @@ case class NativeWindowExec( override val child: SparkPlan) extends NativeWindowBase(windowExpression, partitionSpec, orderSpec, groupLimit, child) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildInternal(newChild: SparkPlan): SparkPlan = copy(child = newChild) diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronBlockStoreShuffleReader.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronBlockStoreShuffleReader.scala index cd8a0ba06..5e2f3519a 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronBlockStoreShuffleReader.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronBlockStoreShuffleReader.scala @@ -41,7 +41,7 @@ class AuronBlockStoreShuffleReader[K, C]( private val _ = mapOutputTracker override def readBlocks(): Iterator[InputStream] = { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") def fetchIterator = new ShuffleBlockFetcherIterator( context, blockManager.blockStoreClient, diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleManager.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleManager.scala index eba8b15fb..501cfd515 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleManager.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleManager.scala @@ -52,7 +52,7 @@ class AuronShuffleManager(conf: SparkConf) extends ShuffleManager with Logging { sortShuffleManager.registerShuffle(shuffleId, dependency) } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getReader[K, C]( handle: ShuffleHandle, startMapIndex: Int, @@ -67,7 +67,7 @@ class AuronShuffleManager(conf: SparkConf) extends ShuffleManager with Logging { @sparkver("3.2") def shuffleMergeFinalized = baseShuffleHandle.dependency.shuffleMergeFinalized - @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") def shuffleMergeFinalized = baseShuffleHandle.dependency.isShuffleMergeFinalizedMarked val (blocksByAddress, canEnableBatchFetch) = diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleWriter.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleWriter.scala index 7e670b71e..2b9d7c5e9 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleWriter.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleWriter.scala @@ -23,6 +23,6 @@ import org.apache.auron.sparkver class AuronShuffleWriter[K, V](metrics: ShuffleWriteMetricsReporter) extends AuronShuffleWriterBase[K, V](metrics) { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def getPartitionLengths(): Array[Long] = partitionLengths } diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeBroadcastJoinExec.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeBroadcastJoinExec.scala index f04770f08..d6be83e72 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeBroadcastJoinExec.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeBroadcastJoinExec.scala @@ -48,7 +48,7 @@ case class NativeBroadcastJoinExec( override val condition: Option[Expression] = None - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def buildSide: org.apache.spark.sql.catalyst.optimizer.BuildSide = broadcastSide match { case JoinBuildLeft => org.apache.spark.sql.catalyst.optimizer.BuildLeft @@ -61,7 +61,7 @@ case class NativeBroadcastJoinExec( case JoinBuildRight => org.apache.spark.sql.execution.joins.BuildRight } - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def requiredChildDistribution : List[org.apache.spark.sql.catalyst.plans.physical.Distribution] = { import org.apache.spark.sql.catalyst.plans.physical.BroadcastDistribution @@ -80,22 +80,22 @@ case class NativeBroadcastJoinExec( override def rewriteKeyExprToLong(exprs: Seq[Expression]): Seq[Expression] = HashJoin.rewriteKeyExpr(exprs) - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def supportCodegen: Boolean = false - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override def inputRDDs(): Nothing = { throw new NotImplementedError("NativeBroadcastJoin dose not support codegen") } - @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.1 / 3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def prepareRelation( ctx: org.apache.spark.sql.catalyst.expressions.codegen.CodegenContext) : org.apache.spark.sql.execution.joins.HashedRelationInfo = { throw new NotImplementedError("NativeBroadcastJoin dose not support codegen") } - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") override protected def withNewChildrenInternal( newLeft: SparkPlan, newRight: SparkPlan): SparkPlan = diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeShuffledHashJoinExecProvider.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeShuffledHashJoinExecProvider.scala index a322f8f71..44bb1ba76 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeShuffledHashJoinExecProvider.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeShuffledHashJoinExecProvider.scala @@ -29,7 +29,7 @@ import org.apache.auron.sparkver case object NativeShuffledHashJoinExecProvider { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") def provide( left: SparkPlan, right: SparkPlan, diff --git a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeSortMergeJoinExecProvider.scala b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeSortMergeJoinExecProvider.scala index 91cf5dba5..7060b486a 100644 --- a/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeSortMergeJoinExecProvider.scala +++ b/spark-extension-shims-spark/src/main/scala/org/apache/spark/sql/execution/joins/auron/plan/NativeSortMergeJoinExecProvider.scala @@ -25,7 +25,7 @@ import org.apache.auron.sparkver case object NativeSortMergeJoinExecProvider { - @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1") + @sparkver("3.2 / 3.3 / 3.4 / 3.5 / 4.0 / 4.1 / 4.2") def provide( left: SparkPlan, right: SparkPlan, diff --git a/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala b/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala index 678faea82..3303cd665 100644 --- a/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala +++ b/spark-extension-shims-spark/src/test/scala/org/apache/spark/sql/AuronQueryTest.scala @@ -19,17 +19,19 @@ package org.apache.spark.sql import org.apache.spark.sql.auron.NativeSupports import org.apache.spark.sql.execution.{LeafExecNode, SparkPlan, UnaryExecNode} import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper -import org.apache.spark.sql.test.SQLTestUtils import org.scalatest.BeforeAndAfterEach +import org.apache.auron.sparkverExcludeParents + /** - * Base test class under org.apache.spark.sql to use package-private [[SQLTestUtils]]; extends - * [[QueryTest]] for comparisons and checks. + * Base test class under org.apache.spark.sql to extends [[QueryTest]] for comparisons and checks. + * Before spark-4.1 also extends package-private [[org.apache.spark.sql.test.SQLTestUtils]]. */ +@sparkverExcludeParents("4.2", "org.apache.spark.sql.test.SQLTestUtils") abstract class AuronQueryTest extends QueryTest - with SQLTestUtils with BeforeAndAfterEach + with org.apache.spark.sql.test.SQLTestUtils with AdaptiveSparkPlanHelper { /** diff --git a/spark-extension/src/main/java/org/apache/auron/spark/configuration/SparkAuronConfiguration.java b/spark-extension/src/main/java/org/apache/auron/spark/configuration/SparkAuronConfiguration.java index 366f32c13..4d3425377 100644 --- a/spark-extension/src/main/java/org/apache/auron/spark/configuration/SparkAuronConfiguration.java +++ b/spark-extension/src/main/java/org/apache/auron/spark/configuration/SparkAuronConfiguration.java @@ -23,8 +23,8 @@ import org.apache.auron.configuration.ConfigOption; import org.apache.spark.SparkContext; import org.apache.spark.SparkEnv; +import org.apache.spark.auron.spark.configurations.ConfigEntryWithDefaultFunction; import org.apache.spark.internal.config.ConfigEntry; -import org.apache.spark.internal.config.ConfigEntryWithDefaultFunction; import org.apache.spark.sql.internal.SQLConf; import scala.Option; import scala.collection.mutable.ListBuffer; diff --git a/spark-extension/src/main/scala/org/apache/spark/auron/spark/configurations/ConfigEntryHelper.scala b/spark-extension/src/main/scala/org/apache/spark/auron/spark/configurations/ConfigEntryHelper.scala new file mode 100644 index 000000000..66e9a041a --- /dev/null +++ b/spark-extension/src/main/scala/org/apache/spark/auron/spark/configurations/ConfigEntryHelper.scala @@ -0,0 +1,52 @@ +/* + * 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.auron.spark.configurations + +import org.apache.spark.internal.config.ConfigEntry +import org.apache.spark.internal.config.ConfigReader + +// org.apache.spark.internal.config.ConfigEntryWithDefaultFunction +private class ConfigEntryWithDefaultFunction[T]( + key: String, + prependedKey: Option[String], + prependSeparator: String, + alternatives: List[String], + _defaultFunction: () => T, + valueConverter: String => T, + stringConverter: T => String, + doc: String, + isPublic: Boolean, + version: String) + extends ConfigEntry( + key, + prependedKey, + prependSeparator, + alternatives, + valueConverter, + stringConverter, + doc, + isPublic, + version) { + + override def defaultValue: Option[T] = Some(_defaultFunction()) + + override def defaultValueString: String = stringConverter(_defaultFunction()) + + def readFrom(reader: ConfigReader): T = { + readString(reader).map(valueConverter).getOrElse(_defaultFunction()) + } +} diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala b/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala index 378a8d662..4a425d681 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/auron/NativeConverters.scala @@ -1276,12 +1276,12 @@ object NativeConverters extends Logging { }) aggBuilder.addChildren(convertExpr(child)) - case CollectList(child, _, _) => + case c: CollectList => aggBuilder.setAggFunction(pb.AggFunction.COLLECT_LIST) - aggBuilder.addChildren(convertExpr(child)) - case CollectSet(child, _, _) => + aggBuilder.addChildren(convertExpr(c.child)) + case c: CollectSet => aggBuilder.setAggFunction(pb.AggFunction.COLLECT_SET) - aggBuilder.addChildren(convertExpr(child)) + aggBuilder.addChildren(convertExpr(c.child)) // brickhouse UDAFs case udaf diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarArray.scala b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarArray.scala index 7d998b0a4..c495bcdec 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarArray.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarArray.scala @@ -166,8 +166,13 @@ class AuronColumnarArray(data: AuronColumnVector, offset: Int, length: Int) exte throw new UnsupportedOperationException } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override def getVariant(i: Int): org.apache.spark.unsafe.types.VariantVal = { throw new UnsupportedOperationException } + + @sparkver("4.2") + override def getBinaryView(i: Int): org.apache.spark.unsafe.types.BinaryView = { + throw new UnsupportedOperationException + } } diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarBatchRow.scala b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarBatchRow.scala index 9908b3466..40c1d19e2 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarBatchRow.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarBatchRow.scala @@ -146,8 +146,13 @@ class AuronColumnarBatchRow(columns: Array[AuronColumnVector], var rowId: Int = throw new UnsupportedOperationException } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override def getVariant(i: Int): org.apache.spark.unsafe.types.VariantVal = { throw new UnsupportedOperationException } + + @sparkver("4.2") + override def getBinaryView(i: Int): org.apache.spark.unsafe.types.BinaryView = { + throw new UnsupportedOperationException + } } diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarStruct.scala b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarStruct.scala index 8cee4d67f..3ac6fce81 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarStruct.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/columnar/AuronColumnarStruct.scala @@ -155,8 +155,13 @@ class AuronColumnarStruct(data: AuronColumnVector, rowId: Int) extends InternalR throw new UnsupportedOperationException } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") override def getVariant(i: Int): org.apache.spark.unsafe.types.VariantVal = { throw new UnsupportedOperationException } + + @sparkver("4.2") + override def getBinaryView(i: Int): org.apache.spark.unsafe.types.BinaryView = { + throw new UnsupportedOperationException + } } diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeBase.scala b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeBase.scala index 9d710a0db..5c2afacd5 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeBase.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/plan/NativeBroadcastExchangeBase.scala @@ -290,7 +290,7 @@ abstract class NativeBroadcastExchangeBase(mode: BroadcastMode, override val chi } } - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") private def getRelationFuture = { SQLExecution.withThreadLocalCaptured[Broadcast[Any]]( this.session.sqlContext.sparkSession, diff --git a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleDependency.scala b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleDependency.scala index cc1c0d5e6..4c5b47c1d 100644 --- a/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleDependency.scala +++ b/spark-extension/src/main/scala/org/apache/spark/sql/execution/auron/shuffle/AuronShuffleDependency.scala @@ -53,7 +53,7 @@ class AuronShuffleDependency[K: ClassTag, V: ClassTag, C: ClassTag]( def getInputRdd: RDD[_ <: Product2[K, V]] = null // For Spark 4+ compatibility: _rdd is required to create NativeRDD.ShuffleWrite in ShuffleWriteProcessor.write - @sparkver("4.0 / 4.1") + @sparkver("4.0 / 4.1 / 4.2") def getInputRdd: RDD[_ <: Product2[K, V]] = _rdd } diff --git a/spark-version-annotation-macros/src/main/scala/org/apache/auron/sparkver.scala b/spark-version-annotation-macros/src/main/scala/org/apache/auron/sparkver.scala index 090c921d9..464ca901d 100644 --- a/spark-version-annotation-macros/src/main/scala/org/apache/auron/sparkver.scala +++ b/spark-version-annotation-macros/src/main/scala/org/apache/auron/sparkver.scala @@ -86,6 +86,40 @@ object sparkver { c.Expr(q"$head; ..${annottees.tail}") } } + + def verExcludeParents(c: whitebox.Context)(annottees: c.Expr[Any]*): c.Expr[Any] = { + import c.universe._ + + val (versions, excludes) = c.macroApplication match { + case Apply(Select(Apply(_, List(vs, ps)), _), _) => + (c.eval(c.Expr[String](q"$vs")), c.eval(c.Expr[String](q"$ps"))) + } + + if (!matchVersion(versions)) { + return c.Expr[Any](q"..$annottees") + } + + val excluded = excludes.split("/").map(_.trim.split('.').last).toSet + + def writtenName(t: Tree): String = t match { + case Ident(name) => name.toString + case Select(qual, name) => writtenName(qual) + "." + name + case Apply(fun, _) => writtenName(fun) + case AppliedTypeTree(tpt, _) => writtenName(tpt) + case _ => t.toString + } + + val head = annottees.head.tree match { + case ClassDef(mods, name, tparams, Template(parents, self, body)) => + val kept = parents.filterNot(p => excluded.contains(writtenName(p).split('.').last)) + ClassDef(mods, name, tparams, Template(kept, self, body)) + case other => + c.abort( + c.enclosingPosition, + s"@sparkverExcludeParents can only annotate a class, got: $other") + } + c.Expr[Any](q"$head; ..${annottees.tail}") + } } } @@ -106,3 +140,9 @@ final class sparkverEnableMembers(vers: String) extends StaticAnnotation { final class sparkverEnableOverride(vers: String) extends StaticAnnotation { def macroTransform(annottees: Any*): Any = macro sparkver.Macros.verEnableOverride } + +@nowarn("cat=unused") // 'vers' and 'parents' are used by macro +@compileTimeOnly("enable macro paradise to expand macro annotations") +final class sparkverExcludeParents(vers: String, parents: String) extends StaticAnnotation { + def macroTransform(annottees: Any*): Any = macro sparkver.Macros.verExcludeParents +}