Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
package org.apache.gluten.execution

import org.apache.gluten.config.GlutenConfig
import org.apache.gluten.utils.VeloxFileSystemValidationJniWrapper

import org.apache.spark.sql.Row
import org.apache.spark.sql.connector.catalog.{Identifier, TableCatalog}
Expand Down Expand Up @@ -64,4 +65,42 @@ class VeloxIcebergSuite extends IcebergSuite {
}
}
}

test("iceberg root paths on an unsupported scheme fail native scheme validation") {
// See https://github.com/apache/gluten/issues/12712: IcebergScanTransformer used to always
// report Seq.empty root paths, so VeloxBackend.validateScanExec's scheme check was silently
// skipped for every Iceberg scan, regardless of the scan's actual filesystem.
withTable("iceberg_root_paths_scheme_tb") {
spark.sql("""
|CREATE TABLE iceberg_root_paths_scheme_tb (id INT)
|USING iceberg
|""".stripMargin)
spark.sql("INSERT INTO iceberg_root_paths_scheme_tb VALUES (1), (2)")

runQueryAndCompare("SELECT * FROM iceberg_root_paths_scheme_tb") {
df =>
val scans = getExecutedPlan(df).collect { case i: IcebergScanTransformer => i }
assert(scans.size == 1)
val rootPaths = scans.head.getRootPathsInternal
assert(rootPaths.nonEmpty)
// Confirm these are real per-file scan paths (registered as `file` scheme), which
// VeloxBackendSettings.distinctRootPaths always excludes from scheme validation (the
// local filesystem is always registered) -- so a real root path from this suite
// passes validation, as expected for a supported local table.
assert(
VeloxFileSystemValidationJniWrapper.allSupportedByRegisteredFileSystems(
rootPaths.toArray))

// Since this suite only exercises the local `file` scheme, which is always
// considered supported, separately confirm -- with a clean, synthetic URI, not
// derived from the real paths above -- that Velox's own native filesystem check (the
// same one VeloxBackend.validateScanExec relies on to decide whether to fall back)
// correctly rejects a genuinely unsupported scheme.
assert(
!VeloxFileSystemValidationJniWrapper.allSupportedByRegisteredFileSystems(
Array("unsupported-test-scheme://bucket/path/file.parquet")),
"expected an unsupported scheme to fail native filesystem validation")
}
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -192,8 +192,8 @@ case class IcebergScanTransformer(

override def getDataSchema: StructType = new StructType()

// TODO: get root paths from table.
override def getRootPathsInternal: Seq[String] = Seq.empty
override def getRootPathsInternal: Seq[String] =
GlutenIcebergSourceUtil.getRootPaths(finalPartitions)

private lazy val readSchemaFields =
scan.readSchema().fieldNames.map(_.toLowerCase(Locale.ROOT)).toSet
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ import org.apache.gluten.execution.SparkDataSourceRDDPartition
import org.apache.gluten.substrait.rel.{IcebergLocalFilesBuilder, SplitInfo}
import org.apache.gluten.substrait.rel.LocalFilesNode.ReadFileFormat

import org.apache.spark.Partition
import org.apache.spark.softaffinity.SoftAffinity
import org.apache.spark.sql.catalyst.catalog.ExternalCatalogUtils
import org.apache.spark.sql.connector.read.Scan
Expand Down Expand Up @@ -112,6 +113,31 @@ object GlutenIcebergSourceUtil {
)
}

/**
* Returns one representative data file path per planned Spark input partition, so callers can
* validate the underlying filesystem scheme(s) without enumerating every data file.
*
* We deliberately read the actual scanned file paths (task.file().path()) rather than the
* Iceberg table's declared base location (Table.location()): Iceberg lets a table relocate its
* actual data files independently of the table location via the write.data.path property or a
* custom LocationProvider, so the table location is not guaranteed to reflect the filesystem
* the data is actually read from. This also works uniformly across all supported Spark
* versions, since it does not depend on BatchScanExec exposing its Table (only added in Spark
* 3.4; see SparkShims.getBatchScanExecTable).
*/
def getRootPaths(partitions: Seq[Partition]): Seq[String] = {
partitions.flatMap {
case p: SparkDataSourceRDDPartition =>
p.inputPartitions.flatMap {
case ip: SparkInputPartition =>
val tasks = ip.taskGroup[ScanTask]().tasks().asScala
asFileScanTask(tasks.toList).headOption.map(_.file().path().toString)
case _ => None
}
case _ => Seq.empty
}
}

def getFieldIds(sparkScan: Scan): JHashMap[String, Integer] = {
val fieldIds = new JHashMap[String, Integer]()
sparkScan match {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,32 @@ abstract class IcebergSuite extends WholeStageTransformerSuite {
}
}

test("iceberg getRootPathsInternal returns actual scanned file paths") {
// See https://github.com/apache/gluten/issues/12712: getRootPathsInternal used to always
// return Seq.empty for Iceberg scans, silently skipping native filesystem scheme validation.
withTable("iceberg_root_paths_tb") {
spark.sql("""
|CREATE TABLE iceberg_root_paths_tb (id INT)
|USING iceberg
|""".stripMargin)
spark.sql("INSERT INTO iceberg_root_paths_tb VALUES (1), (2)")

runQueryAndCompare("SELECT * FROM iceberg_root_paths_tb") {
df =>
val scans = getExecutedPlan(df).collect { case i: IcebergScanTransformer => i }
assert(scans.size == 1)
val rootPaths = scans.head.getRootPathsInternal
assert(rootPaths.nonEmpty, "getRootPathsInternal should not be empty for Iceberg tables")
// Assert the paths point at the actual scanned data files, not just some non-empty
// placeholder: Iceberg tables can relocate real data files independently of the
// table's declared base location (e.g. via write.data.path / a custom
// LocationProvider), so only the real per-file scan path reliably reflects the
// filesystem the data is actually read from.
assert(rootPaths.forall(p => p.contains("/data/") && p.endsWith(".parquet")))
}
}
}

test("iceberg input_file_name") {
withTable("iceberg_input_file_tb") {
spark.sql("""
Expand Down