This is an automated email from the ASF dual-hosted git repository.
jackylee-ch 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 05087f9d83 [GLUTEN-13003][CORE] Move version-independent helpers out
of SparkShims into utils (#13006)
05087f9d83 is described below
commit 05087f9d8388a32f52106196f4ca865b3095e6ca
Author: YangJie <[email protected]>
AuthorDate: Mon Sep 21 01:41:52 2026 -0400
[GLUTEN-13003][CORE] Move version-independent helpers out of SparkShims
into utils (#13006)
---
.../execution/ColumnarPartialGenerateExec.scala | 4 ++--
.../execution/ColumnarPartialProjectExec.scala | 3 +--
.../apache/spark/sql/execution/BroadcastUtils.scala | 8 ++++----
.../sql/execution/ColumnarBuildSideRelation.scala | 6 +++---
.../unsafe/UnsafeColumnarBuildSideRelation.scala | 6 +++---
.../spark/sql/execution/utils/PushDownUtil.scala | 3 ++-
.../execution/ColumnarPartialGenerateExec.scala | 4 ++--
.../execution/ColumnarPartialProjectExec.scala | 3 +--
.../apache/spark/sql/execution/BroadcastUtils.scala | 8 ++++----
.../sql/execution/ColumnarBuildSideRelation.scala | 10 +++++-----
.../unsafe/UnsafeColumnarBuildSideRelation.scala | 10 +++++-----
.../apache/gluten/expression/ExpressionUtils.scala | 12 ++++++++++--
.../org/apache/gluten/sql/shims/SparkShims.scala | 6 +-----
.../org/apache/gluten/utils/ExceptionUtils.scala | 11 +++++++++++
.../gluten/sql/shims/spark34/Spark34Shims.scala | 20 +-------------------
.../gluten/sql/shims/spark35/Spark35Shims.scala | 20 ++------------------
.../gluten/sql/shims/spark40/Spark40Shims.scala | 20 ++------------------
.../gluten/sql/shims/spark41/Spark41Shims.scala | 20 ++------------------
18 files changed, 61 insertions(+), 113 deletions(-)
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
b/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
index 16a5cb95c3..e7336078fc 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
@@ -22,7 +22,7 @@ import org.apache.gluten.expression.InterpretedArrowGenerate
import org.apache.gluten.extension.columnar.transition.Convention
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
-import org.apache.gluten.sql.shims.SparkShimLoader
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.vectorized.{ArrowColumnarRow,
ArrowWritableColumnVector}
import org.apache.spark.rdd.RDD
@@ -61,7 +61,7 @@ case class ColumnarPartialGenerateExec(generateExec:
GenerateExec, child: SparkP
private var hasUnsupportedDataType = false
private val rightSchema =
-
SparkShimLoader.getSparkShims.structFromAttributes(generateExec.generatorOutput)
+ ExpressionUtils.structFromAttributes(generateExec.generatorOutput)
getColumnIndexInChildOutput(
pruneChildAttributes,
diff --git
a/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
b/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
index 81bfdab890..5cbed106d8 100644
---
a/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
+++
b/backends-bolt/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
@@ -23,7 +23,6 @@ import org.apache.gluten.expression.{ArrowProjection,
ExpressionMappings, Expres
import org.apache.gluten.extension.columnar.transition.Convention
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.vectorized.{ArrowColumnarRow,
ArrowWritableColumnVector}
import org.apache.spark.rdd.RDD
@@ -226,7 +225,7 @@ case class ColumnarPartialProjectExec(projectList:
Seq[Expression], child: Spark
c2a += System.currentTimeMillis() - start
val schema =
-
SparkShimLoader.getSparkShims.structFromAttributes(replacedAlias.map(_.toAttribute))
+ ExpressionUtils.structFromAttributes(replacedAlias.map(_.toAttribute))
val vectors: Array[ArrowWritableColumnVector] = ArrowWritableColumnVector
.allocateColumns(numRows, schema)
.map {
diff --git
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
index b588a10e4b..90c301c7e0 100644
---
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
+++
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
@@ -20,7 +20,7 @@ import org.apache.gluten.backendsapi.BackendsApiManager
import org.apache.gluten.columnarbatch.ColumnarBatches
import org.apache.gluten.config.BoltConfig
import org.apache.gluten.runtime.Runtimes
-import org.apache.gluten.sql.shims.SparkShimLoader
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.vectorized.{ColumnarBatchSerializeResult,
ColumnarBatchSerializerJniWrapper}
import org.apache.spark.SparkContext
@@ -101,7 +101,7 @@ object BroadcastUtils {
serializeStream(batchItr()) match {
case ColumnarBatchSerializeResult.EMPTY =>
ColumnarBuildSideRelation(
- SparkShimLoader.getSparkShims.attributesFromStruct(schema),
+ ExpressionUtils.attributesFromStruct(schema),
Array[Array[Byte]](),
mode)
case result: ColumnarBatchSerializeResult =>
@@ -118,12 +118,12 @@ object BroadcastUtils {
bytes
}.toArray
UnsafeColumnarBuildSideRelation(
- SparkShimLoader.getSparkShims.attributesFromStruct(schema),
+ ExpressionUtils.attributesFromStruct(schema),
serialized,
mode)
} else {
ColumnarBuildSideRelation(
- SparkShimLoader.getSparkShims.attributesFromStruct(schema),
+ ExpressionUtils.attributesFromStruct(schema),
result.onHeapData().asScala.toArray,
mode)
}
diff --git
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
index 59a9cb2b00..9a2eb13dff 100644
---
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
+++
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
@@ -21,7 +21,7 @@ import org.apache.gluten.columnarbatch.ColumnarBatches
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
import org.apache.gluten.runtime.Runtimes
-import org.apache.gluten.sql.shims.SparkShimLoader
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.utils.ArrowAbiUtil
import org.apache.gluten.vectorized.{ColumnarBatchSerializerJniWrapper,
NativeColumnarToRowInfo, NativeColumnarToRowJniWrapper}
@@ -98,7 +98,7 @@ case class ColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -149,7 +149,7 @@ case class ColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = serializerJniWrapper.init(cSchema.memoryAddress())
diff --git
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
index 80e92a1537..86e039a67a 100644
---
a/backends-bolt/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
+++
b/backends-bolt/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
@@ -21,7 +21,7 @@ import org.apache.gluten.columnarbatch.ColumnarBatches
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
import org.apache.gluten.runtime.Runtimes
-import org.apache.gluten.sql.shims.SparkShimLoader
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.utils.ArrowAbiUtil
import org.apache.gluten.vectorized.{ColumnarBatchSerializerJniWrapper,
NativeColumnarToRowInfo, NativeColumnarToRowJniWrapper}
@@ -232,7 +232,7 @@ case class UnsafeColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -279,7 +279,7 @@ case class UnsafeColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = serializerJniWrapper.init(cSchema.memoryAddress())
diff --git
a/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/utils/PushDownUtil.scala
b/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/utils/PushDownUtil.scala
index fdad42cabd..a20b82b075 100644
---
a/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/utils/PushDownUtil.scala
+++
b/backends-clickhouse/src/main/scala/org/apache/spark/sql/execution/utils/PushDownUtil.scala
@@ -16,6 +16,7 @@
*/
package org.apache.spark.sql.execution.utils
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.spark.sql.catalyst.expressions.{Attribute, Expression}
@@ -32,7 +33,7 @@ object PushDownUtil {
filter: Expression
): Boolean = {
val schema = new SparkToParquetSchemaConverter(conf).convert(
- SparkShimLoader.getSparkShims.structFromAttributes(output))
+ ExpressionUtils.structFromAttributes(output))
val parquetFilters =
SparkShimLoader.getSparkShims.createParquetFilters(conf, schema)
DataSourceStrategy.translateFilter(filter, supportNestedPredicatePushdown
= true) match {
case Some(sources.StringStartsWith(_, _)) => false
diff --git
a/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
b/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
index 6729e35c6b..f8194a89c0 100644
---
a/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
+++
b/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialGenerateExec.scala
@@ -18,11 +18,11 @@ package org.apache.gluten.execution
import org.apache.gluten.backendsapi.BackendsApiManager
import org.apache.gluten.columnarbatch.{ColumnarBatches, VeloxColumnarBatches}
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.expression.InterpretedArrowGenerate
import org.apache.gluten.extension.columnar.transition.Convention
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.vectorized.{ArrowColumnarBatch, ArrowColumnarRow,
ArrowWritableColumnVector}
import org.apache.spark.rdd.RDD
@@ -61,7 +61,7 @@ case class ColumnarPartialGenerateExec(generateExec:
GenerateExec, child: SparkP
private var hasUnsupportedDataType = false
private val rightSchema =
-
SparkShimLoader.getSparkShims.structFromAttributes(generateExec.generatorOutput)
+ ExpressionUtils.structFromAttributes(generateExec.generatorOutput)
getColumnIndexInChildOutput(
pruneChildAttributes,
diff --git
a/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
b/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
index f2884af5ce..b8f7d34663 100644
---
a/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
+++
b/backends-velox/src/main/scala/org/apache/gluten/execution/ColumnarPartialProjectExec.scala
@@ -23,7 +23,6 @@ import org.apache.gluten.expression.{ArrowProjection,
ConverterUtils, Expression
import org.apache.gluten.extension.columnar.transition.Convention
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.substrait.`type`.TypeBuilder
import org.apache.gluten.substrait.SubstraitContext
import org.apache.gluten.vectorized.{ArrowColumnarRow,
ArrowWritableColumnVector}
@@ -233,7 +232,7 @@ case class ColumnarPartialProjectExec(projectList:
Seq[Expression], child: Spark
c2a += System.currentTimeMillis() - start
val schema =
-
SparkShimLoader.getSparkShims.structFromAttributes(replacedAlias.map(_.toAttribute))
+ ExpressionUtils.structFromAttributes(replacedAlias.map(_.toAttribute))
val vectors: Array[ArrowWritableColumnVector] = ArrowWritableColumnVector
.allocateColumns(numRows, schema)
.map {
diff --git
a/backends-velox/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
b/backends-velox/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
index 9c092e1070..3746f9529b 100644
---
a/backends-velox/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
+++
b/backends-velox/src/main/scala/org/apache/spark/sql/execution/BroadcastUtils.scala
@@ -19,8 +19,8 @@ package org.apache.spark.sql.execution
import org.apache.gluten.backendsapi.BackendsApiManager
import org.apache.gluten.columnarbatch.ColumnarBatches
import org.apache.gluten.config.VeloxConfig
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.runtime.Runtimes
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.vectorized.{ColumnarBatchSerializeResult,
ColumnarBatchSerializerJniWrapper}
import org.apache.spark.SparkContext
@@ -100,20 +100,20 @@ object BroadcastUtils {
serializeStream(batchItr()) match {
case ColumnarBatchSerializeResult.EMPTY =>
ColumnarBuildSideRelation(
- SparkShimLoader.getSparkShims.attributesFromStruct(schema),
+ ExpressionUtils.attributesFromStruct(schema),
Array[Array[Byte]](),
mode)
case result: ColumnarBatchSerializeResult =>
if (result.isOffHeap) {
UnsafeColumnarBuildSideRelation(
- SparkShimLoader.getSparkShims.attributesFromStruct(schema),
+ ExpressionUtils.attributesFromStruct(schema),
result.offHeapData().asScala.toSeq,
mode,
Seq.empty,
result.isOffHeap)
} else {
ColumnarBuildSideRelation(
- SparkShimLoader.getSparkShims.attributesFromStruct(schema),
+ ExpressionUtils.attributesFromStruct(schema),
result.onHeapData().asScala.toArray,
mode)
}
diff --git
a/backends-velox/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
b/backends-velox/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
index dc1be02d74..55616fc8a9 100644
---
a/backends-velox/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
+++
b/backends-velox/src/main/scala/org/apache/spark/sql/execution/ColumnarBuildSideRelation.scala
@@ -21,10 +21,10 @@ import org.apache.gluten.columnarbatch.ColumnarBatches
import org.apache.gluten.config.GlutenConfig
import org.apache.gluten.execution.BroadcastHashJoinContext
import org.apache.gluten.expression.ConverterUtils
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
import org.apache.gluten.runtime.Runtimes
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.utils.{ArrowAbiUtil, SubstraitUtil}
import org.apache.gluten.vectorized.{ColumnarBatchSerializerJniWrapper,
HashJoinBuilder, NativeColumnarToRowInfo, NativeColumnarToRowJniWrapper}
@@ -132,7 +132,7 @@ case class ColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -185,7 +185,7 @@ case class ColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -280,7 +280,7 @@ case class ColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.globalInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -377,7 +377,7 @@ case class ColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = serializerJniWrapper.init(cSchema.memoryAddress())
diff --git
a/backends-velox/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
b/backends-velox/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
index 1c2adbc57f..e0f44d350f 100644
---
a/backends-velox/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
+++
b/backends-velox/src/main/scala/org/apache/spark/sql/execution/unsafe/UnsafeColumnarBuildSideRelation.scala
@@ -21,10 +21,10 @@ import org.apache.gluten.columnarbatch.ColumnarBatches
import org.apache.gluten.config.GlutenConfig
import org.apache.gluten.execution.BroadcastHashJoinContext
import org.apache.gluten.expression.ConverterUtils
+import org.apache.gluten.expression.ExpressionUtils
import org.apache.gluten.iterator.Iterators
import org.apache.gluten.memory.arrow.alloc.ArrowBufferAllocators
import org.apache.gluten.runtime.Runtimes
-import org.apache.gluten.sql.shims.SparkShimLoader
import org.apache.gluten.utils.{ArrowAbiUtil, SubstraitUtil}
import org.apache.gluten.vectorized.{ColumnarBatchSerializerJniWrapper,
HashJoinBuilder, NativeColumnarToRowInfo, NativeColumnarToRowJniWrapper}
@@ -144,7 +144,7 @@ class UnsafeColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -240,7 +240,7 @@ class UnsafeColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.globalInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -393,7 +393,7 @@ class UnsafeColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = jniWrapper
@@ -442,7 +442,7 @@ class UnsafeColumnarBuildSideRelation(
val allocator = ArrowBufferAllocators.contextInstance()
val cSchema = ArrowSchema.allocateNew(allocator)
val arrowSchema = SparkArrowUtil.toArrowSchema(
- SparkShimLoader.getSparkShims.structFromAttributes(output),
+ ExpressionUtils.structFromAttributes(output),
SQLConf.get.sessionLocalTimeZone)
ArrowAbiUtil.exportSchema(allocator, arrowSchema, cSchema)
val handle = serializerJniWrapper.init(cSchema.memoryAddress())
diff --git
a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionUtils.scala
b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionUtils.scala
index 5aa7fb49ed..575d611b2a 100644
---
a/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionUtils.scala
+++
b/gluten-substrait/src/main/scala/org/apache/gluten/expression/ExpressionUtils.scala
@@ -16,12 +16,20 @@
*/
package org.apache.gluten.expression
-import org.apache.spark.sql.catalyst.expressions.{Add, Cast, Divide, EvalMode,
Expression, IntegralDivide, LeafExpression, Multiply, Subtract}
+import org.apache.spark.sql.catalyst.expressions.{Add, Attribute,
AttributeReference, Cast, Divide, EvalMode, Expression, IntegralDivide,
LeafExpression, Multiply, Subtract}
import org.apache.spark.sql.execution.SparkPlan
-import org.apache.spark.sql.types.{ArrayType, DataType, MapType, StructType}
+import org.apache.spark.sql.types.{ArrayType, DataType, MapType, StructField,
StructType}
object ExpressionUtils {
+ def structFromAttributes(attrs: Seq[Attribute]): StructType =
+ StructType(attrs.map(a => StructField(a.name, a.dataType, a.nullable,
a.metadata)))
+
+ def attributesFromStruct(structType: StructType): Seq[Attribute] =
+ structType.fields.map {
+ field => AttributeReference(field.name, field.dataType, field.nullable,
field.metadata)()
+ }
+
private def getExpressionTreeDepth(expr: Expression): Integer = {
if (expr.isInstanceOf[LeafExpression]) {
return 0
diff --git
a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
index d7348f233d..18fba65f5c 100644
--- a/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
+++ b/shims/common/src/main/scala/org/apache/gluten/sql/shims/SparkShims.scala
@@ -22,7 +22,7 @@ import org.apache.gluten.expression.Sig
import org.apache.spark.{SparkContext, SparkException}
import org.apache.spark.sql.{AnalysisException, SparkSession}
import org.apache.spark.sql.catalyst.InternalRow
-import org.apache.spark.sql.catalyst.expressions.{Attribute, BinaryArithmetic,
Expression, RaiseError}
+import org.apache.spark.sql.catalyst.expressions.{BinaryArithmetic,
Expression, RaiseError}
import org.apache.spark.sql.catalyst.plans.JoinType
import org.apache.spark.sql.catalyst.plans.QueryPlan
import org.apache.spark.sql.catalyst.plans.logical.LogicalPlan
@@ -122,10 +122,6 @@ trait SparkShims {
partitionValues: InternalRow,
metadata: Map[String, Any] = Map.empty): Seq[PartitionedFile]
- def structFromAttributes(attrs: Seq[Attribute]): StructType
-
- def attributesFromStruct(structType: StructType): Seq[Attribute]
-
// For compatibility with Spark-3.5.
def getAnalysisExceptionPlan(ae: AnalysisException): Option[LogicalPlan]
diff --git
a/shims/common/src/main/scala/org/apache/gluten/utils/ExceptionUtils.scala
b/shims/common/src/main/scala/org/apache/gluten/utils/ExceptionUtils.scala
index f86071186d..6f2b21ac1e 100644
--- a/shims/common/src/main/scala/org/apache/gluten/utils/ExceptionUtils.scala
+++ b/shims/common/src/main/scala/org/apache/gluten/utils/ExceptionUtils.scala
@@ -16,8 +16,19 @@
*/
package org.apache.gluten.utils
+import org.apache.spark.SparkException
+
object ExceptionUtils {
+ // Spark's QueryExecutionErrors.invalidBucketFile is private[sql] and cannot
be reached from
+ // this package, so build the same exception here.
+ // https://issues.apache.org/jira/browse/SPARK-40400
+ def invalidBucketFile(path: String): Throwable =
+ new SparkException(
+ errorClass = "INVALID_BUCKET_FILE",
+ messageParameters = Map("path" -> path),
+ cause = null)
+
/**
* Utility to check the exception for the specified type.
*
diff --git
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
index b16361a087..cfd06b6e68 100644
---
a/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
+++
b/shims/spark34/src/main/scala/org/apache/gluten/sql/shims/spark34/Spark34Shims.scala
@@ -91,18 +91,10 @@ class Spark34Shims extends SparkShims {
f =>
BucketingUtils
.getBucketId(f.toPath.getName)
- .getOrElse(throw invalidBucketFile(f.urlEncodedPath))
+ .getOrElse(throw
ExceptionUtils.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)
- }
-
def setJobDescriptionOrTagForBroadcastExchange(
sc: SparkContext,
broadcastExchange: BroadcastExchangeLike): Unit = {
@@ -170,16 +162,6 @@ class Spark34Shims extends SparkShims {
partitionValues)
}
- def structFromAttributes(attrs: Seq[Attribute]): StructType = {
- StructType(attrs.map(a => StructField(a.name, a.dataType, a.nullable,
a.metadata)))
- }
-
- def attributesFromStruct(structType: StructType): Seq[Attribute] = {
- structType.fields.map {
- field => AttributeReference(field.name, field.dataType, field.nullable,
field.metadata)()
- }
- }
-
def getAnalysisExceptionPlan(ae: AnalysisException): Option[LogicalPlan] = {
ae.plan
}
diff --git
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
index 585ba18aeb..29c5195fbb 100644
---
a/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
+++
b/shims/spark35/src/main/scala/org/apache/gluten/sql/shims/spark35/Spark35Shims.scala
@@ -19,6 +19,7 @@ package org.apache.gluten.sql.shims.spark35
import org.apache.gluten.execution.PartitionedFileUtilShim
import org.apache.gluten.expression.{ExpressionNames, Sig}
import org.apache.gluten.sql.shims.SparkShims
+import org.apache.gluten.utils.ExceptionUtils
import org.apache.spark._
import org.apache.spark.sql.{AnalysisException, SparkSession}
@@ -28,7 +29,6 @@ import org.apache.spark.sql.catalyst.expressions.aggregate._
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.{KeyGroupedPartitioning,
Partitioning}
-import org.apache.spark.sql.catalyst.types.DataTypeUtils
import org.apache.spark.sql.catalyst.util.InternalRowComparableWrapper
import org.apache.spark.sql.catalyst.util.RebaseDateTime.RebaseSpec
import org.apache.spark.sql.connector.read.{HasPartitionKey, InputPartition,
Scan}
@@ -97,18 +97,10 @@ class Spark35Shims extends SparkShims {
f =>
BucketingUtils
.getBucketId(f.toPath.getName)
- .getOrElse(throw invalidBucketFile(f.urlEncodedPath))
+ .getOrElse(throw
ExceptionUtils.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
@@ -209,14 +201,6 @@ class Spark35Shims extends SparkShims {
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 =>
diff --git
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
index af563e0bc9..bf9369c6fa 100644
---
a/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
+++
b/shims/spark40/src/main/scala/org/apache/gluten/sql/shims/spark40/Spark40Shims.scala
@@ -19,6 +19,7 @@ package org.apache.gluten.sql.shims.spark40
import org.apache.gluten.execution.PartitionedFileUtilShim
import org.apache.gluten.expression.{ExpressionNames, Sig}
import org.apache.gluten.sql.shims.SparkShims
+import org.apache.gluten.utils.ExceptionUtils
import org.apache.spark._
import org.apache.spark.sql.{AnalysisException, SparkSession}
@@ -29,7 +30,6 @@ 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.{KeyGroupedPartitioning,
KeyGroupedShuffleSpec, 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.classic.ClassicConversions._
@@ -107,18 +107,10 @@ class Spark40Shims extends SparkShims {
f =>
BucketingUtils
.getBucketId(f.toPath.getName)
- .getOrElse(throw invalidBucketFile(f.urlEncodedPath))
+ .getOrElse(throw
ExceptionUtils.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
@@ -225,14 +217,6 @@ class Spark40Shims extends SparkShims {
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 =>
diff --git
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
index faf6bdd088..46ce0ab640 100644
---
a/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
+++
b/shims/spark41/src/main/scala/org/apache/gluten/sql/shims/spark41/Spark41Shims.scala
@@ -19,6 +19,7 @@ package org.apache.gluten.sql.shims.spark41
import org.apache.gluten.execution.PartitionedFileUtilShim
import org.apache.gluten.expression.{ExpressionNames, Sig}
import org.apache.gluten.sql.shims.SparkShims
+import org.apache.gluten.utils.ExceptionUtils
import org.apache.spark._
import org.apache.spark.sql.{AnalysisException, SparkSession}
@@ -29,7 +30,6 @@ 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.{KeyGroupedPartitioning,
KeyGroupedShuffleSpec, 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}
@@ -106,18 +106,10 @@ class Spark41Shims extends SparkShims {
f =>
BucketingUtils
.getBucketId(f.toPath.getName)
- .getOrElse(throw invalidBucketFile(f.urlEncodedPath))
+ .getOrElse(throw
ExceptionUtils.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
@@ -224,14 +216,6 @@ class Spark41Shims extends SparkShims {
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 =>
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]