diff --git a/docs/source/user-guide/latest/iceberg-writes.md b/docs/source/user-guide/latest/iceberg-writes.md index 6f91298183f..df5edfe34fd 100644 --- a/docs/source/user-guide/latest/iceberg-writes.md +++ b/docs/source/user-guide/latest/iceberg-writes.md @@ -174,7 +174,7 @@ A write is eligible only when ALL of the following hold: | `write.metadata.metrics.*` | any value (manifest metrics are re-derived on the JVM with Iceberg's own logic) | | `write.spark.fanout.enabled` | any value (the native writer implements both clustered and fanout modes) | | `write.target-file-size-bytes` | any value (the two writers can choose different roll points; see accepted divergences) | -| data location URI scheme | `file`, `memory`, `s3`, `s3a`, `gs` (`gs` only when the `FileIO` opening the data location is a `GCSFileIO`; see below) | +| data location URI scheme | `file`, `memory`, `s3`, `s3a`, `gs`, matched case-sensitively (`S3://` falls back). `s3`, `s3a` and `gs` need a bucket in the authority (`s3://bucket/...`), so a hostless form such as `s3:/bucket/key` falls back. `gs` only when the `FileIO` opening the data location is a `GCSFileIO`; see below | | resolved `table.locationProvider()` | Iceberg's built-in `DefaultLocationProvider` | | partition spec | any | | column types | any except `uuid` (Spark plans it as a string; no Arrow cast reaches `fixed(16)`) | diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 70982f04610..f7517aed6b2 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -402,6 +402,10 @@ fn env_region_present() -> bool { /// S3-compliant scan to the local filesystem. A `/` before the `:` means there is no scheme (the /// `:` sits inside a path segment, e.g. `/tmp/a:b`), so those and truly schemeless paths default /// to `file`. +/// +/// The JVM write gate (`CometIcebergNativeWrite.storageScheme`) mirrors this rule exactly, case +/// included. Change both together, and keep the cases in +/// `scheme_of_extracts_scheme_from_all_uri_forms` in step with its `storageScheme` test. fn scheme_of(path: &str) -> &str { match path.split_once(':') { Some((scheme, _)) if !scheme.is_empty() && !scheme.contains('/') => scheme, @@ -639,7 +643,17 @@ mod tests { assert_eq!(scheme_of("blob://bucket/key"), "blob"); assert_eq!(scheme_of("blob:/bucket/key"), "blob"); assert_eq!(scheme_of("s3://bucket/key"), "s3"); + assert_eq!(scheme_of("s3:/bucket/key"), "s3"); + // Hadoop normalises `hdfs:///p` to `hdfs:/p`. The JVM write gate must read both as + // `hdfs` (unsupported) rather than admit the hostless form as `file`. + assert_eq!(scheme_of("hdfs:/warehouse/t"), "hdfs"); + assert_eq!(scheme_of("hdfs:///warehouse/t"), "hdfs"); + assert_eq!(scheme_of("hdfs://nn:8020/warehouse/t"), "hdfs"); + assert_eq!(scheme_of("memory:/x"), "memory"); assert_eq!(scheme_of("file:///tmp/x"), "file"); + assert_eq!(scheme_of("file:/tmp/x"), "file"); + // Not lowercased: `storage_factory_for` matches case-sensitively. + assert_eq!(scheme_of("S3://bucket/key"), "S3"); // Schemeless and colon-in-path locals default to the local FS. assert_eq!(scheme_of("/tmp/no-scheme"), "file"); assert_eq!(scheme_of("/tmp/a:b"), "file"); diff --git a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala index 3954e2c67b3..82e2e789b03 100644 --- a/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala +++ b/spark/src/main/scala/org/apache/comet/serde/operator/CometIcebergNativeWrite.scala @@ -87,6 +87,9 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { // `gs` is additionally gated on the resolved FileIO (`requireGcsFileIOForGcsDataLocation`). private val SupportedStorageSchemes: Set[String] = Set("file", "memory", "s3", "s3a", "gs") + // Supported schemes whose native backend is local and needs no host. Every other supported + // scheme reads its bucket from the URL host (`requireSupportedStorageScheme`). + private val LocalStorageSchemes: Set[String] = Set("file", "memory") private val MinUnsupportedFormatVersion = 3 private val ParquetWritePropertyPrefix = "write.parquet." private val ParquetMrPropertyPrefix = "parquet." @@ -331,20 +334,51 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { .find(k => !IgnoredHadoopParquetConfKeys.contains(k)) .map(k => s"Hadoop configuration sets $k (reaches iceberg-java's writer but not native)") - private def storageScheme(location: String): String = - if (location.contains("://")) { - location.substring(0, location.indexOf("://")).toLowerCase(Locale.ROOT) - } else { - "file" - } + /** + * The scheme the native writer picks its storage backend from. Must follow the same rule as + * `scheme_of` in `native/core/src/execution/operators/iceberg_common.rs`: split on the first + * `:`, not `://`, so a hostless `hdfs:/warehouse/t` (as Hadoop normalises `hdfs:///...`) is + * read as `hdfs` rather than admitted as `file`. An empty prefix, or one containing `/` (a `:` + * inside a path segment such as `/tmp/a:b`), means there is no scheme. The scheme is kept as + * written, not lowercased: `storage_factory_for` matches it case-sensitively, so `S3://` must + * be declined here rather than fail at execution. + * + * String-based rather than `java.net.URI` (`NativeConfig.lowerScheme`): `URI` throws on + * characters an Iceberg location may carry unencoded, and its scheme grammar is not the + * first-`:` split that `scheme_of` uses. + */ + private[comet] def storageScheme(location: String): String = { + val colon = location.indexOf(':') + val prefix = if (colon > 0) location.substring(0, colon) else "" + if (prefix.isEmpty || prefix.contains('/')) "file" else prefix + } + + /** + * True when `location` carries a non-empty authority (`scheme://host/...`). iceberg-rust's S3 + * and GCS backends take the bucket from the URL host and never from the path, so a hostless + * `s3:/bucket/key` or `s3:///bucket/key` fails natively with a missing-bucket error. + * + * The write-side counterpart of `CometScanRule.hasOpenableAuthority`, which additionally admits + * hostless S3-compliant aliases because the native reader promotes their bucket from the path. + * The write gate admits no aliases, so it needs no such exception. + */ + private[comet] def hasBucketAuthority(location: String): Boolean = { + val rest = location.substring(location.indexOf(':') + 1) + rest.startsWith("//") && rest.length > 2 && rest.charAt(2) != '/' + } private val requireSupportedStorageScheme: TriggerRule = ctx => IcebergReflection.getDataLocation(ctx.table) match { case None => Some("could not resolve the table data location") case Some(location) => val scheme = storageScheme(location) - if (SupportedStorageSchemes.contains(scheme)) None - else Some(s"unsupported storage scheme: $scheme") + if (!SupportedStorageSchemes.contains(scheme)) { + Some(s"unsupported storage scheme: $scheme") + } else if (!LocalStorageSchemes.contains(scheme) && !hasBucketAuthority(location)) { + Some(s"$scheme data location has no bucket in its authority: $location") + } else { + None + } } // HadoopFileIO takes its GCS configuration from `fs.gs.*`, which is not forwarded to the diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala index 2e2cc3f2344..7b34cfb20f5 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -499,6 +499,63 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes } } + test("fall-back: hostless hdfs:/ data location is read as hdfs, not file") { + // Hadoop normalises `hdfs:///p` to `hdfs:/p`; with no `://` the gate used to call it `file`. + withDetectionCatalog { dir => + createTable( + dir, + "hostless_hdfs", + partitionSpec = "", + properties = Some("'write.data.path'='hdfs:/iceberg/db/hostless_hdfs'")) + assertUnsupportedContains( + planInsertWriteExec(s"$catalog.$ns.hostless_hdfs"), + "hostless_hdfs", + "unsupported storage scheme: hdfs") + } + } + + test("fall-back: s3 data location without a bucket in its authority") { + withDetectionCatalog { dir => + createTable( + dir, + "hostless_s3", + partitionSpec = "", + properties = Some("'write.data.path'='s3:/nonexistent-bucket/iceberg/db/hostless_s3'")) + assertUnsupportedContains( + planInsertWriteExec(s"$catalog.$ns.hostless_s3"), + "hostless_s3", + "s3 data location has no bucket") + } + } + + test("storageScheme follows the native scheme_of rule") { + // Keep in step with `scheme_of_extracts_scheme_from_all_uri_forms` in iceberg_common.rs. + Seq( + "hdfs:/warehouse/t" -> "hdfs", + "hdfs:///warehouse/t" -> "hdfs", + "hdfs://nn:8020/warehouse/t" -> "hdfs", + "s3://bucket/key" -> "s3", + "s3:/bucket/key" -> "s3", + "blob:/bucket/key" -> "blob", + "memory:/x" -> "memory", + "file:///tmp/x" -> "file", + "file:/tmp/x" -> "file", + "/tmp/no-scheme" -> "file", + "/tmp/a:b" -> "file", + "S3://bucket/key" -> "S3").foreach { case (location, expected) => + assert(CometIcebergNativeWrite.storageScheme(location) == expected, location) + } + } + + test("hasBucketAuthority requires a non-empty host after //") { + Seq("s3://bucket/key", "s3a://bucket", "gs://bucket/x").foreach { location => + assert(CometIcebergNativeWrite.hasBucketAuthority(location), location) + } + Seq("s3:/bucket/key", "s3:///bucket/key", "s3:bucket/key", "gs://").foreach { location => + assert(!CometIcebergNativeWrite.hasBucketAuthority(location), location) + } + } + test("Compatible when the data location scheme is s3") { withDetectionCatalog { dir => createTable(