diff --git a/docs/source/user-guide/latest/datasources.md b/docs/source/user-guide/latest/datasources.md index eb56d1e8970..993f7c7e7af 100644 --- a/docs/source/user-guide/latest/datasources.md +++ b/docs/source/user-guide/latest/datasources.md @@ -289,7 +289,10 @@ URLs. This is opt-in and disabled by default. Enable it by listing the schemes to treat as S3-compliant aliases in `spark.hadoop.fs.comet.s3Compliant.schemes` (Hadoop key `fs.comet.s3Compliant.schemes`), a comma-separated, case-insensitive list. This mirrors the -existing `fs.comet.libhdfs.schemes` config. +existing `fs.comet.libhdfs.schemes` config. The list entries are case-insensitive, but an Iceberg +table location must write the alias in lowercase (`blob://`, not `BLOB://`): the native Iceberg +reader opens each recorded location as written, and its S3 backend accepts only a lowercase +scheme prefix, so Comet declines such a location up front and leaves that scan to Spark. ```shell --conf spark.hadoop.fs.comet.s3Compliant.schemes=blob diff --git a/native/core/src/execution/operators/iceberg_common.rs b/native/core/src/execution/operators/iceberg_common.rs index 262325c30cb..a17eb076e0d 100644 --- a/native/core/src/execution/operators/iceberg_common.rs +++ b/native/core/src/execution/operators/iceberg_common.rs @@ -48,16 +48,12 @@ const ICEBERG_PROVIDER_CLASS_PROPERTY: &str = "s3.comet.credential.provider.clas /// iceberg-storage-opendal's own settings, such as `opendal.io-timeout-ms`. const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client.", "opendal."]; -/// Pick an OpenDAL storage backend from a URI's scheme. `file` (or no scheme) falls through to -/// the local file system. `memory` is used by the write path to assemble manifest bytes that -/// stay entirely in-process. For S3, the Comet credential bridge is wired in when a provider -/// class is configured; `access_mode` is forwarded to the JVM SPI so the read and write paths can -/// be granted different (e.g. read-only vs read-write) credentials. +/// Pick an OpenDAL storage backend for a URI whose scheme `builtin_storage_schemes` lists for +/// `access_mode`, or that is an opted-in S3-compliant alias written in lowercase; anything else is +/// rejected before the match, so the list alone decides what each mode admits. For S3, the Comet +/// credential bridge is wired in when a provider class is configured and `access_mode` is +/// forwarded to the JVM SPI. /// -/// The JVM planner mirrors the schemes admitted here by hand, so that it can decline a scan -/// cleanly instead of failing at execution. Changing the arms below means updating -/// `CometScanRule.icebergReadableSchemes` (reads) and -/// `CometIcebergNativeWrite.SupportedStorageSchemes` (writes). /// The bool is false when a read fell back to opendal's default chain because the configured S3 /// access provider failed to initialise; such a `FileIO` must not be cached, so the next task /// retries. @@ -67,24 +63,22 @@ pub(crate) fn storage_factory_for( catalog_name: &str, access_mode: AccessMode, ) -> Result<(Arc, bool), DataFusionError> { + // Verbatim match: OpenDAL strips the scheme prefix from every path case-sensitively at open + // time, so admitting `S3://` here would only defer the failure. Aliases are held to the same + // rule (`is_iceberg_alias_scheme`). The JVM gates match verbatim. let scheme = scheme_of(path); + if !builtin_storage_schemes(access_mode).contains(&scheme) + && !is_s3_family_scheme(scheme, catalog_properties) + { + return Err(DataFusionError::Execution(format!( + "Unsupported storage scheme: {scheme}" + ))); + } match scheme { "file" => Ok((Arc::new(OpenDalStorageFactory::Fs), true)), "memory" => Ok((Arc::new(OpenDalStorageFactory::Memory), true)), "gs" => Ok((Arc::new(OpenDalStorageFactory::Gcs), true)), - // Reads keep the OSS backend they have always had (CometScanRule admits `oss` scan - // locations through HadoopFileIO). Writes fail closed: Comet does not forward `oss.*` - // properties into the FileIO and no test covers the write path, so OSS-specific - // endpoint/credential configuration could silently be dropped. The JVM write gate - // already declines `oss` locations; this is the native-side backstop. - "oss" => match access_mode { - AccessMode::Read => Ok((Arc::new(OpenDalStorageFactory::Oss), true)), - AccessMode::Write => Err(DataFusionError::Execution( - "OSS is not supported for native Iceberg writes (oss.* properties are not \ - forwarded to the native FileIO)" - .to_string(), - )), - }, + "oss" => Ok((Arc::new(OpenDalStorageFactory::Oss), true)), // s3, s3a, and any opted-in s3-compliant alias (e.g. blob) route to the S3 backend. Listed // last so the built-in backends above stay authoritative even if one of their schemes is // also named in `fs.comet.s3Compliant.schemes`. An alias additionally gets a wrapper that @@ -93,7 +87,7 @@ pub(crate) fn storage_factory_for( // location-scoped provider gets a storage that serves each file with the credential of // its policy location -- see iceberg_location_scoped. s if is_s3_family_scheme(s, catalog_properties) => { - let alias = is_s3_compliant_alias_scheme(s, catalog_properties); + let alias = is_iceberg_alias_scheme(s, catalog_properties); let (access, cacheable) = build_s3_access(path, catalog_properties, catalog_name, access_mode)?; let factory: Arc = match access { @@ -111,6 +105,7 @@ pub(crate) fn storage_factory_for( }; Ok((factory, cacheable)) } + // Only a listed scheme without an arm reaches here; the tests below fail on that. _ => Err(DataFusionError::Execution(format!( "Unsupported storage scheme: {scheme}" ))), @@ -303,6 +298,17 @@ fn cached_file_io( Ok(file_io) } +/// The single point of change for the built-in schemes `storage_factory_for` admits per access +/// mode; the JVM read and write gates load these lists over JNI. `memory` is write-only: an OpenDAL +/// memory backend is a fresh empty in-process store the write path assembles manifests in, so a +/// read finds nothing. `oss` is read-only: no `oss.*` property is forwarded and no test covers it. +pub(crate) fn builtin_storage_schemes(access_mode: AccessMode) -> &'static [&'static str] { + match access_mode { + AccessMode::Read => &["file", "gs", "oss", "s3", "s3a"], + AccessMode::Write => &["file", "memory", "gs", "s3", "s3a"], + } +} + pub(crate) fn load_file_io( catalog_properties: &HashMap, reference_path: &str, @@ -587,7 +593,19 @@ pub(crate) fn scheme_of(path: &str) -> &str { /// iceberg-rust's storage factory has no `s3n` backend and the Scala Iceberg scheme gate /// (`isIcebergReadableScheme`) rejects `s3n` before a path ever reaches this operator. fn is_s3_family_scheme(scheme: &str, catalog_properties: &HashMap) -> bool { - matches!(scheme, "s3" | "s3a") || is_s3_compliant_alias_scheme(scheme, catalog_properties) + matches!(scheme, "s3" | "s3a") || is_iceberg_alias_scheme(scheme, catalog_properties) +} + +/// True if `scheme` is an opted-in S3-compliant alias written exactly as the lowercase form of its +/// list entry. The shared `is_s3_compliant_alias_scheme` is case-insensitive, which is safe for +/// the Parquet path because it rewrites an alias URL to `s3://` before anything opens it. The +/// Iceberg path opens the recorded location as written, and OpenDAL's S3 backend checks it against +/// a `scheme://bucket/` prefix whose scheme comes from `Url::parse` and so is lowercase, so +/// `BLOB://bucket/key` would pass a case-insensitive gate here and fail at open time. The JVM +/// Iceberg gate matches the same way. +fn is_iceberg_alias_scheme(scheme: &str, catalog_properties: &HashMap) -> bool { + !scheme.bytes().any(|b| b.is_ascii_uppercase()) + && is_s3_compliant_alias_scheme(scheme, catalog_properties) } #[cfg(test)] @@ -823,7 +841,25 @@ mod tests { // oss.* property forwarding exists and is tested. assert!(factory_result("oss://bucket/db/table", AccessMode::Read).is_ok()); let err = factory_result("oss://bucket/db/table", AccessMode::Write).unwrap_err(); - assert!(err.contains("OSS"), "unexpected error: {err}"); + assert!( + err.contains("Unsupported storage scheme: oss"), + "unexpected error: {err}" + ); + } + + #[test] + fn memory_scheme_is_writable_but_not_readable() { + // The write path assembles manifests in a fresh in-process memory store. A read against + // a new memory store can never find data, so the factory must decline it up front. + assert!(factory_result("memory:manifest.avro", AccessMode::Write).is_ok()); + assert!(factory_result("memory:///manifest.avro", AccessMode::Write).is_ok()); + for path in ["memory:manifest.avro", "memory:///key.parquet"] { + let err = factory_result(path, AccessMode::Read).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme: memory"), + "unexpected error for {path}: {err}" + ); + } } #[test] @@ -831,13 +867,85 @@ mod tests { for mode in [AccessMode::Read, AccessMode::Write] { assert!(factory_result("file:///tmp/x", mode).is_ok()); assert!(factory_result("/tmp/no-scheme", mode).is_ok()); - assert!(factory_result("memory:manifest.avro", mode).is_ok()); // No credential provider configured: the default chain applies in both modes. assert!(factory_result("s3://bucket/db/table", mode).is_ok()); assert!(factory_result("gs://bucket/db/table", mode).is_ok()); } } + #[test] + fn scheme_list_pre_check_is_load_bearing() { + // The oss and memory arms build a backend for either mode; only their absence from the + // list for the other mode rejects them. Both facts are asserted so that adding an + // access-mode branch back into an arm, or listing the scheme, breaks this test. + assert!(factory_result("oss://bucket/path", AccessMode::Read).is_ok()); + assert!(!builtin_storage_schemes(AccessMode::Write).contains(&"oss")); + let err = factory_result("oss://bucket/path", AccessMode::Write).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme: oss"), + "unexpected error: {err}" + ); + assert!(factory_result("memory:///path", AccessMode::Write).is_ok()); + assert!(!builtin_storage_schemes(AccessMode::Read).contains(&"memory")); + let err = factory_result("memory:///path", AccessMode::Read).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme: memory"), + "unexpected error: {err}" + ); + } + + #[test] + fn mixed_case_scheme_is_rejected_for_both_modes() { + // OpenDAL strips the scheme prefix from every path case-sensitively at open time + // (`S3://bucket/key` fails its `s3://bucket/` prefix check), so the factory must not + // admit what the open cannot serve. The JVM gates match the built-in set verbatim too. + for mode in [AccessMode::Read, AccessMode::Write] { + for path in ["S3://bucket/key", "File:///tmp/x", "GS://bucket/key"] { + let err = factory_result(path, mode).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme"), + "unexpected error for {path} in {mode:?}: {err}" + ); + } + } + } + + #[test] + fn alias_scheme_is_matched_verbatim() { + // An alias location is opened as written, and OpenDAL's S3 backend checks it against a + // lowercase `scheme://bucket/` prefix, so a mixed-case alias that passed a + // case-insensitive gate here would only fail at open time. The list entry may be written + // in any case; the location's scheme must match its lowercase form exactly. + for listed in ["blob", " BLOB ", "minio,Blob"] { + let props = HashMap::from([( + "fs.comet.s3Compliant.schemes".to_string(), + listed.to_string(), + )]); + for mode in [AccessMode::Read, AccessMode::Write] { + assert!( + storage_factory_for("blob://bucket/key", &props, "test_cat", mode).is_ok(), + "blob://bucket/key must be admitted for {mode:?} with list {listed:?}" + ); + for path in [ + "BLOB://bucket/key", + "Blob://bucket/key", + "BLOB:///bucket/key", + ] { + let err = storage_factory_for(path, &props, "test_cat", mode) + .map(|_| ()) + .expect_err(&format!( + "{path} must be rejected for {mode:?} with list {listed:?}" + )) + .to_string(); + assert!( + err.contains("Unsupported storage scheme"), + "unexpected error for {path} in {mode:?}: {err}" + ); + } + } + } + } + #[test] fn unknown_scheme_is_rejected() { let err = factory_result("hdfs://nn/db/table", AccessMode::Read).unwrap_err(); @@ -869,6 +977,43 @@ mod tests { assert!(has_explicit_s3_credentials(&blank)); } + #[test] + fn listed_schemes_are_accepted_by_storage_factory() { + for mode in [AccessMode::Read, AccessMode::Write] { + for scheme in builtin_storage_schemes(mode) { + let url = format!("{scheme}://bucket/path"); + assert!( + factory_result(&url, mode).is_ok(), + "{scheme} is listed for {mode:?} but the factory rejects it" + ); + } + } + assert!(builtin_storage_schemes(AccessMode::Read).contains(&"oss")); + assert!(!builtin_storage_schemes(AccessMode::Write).contains(&"oss")); + assert!(!builtin_storage_schemes(AccessMode::Read).contains(&"memory")); + assert!(builtin_storage_schemes(AccessMode::Write).contains(&"memory")); + } + + #[test] + fn unlisted_schemes_are_rejected_for_both_modes() { + let unlisted = [ + "hdfs", "abfs", "abfss", "wasb", "wasbs", "gcs", "http", "https", "azure", + ]; + for mode in [AccessMode::Read, AccessMode::Write] { + for scheme in unlisted { + let err = factory_result(&format!("{scheme}://bucket/path"), mode).unwrap_err(); + assert!( + err.contains("Unsupported storage scheme"), + "unexpected error for {scheme} in {mode:?}: {err}" + ); + assert!( + !builtin_storage_schemes(mode).contains(&scheme), + "{scheme} must not be listed for {mode:?}" + ); + } + } + } + #[test] fn scheme_of_extracts_scheme_from_all_uri_forms() { // Host-bearing and hostless/opaque vendor forms must resolve to the same scheme, so an diff --git a/native/core/src/lib.rs b/native/core/src/lib.rs index b712a65c115..5eb56aaec01 100644 --- a/native/core/src/lib.rs +++ b/native/core/src/lib.rs @@ -31,6 +31,7 @@ extern crate datafusion_comet_jni_bridge; use jni::{ objects::{JClass, JString}, + sys::{jboolean, jstring, JNI_FALSE}, EnvUnowned, }; use log::info; @@ -52,6 +53,9 @@ pub mod jvm_bridge { use errors::{try_unwrap_or_throw, CometError, CometResult}; +use crate::cloud::s3::credential_bridge::AccessMode; +use crate::execution::operators::iceberg_common::builtin_storage_schemes; + pub mod alloc_accounting; pub mod cloud; pub mod comet_native_udf_bridge; @@ -270,6 +274,27 @@ pub fn tzdata_version() -> &'static str { chrono_tz::IANA_TZDB_VERSION } +/// JNI: the comma-joined built-in storage schemes the native Iceberg scan (`for_write` false) or +/// write (`for_write` true) path admits. This is the source of truth the JVM read and write gates +/// load. Opt-in S3-compliant alias schemes are not answered here; the JVM adds them from catalog +/// properties. +#[no_mangle] +pub extern "system" fn Java_org_apache_comet_NativeBase_icebergStorageSchemes( + env: EnvUnowned, + _: JClass, + for_write: jboolean, +) -> jstring { + try_unwrap_or_throw(&env, |env| { + let access_mode = if for_write != JNI_FALSE { + AccessMode::Write + } else { + AccessMode::Read + }; + let joined = builtin_storage_schemes(access_mode).join(","); + Ok(env.new_string(joined)?.into_raw()) + }) +} + // Creates a default log4rs config, which logs to console with log level. fn default_logger_config(log_level: &str) -> CometResult { let console_append = ConsoleAppender::builder() diff --git a/native/core/src/parquet/objectstore/s3_blob_fs_support.rs b/native/core/src/parquet/objectstore/s3_blob_fs_support.rs index f455f8c3849..ab09e7d24c3 100644 --- a/native/core/src/parquet/objectstore/s3_blob_fs_support.rs +++ b/native/core/src/parquet/objectstore/s3_blob_fs_support.rs @@ -98,8 +98,10 @@ pub(crate) fn normalize_object_store_url( /// True if `scheme` is a configured s3-compliant alias (`fs.comet.s3Compliant.schemes`; empty or /// unset means none). These route to `AmazonS3` but are not recognized by /// `ObjectStoreScheme::parse`. `s3a` is excluded -- object_store knows it and callers special-case -/// it. Shared by the Parquet normalizer above and the Iceberg storage-factory gate -/// (`iceberg_common::storage_factory_for`) so both admit the same schemes. +/// it. Case-insensitive, which is safe for the Parquet normalizer above because it rewrites the +/// alias to `s3://` before anything opens the path. The Iceberg storage-factory gate opens the +/// recorded location as written, so it wraps this in `iceberg_common::is_iceberg_alias_scheme`, +/// which additionally requires the scheme to be lowercase. pub(crate) fn is_s3_compliant_alias_scheme( scheme: &str, object_store_configs: &HashMap, diff --git a/spark/src/main/java/org/apache/comet/NativeBase.java b/spark/src/main/java/org/apache/comet/NativeBase.java index ab4eee856ba..7952e736116 100644 --- a/spark/src/main/java/org/apache/comet/NativeBase.java +++ b/spark/src/main/java/org/apache/comet/NativeBase.java @@ -384,4 +384,16 @@ public static void releaseNative() throws Throwable { * @return the version compiled into libcomet */ public static native String getTzdataVersion(); + + /** + * The comma-joined URL schemes the native Iceberg storage factory can open, for reads ({@code + * forWrite} false) or writes ({@code forWrite} true). This is the authoritative list the JVM + * Iceberg scan and write gates load, so the planner never hardcodes (and drifts from) the set the + * native factory actually builds. Opt-in S3-compliant alias schemes are not included; the JVM + * adds them from {@code fs.comet.s3Compliant.schemes}. + * + * @param forWrite true for the write path's schemes, false for the scan path's + * @return the supported schemes joined by commas, e.g. "file,memory,gs,s3,s3a" + */ + public static native String icebergStorageSchemes(boolean forWrite); } diff --git a/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala new file mode 100644 index 00000000000..51f4b5f7019 --- /dev/null +++ b/spark/src/main/scala/org/apache/comet/iceberg/IcebergStorageSchemes.scala @@ -0,0 +1,77 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.iceberg + +import java.util.Locale + +import org.apache.spark.internal.Logging + +import org.apache.comet.NativeBase + +/** + * The storage schemes the native Iceberg storage factory publishes over JNI, so the JVM scan and + * write gates decline what native cannot open instead of failing at execution. Native is the only + * source of these lists. Every caller sits behind `isCometLoaded`, so a JVM where the library is + * not loaded never plans a native scan or write, and both sets are simply empty there. + */ +private[comet] object IcebergStorageSchemes extends Logging { + + lazy val read: Set[String] = load(forWrite = false) + lazy val write: Set[String] = load(forWrite = true) + + /** + * Splits the comma-joined native list, trimming and lowercasing each entry. A null or blank + * list throws: `builtin_storage_schemes` is a fixed non-empty constant natively, so a loaded + * library that publishes nothing can only be a build or marshalling bug, and a warning here + * would silently disable every native Iceberg scan and write. + */ + private[comet] def parse(joined: String, forWrite: Boolean): Set[String] = { + val schemes = Option(joined).toSeq + .flatMap(_.split(",")) + .map(_.trim.toLowerCase(Locale.ROOT)) + .filter(_.nonEmpty) + .toSet + if (schemes.isEmpty) { + val path = if (forWrite) "write" else "read" + throw new IllegalStateException( + s"Comet native library published no Iceberg $path scheme list") + } + schemes + } + + /** + * Loads one list over JNI, or answers the empty set with a warning when the library is not + * loaded; that branch never runs a plan, because every caller sits behind `isCometLoaded`. + * `isLoaded` is a parameter so the unloaded answer can be tested without unloading the library. + * Once loaded, anything but a non-empty list is a build bug that must fail loudly: a native + * fault while answering the probe propagates, and a blank answer throws in `parse`. + */ + private[comet] def load( + forWrite: Boolean, + isLoaded: Boolean = NativeBase.isLoaded): Set[String] = { + if (isLoaded) { + parse(NativeBase.icebergStorageSchemes(forWrite), forWrite) + } else { + val path = if (forWrite) "write" else "read" + logWarning(s"Comet native library is not loaded; the Iceberg $path scheme list is empty") + Set.empty + } + } +} diff --git a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala index b109fa158ed..983167ad52f 100644 --- a/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala +++ b/spark/src/main/scala/org/apache/comet/rules/CometScanRule.scala @@ -47,7 +47,7 @@ import org.apache.spark.sql.types._ import org.apache.comet.{CometConf, DataTypeSupport, NativeBase} import org.apache.comet.CometConf._ import org.apache.comet.CometSparkSessionExtensions.{isCometLoaded, isSpark35Plus, withFallbackReason, withFallbackReasons} -import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection} +import org.apache.comet.iceberg.{CometIcebergNativeScanMetadata, IcebergReflection, IcebergStorageSchemes} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.parquet.CometParquetUtils.{encryptionEnabled, isEncryptionConfigSupported, readFieldId} import org.apache.comet.serde.operator.{CometIcebergNativeScan, CometNativeScan} @@ -1276,20 +1276,25 @@ object CometScanRule extends Logging { org.apache.spark.sql.catalyst.trees.TreeNodeTag[Unit]("comet.skipCometScan") /** - * Schemes Comet's native Iceberg scan can actually open, mirroring the match arms in - * `native/core/src/execution/operators/iceberg_common.rs::storage_factory_for`. Deliberately - * NOT delegated to `isNativelyReadableScheme`: object_store recognizes schemes (http/https, - * azure, memory) that iceberg-rust's OpenDAL storage factory cannot build, and admitting them - * here turns a clean JVM fallback into a native runtime "Unsupported storage scheme" error. Add - * here what you add to `storage_factory_for` (currently Aliyun `oss` and GCS `gs`). - * S3-compliant aliases like `blob` are opt-in via `fs.comet.s3Compliant.schemes` (see - * `isIcebergReadableScheme`), not hardcoded, since the native planner opens them via S3. The - * write path keeps its own list (`CometIcebergNativeWrite.SupportedStorageSchemes`), which - * differs deliberately: it excludes `oss` (fails closed, see `storage_factory_for`) and - * includes `memory`. + * Schemes Comet's native Iceberg scan can open, loaded from the native storage factory over JNI + * so this gate cannot drift from `storage_factory_for`. Lazy so that constructing the rule does + * not touch the native library before `isCometLoaded` has been consulted. Opt-in aliases from + * `fs.comet.s3Compliant.schemes` are additive (see `isIcebergReadableScheme`); the write path + * loads its own set (`CometIcebergNativeWrite.SupportedStorageSchemes`). */ - private val icebergReadableSchemes: Set[String] = - Set("file", "s3", "s3a", "gs", "oss") + private lazy val icebergReadableSchemes: Set[String] = IcebergStorageSchemes.read + + /** + * True when the Iceberg scan gate admits `scheme`, written exactly as recorded. Native opens a + * location by its raw scheme, and OpenDAL's S3 backend checks it against a `scheme://bucket/` + * prefix whose scheme comes from `Url::parse` and so is lowercase: `S3://` is not `s3://`, and + * `BLOB://` is not an opted-in `blob`. The alias set is lowercase + * (`NativeConfig.parseSchemeSet`), so an alias is admitted only when the location writes it in + * lowercase. The Parquet gate stays case-insensitive because it rewrites alias URLs to `s3://` + * before anything opens them. + */ + private def isAdmittedIcebergScheme(scheme: String, s3CompliantSchemes: Set[String]): Boolean = + icebergReadableSchemes.contains(scheme) || s3CompliantSchemes.contains(scheme) /** * "Supported schemes: ..." suffix shared by the Iceberg scheme-fallback messages. Lists the @@ -1310,8 +1315,7 @@ object CometScanRule extends Logging { s3CompliantSchemes: Set[String]): Boolean = { val scheme = uri.getScheme if (scheme == null) return true - val lower = scheme.toLowerCase(Locale.ROOT) - icebergReadableSchemes.contains(lower) || s3CompliantSchemes.contains(lower) + isAdmittedIcebergScheme(scheme, s3CompliantSchemes) } /** @@ -1373,9 +1377,6 @@ object CometScanRule extends Logging { // hasOpenableAuthority); non-empty => decline. One example suffices for the message. var hostlessLocation: Option[String] = None - // Union of the two admitted scheme sets, built once so the per-file loop does one lookup. - val openableSchemes = icebergReadableSchemes ++ s3CompliantSchemes - // Classify one data/delete file location; see `icebergReadableSchemes` for why that allowlist // is narrower than the Parquet native gate. Runs per data and delete file, so the scheme and // the bucket are each derived once and threaded down. @@ -1386,9 +1387,8 @@ object CometScanRule extends Logging { // A schemeless local path routes to iceberg-rust's LocalFs and needs no host. val scheme = uri.getScheme if (scheme == null) return - val lower = scheme.toLowerCase(Locale.ROOT) - if (!openableSchemes.contains(lower)) { - unsupportedSchemes += lower + if (!isAdmittedIcebergScheme(scheme, s3CompliantSchemes)) { + unsupportedSchemes += scheme } else if (!hasOpenableAuthority(uri, s3CompliantSchemes)) { if (hostlessLocation.isEmpty) hostlessLocation = Some(rawPath) } else { 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 e55014812be..34c07ff8dd5 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 @@ -30,7 +30,7 @@ import org.apache.spark.sql.comet.{CometIcebergWriteExec, CometNativeExec, Icebe import org.apache.comet.{CometConf, ConfigEntry} import org.apache.comet.CometSparkSessionExtensions.withFallbackReason -import org.apache.comet.iceberg.{IcebergReflection, PositionDeltaWrite, ReplaceDataWrite} +import org.apache.comet.iceberg.{IcebergReflection, IcebergStorageSchemes, PositionDeltaWrite, ReplaceDataWrite} import org.apache.comet.objectstore.NativeConfig import org.apache.comet.serde.{CometOperatorSerde, Compatible, OperatorOuterClass, SupportLevel, Unsupported} import org.apache.comet.serde.OperatorOuterClass.Operator @@ -86,12 +86,11 @@ object CometIcebergNativeWrite extends CometOperatorSerde[IcebergWriteExec] { // `timestamp_ns`, `geometry` or `geography` today, so those are declined in case it learns to. private val UnsupportedWriteTypeIds: Set[String] = Set("UUID", "VARIANT", "UNKNOWN", "TIMESTAMP_NANO", "GEOMETRY", "GEOGRAPHY") - // `oss` is deliberately absent: iceberg-rust has an OSS backend, but Comet does not forward - // `oss.*` catalog properties to it and no functional test covers the path, so an OSS write - // could silently drop endpoint/credential configuration. Fail closed until it is covered. - // `gs` is additionally gated on the resolved FileIO (`requireGcsFileIOForGcsDataLocation`). - private val SupportedStorageSchemes: Set[String] = - Set("file", "memory", "s3", "s3a", "gs") + // Loaded from the native storage factory: `builtin_storage_schemes` in + // `native/core/src/execution/operators/iceberg_common.rs` is the single point of change and + // explains why `oss` is read-only and `memory` write-only. Lazy so constructing the serde does + // not touch the native library. `gs` is additionally gated on the resolved FileIO below. + private lazy val SupportedStorageSchemes: Set[String] = IcebergStorageSchemes.write // 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") diff --git a/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java b/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java new file mode 100644 index 00000000000..95ff00b818a --- /dev/null +++ b/spark/src/test/java/org/apache/comet/hadoop/fs/FakeWasbSchemeFileSystem.java @@ -0,0 +1,51 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.comet.hadoop.fs; + +import java.net.URI; + +import org.apache.hadoop.fs.RawLocalFileSystem; + +/** + * A local-disk-backed FileSystem that reports the {@code wasb} scheme, so a test can create and + * read an Iceberg table under a {@code wasb://} warehouse without Azure. Used to assert that the + * Iceberg scan and write gates decline a table whose storage scheme the native Iceberg storage + * factory has no arm for, instead of claiming it and failing at execution. + */ +public class FakeWasbSchemeFileSystem extends RawLocalFileSystem { + + public static final String PREFIX = "wasb://fake-container"; + + public FakeWasbSchemeFileSystem() { + // Avoid `URI scheme is not "file"` error on + // RawLocalFileSystem$DeprecatedRawLocalFileStatus.getOwner + RawLocalFileSystem.useStatIfAvailable(); + } + + @Override + public String getScheme() { + return "wasb"; + } + + @Override + public URI getUri() { + return URI.create(PREFIX); + } +} diff --git a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala index 49238eee228..922b5d6073f 100644 --- a/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala +++ b/spark/src/test/scala/org/apache/comet/CometIcebergWriteDetectionSuite.scala @@ -545,6 +545,20 @@ class CometIcebergWriteDetectionSuite extends CometTestBase with CometIcebergTes } } + test("fall-back: mixed-case data location scheme (native opens the location verbatim)") { + // OpenDAL strips the scheme prefix from a path case-sensitively, so `S3://` cannot be opened + // natively even though `s3://` can. The gate must match the scheme verbatim and decline it. + withDetectionCatalog { dir => + createTable( + dir, + "mixed_case_scheme", + partitionSpec = "", + properties = + Some("'write.data.path'='S3://nonexistent-bucket/iceberg/db/mixed_case_scheme'")) + assertUnsupportedContainsAllowingWriteFailure("mixed_case_scheme", "storage scheme", "S3") + } + } + 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 => diff --git a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala index 5a14f9f592e..21011fae65d 100644 --- a/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala +++ b/spark/src/test/scala/org/apache/comet/rules/CometScanSchemeFallbackSuite.scala @@ -27,11 +27,12 @@ import java.util.UUID import org.apache.commons.io.FileUtils import org.apache.spark.SparkConf import org.apache.spark.sql.{CometTestBase, SaveMode} -import org.apache.spark.sql.comet.CometScanExec +import org.apache.spark.sql.comet.{CometIcebergNativeScanExec, CometIcebergWriteExec, CometScanExec} import org.apache.spark.sql.execution.{FileSourceScanExec, SparkPlan} -import org.apache.comet.CometConf -import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem} +import org.apache.comet.{CometConf, CometIcebergTestBase, ExtendedExplainInfo, NativeBase} +import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem, FakeWasbSchemeFileSystem} +import org.apache.comet.iceberg.IcebergStorageSchemes /** * Comet's native readers go through object_store, which understands a fixed set of URL schemes. A @@ -43,7 +44,7 @@ import org.apache.comet.hadoop.fs.{FakeHDFSFileSystem, FakeHdfsSchemeFileSystem} * native probe is cached per probe URL so an authorityless URL can't poison the authority-bearing * form of the same scheme. */ -class CometScanSchemeFallbackSuite extends CometTestBase { +class CometScanSchemeFallbackSuite extends CometTestBase with CometIcebergTestBase { private var fakeRootDir: File = _ @@ -54,6 +55,9 @@ class CometScanSchemeFallbackSuite extends CometTestBase { // cluster. `hdfs` is natively readable by default, so this scan must be CLAIMED, not declined. conf.set("spark.hadoop.fs.hdfs.impl", "org.apache.comet.hadoop.fs.FakeHdfsSchemeFileSystem") conf.set("spark.hadoop.fs.defaultFS", FakeHDFSFileSystem.PREFIX) + // Back `wasb` with a local FS so an Iceberg table can live under a `wasb://` warehouse. The + // native Iceberg storage factory has no arm for it, so that scan must be DECLINED, not claimed. + conf.set("spark.hadoop.fs.wasb.impl", classOf[FakeWasbSchemeFileSystem].getName) // Intentionally NOT setting CometConf.COMET_LIBHDFS_SCHEMES -- `fake` is not natively readable, // and `hdfs` must still be claimed by default (mirrors the native `is_hdfs_scheme` default). conf @@ -120,10 +124,65 @@ class CometScanSchemeFallbackSuite extends CometTestBase { "gs://bucket/... must be admitted even after an authorityless gs:// URI was probed first") } + test("iceberg gate: the JNI scheme list parser trims, lowercases and rejects a blank list") { + assert( + IcebergStorageSchemes.parse("file, S3,,s3a,", forWrite = false) == Set("file", "s3", "s3a")) + // A loaded library publishes a fixed non-empty list, so a blank answer is a build bug and + // must fail loudly rather than quietly disable every native Iceberg scan and write. + Seq(("", false), (" , ", true), (null, false)).foreach { case (joined, forWrite) => + val path = if (forWrite) "write" else "read" + val e = intercept[IllegalStateException](IcebergStorageSchemes.parse(joined, forWrite)) + assert( + e.getMessage == s"Comet native library published no Iceberg $path scheme list", + s"unexpected message for ${Option(joined)}: ${e.getMessage}") + } + } + + test("iceberg gate: an unloaded native library yields no schemes") { + // Every caller sits behind `isCometLoaded`, so this answer never reaches a plan; it only + // guarantees that nothing hand-maintained stands in for the native list. + assert(IcebergStorageSchemes.load(forWrite = false, isLoaded = false) == Set.empty) + assert(IcebergStorageSchemes.load(forWrite = true, isLoaded = false) == Set.empty) + } + + test("iceberg gate: the scheme sets come from native and carry its mode-specific entries") { + // The lazy sets must be what the JNI probe answers, not an empty set from a load that ran + // before the library was ready. `oss` is read-only (no `oss.*` property forwarding for + // writes) and `memory` is write-only (a fresh in-process store per FileIO, so a read finds + // nothing); the Azure schemes and `gcs` have no storage factory arm at all. + assume(NativeBase.isLoaded, "Comet native library not loaded") + val read = IcebergStorageSchemes.read + val write = IcebergStorageSchemes.write + assert(read == IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(false), false)) + assert(write == IcebergStorageSchemes.parse(NativeBase.icebergStorageSchemes(true), true)) + assert(read.nonEmpty, "native published no read schemes") + assert(write.nonEmpty, "native published no write schemes") + assert(read.contains("oss"), s"oss must be in the native read schemes ${read.toSeq.sorted}") + assert( + !write.contains("oss"), + s"oss must not be in the native write schemes ${write.toSeq.sorted}") + assert( + write.contains("memory"), + s"memory must be in the native write schemes ${write.toSeq.sorted}") + assert( + !read.contains("memory"), + s"memory must not be in the native read schemes ${read.toSeq.sorted}") + Seq("abfs", "abfss", "wasb", "wasbs", "gcs").foreach { scheme => + assert( + !read.contains(scheme), + s"$scheme must not be in the native read schemes ${read.toSeq.sorted}") + assert( + !write.contains(scheme), + s"$scheme must not be in the native write schemes ${write.toSeq.sorted}") + } + } + test("iceberg gate: builtin allowlist admitted, unbuildable schemes rejected") { - // The Iceberg gate is an explicit allowlist mirroring `storage_factory_for`'s arms (keep in - // lockstep). Narrower than the Parquet gate: object_store recognizes http/abfs/wasb but - // iceberg-rust can't build them, so reject up-front rather than fail during native setup. + // The Iceberg gate is the allowlist the native storage factory publishes over JNI. Narrower + // than the Parquet gate: object_store recognizes http and the Azure schemes, but the native + // Iceberg storage factory has no arm for them, so reject up-front rather than fail at + // execution. `memory` is write-only (a fresh in-process store per FileIO holds no table), and + // the built-in set matches verbatim because native opens a location by its raw scheme. Seq( "file:///tmp/key.parquet", "s3://bucket/key.parquet", @@ -135,6 +194,9 @@ class CometScanSchemeFallbackSuite extends CometTestBase { s"$u must be iceberg-readable; icebergReadableSchemes has regressed") } Seq( + "memory:///key.parquet", + "S3://bucket/key.parquet", + "File:///tmp/key.parquet", "http://bucket.example.com/key.parquet", "https://bucket.example.com/key.parquet", "abfs://container@acct/key.parquet", @@ -143,10 +205,31 @@ class CometScanSchemeFallbackSuite extends CometTestBase { "wasbs://container@acct/key.parquet").foreach { u => assert( !CometScanRule.isIcebergReadableScheme(new URI(u), Set.empty), - s"$u must not be iceberg-readable; storage_factory_for has no matching arm") + s"$u must not be iceberg-readable; storage_factory_for rejects it for reads") } } + test("iceberg gate: an opt-in alias is matched verbatim, like the built-in schemes") { + // Native opens an alias location as written and OpenDAL's S3 backend checks it against a + // lowercase `scheme://bucket/` prefix, so a `BLOB://` location admitted here would only fail + // at execution. The alias set is lowercase, so a scheme written any other way is declined. + val schemes = Set("blob") + assert( + CometScanRule.isIcebergReadableScheme(new URI("blob://bucket/key.parquet"), schemes), + "blob://bucket/... must be admitted once blob is opted in") + Seq("BLOB://bucket/key.parquet", "Blob://bucket/key.parquet", "BLOB:///bucket/key.parquet") + .foreach { u => + assert( + !CometScanRule.isIcebergReadableScheme(new URI(u), schemes), + s"$u must be declined: native opens the location as written and the S3 backend " + + "rejects a scheme prefix that is not lowercase") + } + // The fallback reason names the lowercase alias, which is the form the scan would admit. + assert( + CometScanRule.icebergSupportedSchemesMessage(schemes).contains("blob"), + "the supported-schemes message must list the opted-in alias in its admitted form") + } + test("parquet gate: mixed-bucket alias scan is declined (single object store per partition)") { // Native planning registers one object store per FilePartition and strips the authority from // every file's object key, so files in a second bucket would be read from the first. A scan @@ -218,6 +301,67 @@ class CometScanSchemeFallbackSuite extends CometTestBase { "without opt-in, blob gets no bucket promotion, so a hostless blob location is unopenable") } + test("iceberg scan and write decline a wasb table instead of failing at execution") { + // Regression guard for both scheme gates: the native Iceberg storage factory has no arm for + // `wasb`, so the scan and the write must each be declined with a reason naming the scheme + // and run in Spark, rather than be claimed and fail at execution with "Unsupported storage + // scheme: wasb". + assume(icebergAvailable, "Iceberg not available in classpath") + withTempIcebergDir { dir => + val warehouse = s"${FakeWasbSchemeFileSystem.PREFIX}${dir.getAbsolutePath}/warehouse" + withSQLConf( + "spark.sql.catalog.wasb_catalog" -> "org.apache.iceberg.spark.SparkCatalog", + "spark.sql.catalog.wasb_catalog.type" -> "hadoop", + "spark.sql.catalog.wasb_catalog.warehouse" -> warehouse, + CometConf.COMET_ENABLED.key -> "true", + CometConf.COMET_EXEC_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_NATIVE_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_NATIVE_WRITE_ENABLED.key -> "true", + CometConf.COMET_ICEBERG_WRITE_SPLIT_OPERATOR_ENABLED.key -> "true", + // The VALUES source must be native, or the write is declined for its input before the + // write serde (and its scheme gate) is ever consulted. + CometConf.COMET_EXEC_LOCAL_TABLE_SCAN_ENABLED.key -> "true") { + spark.sql("CREATE TABLE wasb_catalog.db.wasb_table (id INT, name STRING) USING iceberg") + try { + val insertPlans = capturePlans(spark) { + spark.sql( + "INSERT INTO wasb_catalog.db.wasb_table VALUES (1, 'Alice'), (2, 'Bob'), (3, 'Carol')") + } + assert(insertPlans.nonEmpty, "no INSERT plan was captured") + val nativeWrites = insertPlans.flatMap(_.collectWithSubqueries { + case w: CometIcebergWriteExec => w + }) + assert( + nativeWrites.isEmpty, + "`wasb://` has no native Iceberg storage factory arm; the write must fall back to " + + s"Spark, but Comet claimed it:\n${insertPlans.mkString("\n--\n")}") + val writeReasons = insertPlans.flatMap(new ExtendedExplainInfo().getFallbackReasons) + assert( + writeReasons.exists(_.contains("unsupported storage scheme: wasb")), + s"the write must be declined for its scheme, got reasons: $writeReasons") + + val (_, cometPlan) = checkSparkAnswerAndFallbackReason( + "SELECT * FROM wasb_catalog.db.wasb_table ORDER BY id", + "wasb") + val nativeScans = cometPlan.collect { case s: CometIcebergNativeScanExec => s } + assert( + nativeScans.isEmpty, + "`wasb://` has no native Iceberg storage factory arm; the scan must fall back to " + + s"Spark, but Comet claimed it:\n$cometPlan") + // Nothing but the scheme may have caused the fallback: every reason names `wasb`, + // except the one the operator above records about its now-Spark BatchScanExec child. + val scanReasons = new ExtendedExplainInfo().getFallbackReasons(cometPlan) + val unrelated = scanReasons.filterNot(_.contains("wasb")) + assert( + unrelated.forall(_.contains("BatchScanExec")), + s"fallback reasons other than the scheme were recorded: $unrelated") + } finally { + spark.sql("DROP TABLE wasb_catalog.db.wasb_table") + } + } + } + } + test("native scan claims hdfs:// when libhdfs.schemes is unset (native-default lockstep)") { // Native `is_hdfs_scheme` treats `hdfs` as readable when `fs.comet.libhdfs.schemes` is unset, // so the JVM gate must agree and CLAIM `hdfs://`. Guards the `case None => Set("hdfs")` default