This is an automated email from the ASF dual-hosted git repository.
taiyang-li pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/gluten.git
The following commit(s) were added to refs/heads/main by this push:
new 50e28e3f78 [GLUTEN-12998][BOLT] Fix backends-bolt compilation after
the Spark shims cleanup (#12999)
50e28e3f78 is described below
commit 50e28e3f78d02e69427705a7ef4f7ba9dda1fd14
Author: YangJie <[email protected]>
AuthorDate: Wed Sep 16 07:19:47 2026 -0400
[GLUTEN-12998][BOLT] Fix backends-bolt compilation after the Spark shims
cleanup (#12999)
---
.../apache/gluten/backendsapi/bolt/BoltBackend.scala | 14 ++------------
.../gluten/backendsapi/bolt/BoltIteratorApi.scala | 4 ++--
.../gluten/backendsapi/bolt/BoltListenerApi.scala | 3 ---
.../apache/gluten/backendsapi/bolt/BoltRuleApi.scala | 3 ---
.../gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala | 18 ++++++++----------
.../execution/datasources/bolt/BoltBlockStripes.java | 3 +--
6 files changed, 13 insertions(+), 32 deletions(-)
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltBackend.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltBackend.scala
index 7b6fc82a38..c44b4e895b 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltBackend.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltBackend.scala
@@ -25,7 +25,6 @@ import org.apache.gluten.execution.WriteFilesExecTransformer
import org.apache.gluten.expression.WindowFunctionsBuilder
import org.apache.gluten.extension.columnar.cost.{LegacyCoster, LongCoster,
RoughCoster}
import org.apache.gluten.extension.columnar.transition.{Convention,
ConventionFunc}
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.substrait.rel.LocalFilesNode
import org.apache.gluten.substrait.rel.LocalFilesNode.ReadFileFormat
import
org.apache.gluten.substrait.rel.LocalFilesNode.ReadFileFormat.{DwrfReadFormat,
OrcReadFormat, ParquetReadFormat}
@@ -40,8 +39,7 @@ import org.apache.spark.sql.connector.read.Scan
import org.apache.spark.sql.execution.{ColumnarCachedBatchSerializer,
SparkPlan}
import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanExec
import org.apache.spark.sql.execution.columnar.InMemoryTableScanExec
-import
org.apache.spark.sql.execution.command.CreateDataSourceTableAsSelectCommand
-import org.apache.spark.sql.execution.datasources.{FileFormat,
InsertIntoHadoopFsRelationCommand}
+import org.apache.spark.sql.execution.datasources.FileFormat
import org.apache.spark.sql.execution.datasources.parquet.{ParquetFileFormat,
ParquetOptions}
import org.apache.spark.sql.hive.execution.HiveFileFormat
import org.apache.spark.sql.internal.SQLConf
@@ -537,12 +535,6 @@ object BoltBackendSettings extends BackendSettingsApi {
override def insertPostProjectForGenerate(): Boolean = true
- override def skipNativeCtas(ctas: CreateDataSourceTableAsSelectCommand):
Boolean = true
-
- override def skipNativeInsertInto(insertInto:
InsertIntoHadoopFsRelationCommand): Boolean = {
- insertInto.bucketSpec.nonEmpty
- }
-
override def alwaysFailOnMapExpression(): Boolean = true
override def requiredChildOrderingForWindowGroupLimit(): Boolean = false
@@ -550,9 +542,7 @@ object BoltBackendSettings extends BackendSettingsApi {
override def staticPartitionWriteOnly(): Boolean = true
override def enableNativeWriteFiles(): Boolean = {
- GlutenConfig.get.enableNativeWriter.getOrElse(
- SparkShimLoader.getSparkShims.enableNativeWriteFilesByDefault()
- )
+ GlutenConfig.get.enableNativeWriter.getOrElse(true)
}
override def shouldRewriteCount(): Boolean = {
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltIteratorApi.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltIteratorApi.scala
index 10eeac2849..15f99ac184 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltIteratorApi.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltIteratorApi.scala
@@ -26,6 +26,7 @@ import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.substrait.plan.PlanNode
import org.apache.gluten.substrait.rel.{LocalFilesBuilder, LocalFilesNode,
SplitInfo}
import org.apache.gluten.substrait.rel.LocalFilesNode.ReadFileFormat
+import org.apache.gluten.utils.FileMetadataUtil
import org.apache.gluten.vectorized._
import org.apache.spark.{Partition, SparkConf, TaskContext}
@@ -114,8 +115,7 @@ class BoltIteratorApi extends IteratorApi with Logging {
val (fileSizes, modificationTimes) = getFileInfo(partitionFiles).unzip
val partitionColumns = getPartitionColumns(partitionSchema, partitionFiles)
val metadataColumns = partitionFiles
- .map(
- f => SparkShimLoader.getSparkShims.generateMetadataColumns(f,
metadataColumnNames).asJava)
+ .map(f => FileMetadataUtil.generateMetadataColumns(f,
metadataColumnNames).asJava)
val otherMetadataColumns = partitionFiles
.map(f =>
SparkShimLoader.getSparkShims.getOtherConstantMetadataColumnValues(f))
setFileSchemaForLocalFiles(
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltListenerApi.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltListenerApi.scala
index 0eaca4db3c..f6c2f1d67a 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltListenerApi.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltListenerApi.scala
@@ -39,7 +39,6 @@ import org.apache.spark.network.util.ByteUnit
import org.apache.spark.shuffle.{ColumnarShuffleDependency, LookupKey,
ShuffleManagerRegistry}
import org.apache.spark.shuffle.sort.ColumnarShuffleManager
import org.apache.spark.sql.execution.ColumnarCachedBatchSerializer
-import org.apache.spark.sql.execution.datasources.GlutenWriterColumnarRules
import
org.apache.spark.sql.execution.datasources.bolt.{BoltParquetWriterInjects,
BoltRowSplitter}
import org.apache.spark.sql.expression.UDFResolver
import org.apache.spark.sql.internal.{GlutenConfigUtil, StaticSQLConf}
@@ -231,8 +230,6 @@ class BoltListenerApi extends ListenerApi with Logging {
// Inject backend-specific implementations to override spark classes.
GlutenFormatFactory.register(new BoltParquetWriterInjects)
- GlutenFormatFactory.injectPostRuleFactory(
- session => GlutenWriterColumnarRules.NativeWritePostRule(session))
GlutenFormatFactory.register(new BoltRowSplitter())
}
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltRuleApi.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltRuleApi.scala
index 709a592f3f..f1fcf222b2 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltRuleApi.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltRuleApi.scala
@@ -126,9 +126,6 @@ object BoltRuleApi {
// Gluten columnar: Post rules.
injector.injectPost(c => RemoveTopmostColumnarToRow(c.session,
c.caller.isAqe()))
- SparkShimLoader.getSparkShims
- .getExtendedColumnarPostRules()
- .foreach(each => injector.injectPost(c => each(c.session)))
injector.injectPost(c => ColumnarCollapseTransformStages(new
GlutenConfig(c.sqlConf)))
injector.injectPost(c => CudfNodeValidationRule(new
GlutenConfig(c.sqlConf)))
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
index d9d43cc65a..f25f77c336 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/backendsapi/bolt/BoltSparkPlanExecApi.scala
@@ -144,7 +144,7 @@ class BoltSparkPlanExecApi extends SparkPlanExecApi {
right: ExpressionTransformer,
original: Expression,
checkArithmeticExprName: String): ExpressionTransformer = {
- if (SparkShimLoader.getSparkShims.withTryEvalMode(original)) {
+ if (ExpressionUtils.withTryEvalMode(original)) {
original.dataType match {
case LongType | IntegerType | ShortType | ByteType =>
case _ =>
@@ -155,7 +155,7 @@ class BoltSparkPlanExecApi extends SparkPlanExecApi {
ExpressionMappings.expressionsMap(classOf[TryEval]),
Seq(GenericExpressionTransformer(checkArithmeticExprName, Seq(left,
right), original)),
original)
- } else if (SparkShimLoader.getSparkShims.withAnsiEvalMode(original)) {
+ } else if (ExpressionUtils.withAnsiEvalMode(original)) {
GenericExpressionTransformer(checkArithmeticExprName, Seq(left, right),
original)
} else {
GenericExpressionTransformer(substraitExprName, Seq(left, right),
original)
@@ -920,7 +920,7 @@ class BoltSparkPlanExecApi extends SparkPlanExecApi {
substraitExprName: String,
child: ExpressionTransformer,
expr: UnBase64): ExpressionTransformer = {
- if (SparkShimLoader.getSparkShims.unBase64FunctionFailsOnError(expr)) {
+ if (expr.failOnError) {
GlutenExceptionUtil
.throwsNotFullySupported(
ExpressionNames.UNBASE64,
@@ -1135,7 +1135,6 @@ class BoltSparkPlanExecApi extends SparkPlanExecApi {
left: ExpressionTransformer,
right: ExpressionTransformer,
original: Expression): ExpressionTransformer = {
- // Since spark 3.3.0
val extract =
SparkShimLoader.getSparkShims.extractExpressionTimestampAddUnit(original)
if (extract.isEmpty) {
@@ -1149,13 +1148,12 @@ class BoltSparkPlanExecApi extends SparkPlanExecApi {
left: ExpressionTransformer,
right: ExpressionTransformer,
original: Expression): ExpressionTransformer = {
- // Since spark 3.3.0
- val extract =
-
SparkShimLoader.getSparkShims.extractExpressionTimestampDiffUnit(original)
- if (extract.isEmpty) {
- throw new UnsupportedOperationException(s"Not support expression
TimestampDiff.")
+ val unit = original match {
+ case timestampDiff: TimestampDiff => timestampDiff.unit
+ case _ =>
+ throw new UnsupportedOperationException("Not support expression
TimestampDiff.")
}
- TimestampDiffTransformer(substraitExprName, extract.get, left, right,
original)
+ TimestampDiffTransformer(substraitExprName, unit, left, right, original)
}
override def genToUnixTimestampTransformer(
diff --git
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/datasources/bolt/BoltBlockStripes.java
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/datasources/bolt/BoltBlockStripes.java
index 511734486a..c3ef5f491e 100644
---
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/datasources/bolt/BoltBlockStripes.java
+++
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/datasources/bolt/BoltBlockStripes.java
@@ -22,7 +22,6 @@ import org.apache.spark.sql.catalyst.expressions.UnsafeRow;
import org.apache.spark.sql.execution.datasources.BlockStripe;
import org.apache.spark.sql.execution.datasources.BlockStripes;
import org.apache.spark.sql.vectorized.ColumnarBatch;
-import org.jetbrains.annotations.NotNull;
import java.util.Iterator;
@@ -34,7 +33,7 @@ public class BoltBlockStripes extends BlockStripes {
}
@Override
- public @NotNull Iterator<BlockStripe> iterator() {
+ public Iterator<BlockStripe> iterator() {
return new Iterator<BlockStripe>() {
private int index = 0;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]