From a5a3de826cb5d07c59ddd1660e9c3b355bb0352a Mon Sep 17 00:00:00 2001 From: 0lai0 Date: Thu, 1 Oct 2026 12:33:59 +0300 Subject: [PATCH 1/2] fix: read the Iceberg write gate's storage scheme the same way the native factory does --- .../src/execution/operators/iceberg_common.rs | 12 ++++ .../operator/CometIcebergNativeWrite.scala | 51 ++++++++++++++--- .../CometIcebergWriteDetectionSuite.scala | 56 +++++++++++++++++++ 3 files changed, 111 insertions(+), 8 deletions(-) diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 509e3b1f986..5722472a18f 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -219,6 +219,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, except that it +/// lowercases the scheme. 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, @@ -284,7 +288,15 @@ 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"); // 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 94f6b139273..4d87ba0faa4 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." @@ -307,20 +310,52 @@ 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. + * + * Unlike `scheme_of`, this lowercases the scheme, so `S3://bucket/key` is admitted here but + * rejected natively. + * + * 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.toLowerCase(Locale.ROOT) + } + + /** + * 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 37ddb1872e8..d8b6c0e633d 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -444,6 +444,62 @@ 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").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( From 5ab1629c4413ee00a24fc29119c4fec8c888d6ad Mon Sep 17 00:00:00 2001 From: 0lai0 Date: Fri, 2 Oct 2026 23:49:09 +0300 Subject: [PATCH 2/2] address comment --- docs/source/user-guide/latest/iceberg-writes.md | 2 +- native/core/src/execution/operators/iceberg_common.rs | 6 ++++-- .../comet/serde/operator/CometIcebergNativeWrite.scala | 9 ++++----- .../apache/comet/CometIcebergWriteDetectionSuite.scala | 3 ++- 4 files changed, 11 insertions(+), 9 deletions(-) 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 8274b8c1129..f7517aed6b2 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -403,8 +403,8 @@ fn env_region_present() -> bool { /// `:` 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, except that it -/// lowercases the scheme. Change both together, and keep the cases in +/// 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(':') { @@ -652,6 +652,8 @@ mod tests { 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 7e6d991a824..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 @@ -339,10 +339,9 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { * `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. - * - * Unlike `scheme_of`, this lowercases the scheme, so `S3://bucket/key` is admitted here but - * rejected natively. + * 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 @@ -351,7 +350,7 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { 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.toLowerCase(Locale.ROOT) + if (prefix.isEmpty || prefix.contains('/')) "file" else prefix } /** diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala index 93bf60af38e..7b34cfb20f5 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -541,7 +541,8 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes "file:///tmp/x" -> "file", "file:/tmp/x" -> "file", "/tmp/no-scheme" -> "file", - "/tmp/a:b" -> "file").foreach { case (location, expected) => + "/tmp/a:b" -> "file", + "S3://bucket/key" -> "S3").foreach { case (location, expected) => assert(CometIcebergNativeWrite.storageScheme(location) == expected, location) } }