zhouyuan commented on code in PR #12962:
URL: https://github.com/apache/gluten/pull/12962#discussion_r3926559088
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -42,8 +42,13 @@ object GlutenIcebergSourceUtil {
private val InputFileBlockStartCol = "input_file_block_start"
private val InputFileBlockLengthCol = "input_file_block_length"
- def getClassOfSparkBatchQueryScan(): Class[SparkBatchQueryScan] = {
- classOf[SparkBatchQueryScan]
+ def supportsScan(sparkScan: Scan): Boolean = sparkScan match {
+ case _: SparkBatchQueryScan => true
+ case scan: SparkStagedScan =>
+ val tasks = getScanTasks(scan)
Review Comment:
this will bring big perf overhead here, can we avoid this?
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -42,8 +42,13 @@ object GlutenIcebergSourceUtil {
private val InputFileBlockStartCol = "input_file_block_start"
private val InputFileBlockLengthCol = "input_file_block_length"
- def getClassOfSparkBatchQueryScan(): Class[SparkBatchQueryScan] = {
- classOf[SparkBatchQueryScan]
+ def supportsScan(sparkScan: Scan): Boolean = sparkScan match {
+ case _: SparkBatchQueryScan => true
+ case scan: SparkStagedScan =>
+ val tasks = getScanTasks(scan)
+ tasks.nonEmpty &&
Review Comment:
do we need to exclude empty scan here?
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -42,8 +42,13 @@ object GlutenIcebergSourceUtil {
private val InputFileBlockStartCol = "input_file_block_start"
private val InputFileBlockLengthCol = "input_file_block_length"
- def getClassOfSparkBatchQueryScan(): Class[SparkBatchQueryScan] = {
- classOf[SparkBatchQueryScan]
+ def supportsScan(sparkScan: Scan): Boolean = sparkScan match {
+ case _: SparkBatchQueryScan => true
+ case scan: SparkStagedScan =>
+ val tasks = getScanTasks(scan)
+ tasks.nonEmpty &&
+ (tasks.forall(_.isFileScanTask) ||
tasks.forall(_.isInstanceOf[CombinedScanTask]))
Review Comment:
here it also excluded some of scan tasks, is this necessary?
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -42,8 +42,13 @@ object GlutenIcebergSourceUtil {
private val InputFileBlockStartCol = "input_file_block_start"
private val InputFileBlockLengthCol = "input_file_block_length"
- def getClassOfSparkBatchQueryScan(): Class[SparkBatchQueryScan] = {
- classOf[SparkBatchQueryScan]
+ def supportsScan(sparkScan: Scan): Boolean = sparkScan match {
Review Comment:
maybe a clear func name: `isSupportedScan`
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -162,65 +157,72 @@ object GlutenIcebergSourceUtil {
metadataColumns
}
- def getFileFormat(sparkScan: Scan): ReadFileFormat = sparkScan match {
- case scan: SparkBatchQueryScan =>
- val tasks = scan.tasks().asScala
- asFileScanTask(tasks.toList).foreach {
- task =>
- task.file().format() match {
- case FileFormat.PARQUET => return ReadFileFormat.ParquetReadFormat
- case FileFormat.ORC => return ReadFileFormat.OrcReadFormat
- case _ =>
- }
- }
- throw new GlutenNotSupportException("Iceberg Only support parquet and
orc file format.")
- case _ =>
- throw new GlutenNotSupportException("Only support iceberg
SparkBatchQueryScan.")
+ def getFileFormat(sparkScan: Scan): ReadFileFormat = {
+ asFileScanTask(getScanTasks(sparkScan)).foreach {
+ task =>
+ task.file().format() match {
+ case FileFormat.PARQUET => return ReadFileFormat.ParquetReadFormat
+ case FileFormat.ORC => return ReadFileFormat.OrcReadFormat
+ case _ =>
+ }
+ }
+ throw new GlutenNotSupportException("Iceberg Only support parquet and orc
file format.")
}
- def getReadPartitionSchema(sparkScan: Scan): StructType = sparkScan match {
- case scan: SparkBatchQueryScan =>
- val tasks = scan.tasks().asScala
- asFileScanTask(tasks.toList).foreach {
- task =>
- val spec = task.spec()
- if (spec.isPartitioned) {
- val readFields = scan.readSchema().fields.map(_.name).toSet
- // Iceberg will generate some non-table fields as partition
fields, such as x_bucket,
- // which will not appear in readFields, they also cannot be
filtered.
- val tableFields =
spec.schema().columns().asScala.map(_.name()).toSet
- val voidTransformFields = scan
- .table()
- .spec()
+ def getReadPartitionSchema(sparkScan: Scan): StructType = {
+ asFileScanTask(getScanTasks(sparkScan)).foreach {
+ task =>
+ val spec = task.spec()
+ if (spec.isPartitioned) {
+ val readFields = sparkScan.readSchema().fields.map(_.name).toSet
+ // Iceberg will generate some non-table fields as partition fields,
such as x_bucket,
+ // which will not appear in readFields, they also cannot be filtered.
+ val tableFields = spec.schema().columns().asScala.map(_.name()).toSet
+ val voidTransformFields = getTable(sparkScan)
+ .spec()
Review Comment:
we already use `foreach getScanTasks(sparkScan)` to get task.spec, can we
skip the `.spec()` here?
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -31,7 +31,7 @@ import org.apache.spark.sql.types.StructType
import org.apache.iceberg._
import org.apache.iceberg.spark.SparkSchemaUtil
-import java.lang.{Class, Long => JLong}
+import java.lang.{Long => JLong}
Review Comment:
not necessary to use {}
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]