infvg commented on code in PR #12833:
URL: https://github.com/apache/gluten/pull/12833#discussion_r4210885217
##########
gluten-iceberg/src/main/scala/org/apache/iceberg/spark/source/GlutenIcebergSourceUtil.scala:
##########
@@ -155,16 +155,69 @@ object GlutenIcebergSourceUtil {
metadataColumns
}
- def getFileFormat(sparkScan: Scan): ReadFileFormat = {
+ /**
+ * Derives both the file format and, optionally, the root path(s) actually
scanned from the
+ * planned Iceberg scan tasks in a single pass over
`getScanTasks(sparkScan)`.
+ *
+ * Root paths are derived from the real per-file (and, when present,
per-delete-file) paths rather
+ * than table metadata (Table.location() / write.data.path), since data
files can live at a
+ * location the static metadata does not reflect (e.g. after write.data.path
is changed without
+ * moving previously-written files). One representative path per distinct
scheme prefix is kept,
+ * since callers only need this to validate that the native filesystem
implementation supports
+ * every scheme in play, not to enumerate every file.
+ *
+ * NOTE: this deliberately reuses the same `getScanTasks(sparkScan)`
(Iceberg's own
+ * SparkPartitioningAwareScan#tasks()) that `getReadPartitionSchema` already
calls unconditionally
+ * from `doValidateInternal()`, not Spark's
BatchScanExecShim#filteredPartitions (which
+ * additionally calls SupportsRuntimeV2Filtering#filter() and
+ * Scan#toBatch().planInputPartitions(), and was the root cause of the DPP
regression from
+ * https://github.com/apache/gluten/issues/12712 when a prior revision of
this fix read paths from
+ * BatchScanExecTransformerBase#finalPartitions).
+ */
+ def getFileFormatAndRootPaths(
+ sparkScan: Scan,
+ collectRootPaths: Boolean): (ReadFileFormat, Seq[String]) = {
+ var fileFormat: ReadFileFormat = ReadFileFormat.UnknownFormat
+ // Built with plain mutable local collections that are never returned or
captured, then
+ // copied into an immutable Seq at the end: the result is cached in a
`lazy val` on
+ // IcebergScanTransformer (a case class captured by Spark task closures),
so it must not
+ // retain any reference to a non-serializable mutable collection.
+ val seenSchemes = scala.collection.mutable.HashSet.empty[String]
+ val rootPathsBuilder = Seq.newBuilder[String]
+
+ def recordPath(path: String): Unit = {
+ if (collectRootPaths) {
+ // Extract just the scheme prefix (e.g. "s3://", "hdfs://") without
going through
+ // java.net.URI, whose strict RFC 3986 parsing throws
URISyntaxException on file paths
+ // containing characters like spaces or '[' / ']' that Hadoop-style
paths otherwise
+ // tolerate (e.g. from a partition value with a space in it).
+ val schemeSeparator = path.indexOf("://")
+ val scheme = if (schemeSeparator >= 0) path.substring(0,
schemeSeparator) else ""
Review Comment:
i think `new Path(path).toUri.getScheme` would be better here. Paths like
`file:/tmp/data.parquet and hdfs:/warehouse/delete.parquet` are valid Hadoop
paths, but both get an empty scheme with this check.
I tested this locally: with the local data file first, the HDFS delete path
gets dropped. Velox then filters out the remaining file: path, so we skip
filesystem validation and can still offload without HDFS support.
--
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]