infvg commented on code in PR #12833:
URL: https://github.com/apache/gluten/pull/12833#discussion_r4120016377
##########
gluten-iceberg/src/main/scala/org/apache/gluten/execution/IcebergScanTransformer.scala:
##########
@@ -192,8 +192,15 @@ case class IcebergScanTransformer(
override def getDataSchema: StructType = new StructType()
- // TODO: get root paths from table.
- override def getRootPathsInternal: Seq[String] = Seq.empty
+ // On Spark 3.3, SparkShims.getBatchScanExecTable always returns null
(BatchScanExec has no
+ // `table` field until Spark 3.4), so this falls back to the previous
Seq.empty behavior there;
+ // on Spark 3.4+ it returns the Iceberg table's base location.
+ override def getRootPathsInternal: Seq[String] = {
+ table match {
+ case t: SparkTable => Seq(t.table().location())
+ case _ => Seq.empty
+ }
+ }
Review Comment:
My main concern is that if we change the table's write.data.path, it still
doesn't move its old files so we can miss the actual file system. Maybe we
could the existing getScanTasks instead of finalPartitions (since we use it to
get the file format).
try something like this, i tested it locally should work (also we should
check delete file paths too):
```scala
def getFileFormatAndRootPaths(
sparkScan: Scan,
collectRootPaths: Boolean): (ReadFileFormat, Seq[String]) = {
var fileFormat = ReadFileFormat.UnknownFormat
val pathsByScheme = mutable.LinkedHashMap.empty[String, String]
var lastSchemePrefix = ""
asFileScanTask(getScanTasks(sparkScan)).foreach {
task =>
if (fileFormat == ReadFileFormat.UnknownFormat) {
task.file().format() match {
case FileFormat.PARQUET => fileFormat =
ReadFileFormat.ParquetReadFormat
case FileFormat.ORC => fileFormat = ReadFileFormat.OrcReadFormat
case _ =>
}
}
if (collectRootPaths) {
val path = task.file().path().toString
if (lastSchemePrefix.isEmpty ||
!path.startsWith(lastSchemePrefix)) {
val prefix = path.substring(0, path.indexOf(':') + 1)
lastSchemePrefix = if (pathsByScheme.contains(prefix)) {
prefix
} else {
Option(new Path(path).toUri.getScheme).fold("")(_ + ":")
}
pathsByScheme.getOrElseUpdate(lastSchemePrefix, path)
}
} else if (fileFormat != ReadFileFormat.UnknownFormat) {
return (fileFormat, Seq.empty)
}
}
if (fileFormat == ReadFileFormat.UnknownFormat) {
throw new GlutenNotSupportException(
"Iceberg Only support parquet and orc file format.")
}
(fileFormat, pathsByScheme.values.toVector)
}
```
```scala
private lazy val scanFileInfo =
GlutenIcebergSourceUtil.getFileFormatAndRootPaths(
scan,
GlutenConfig.get.scanFileSchemeValidationEnabled)
override lazy val fileFormat: ReadFileFormat = scanFileInfo._1
override def getRootPathsInternal: Seq[String] = scanFileInfo._2
```
--
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]