Skip to content
Merged
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
2 changes: 1 addition & 1 deletion docs/source/user-guide/latest/iceberg-writes.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)`) |
Expand Down
14 changes: 14 additions & 0 deletions native/core/src/execution/operators/iceberg_common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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."
Expand Down Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
Loading