Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
41 commits
Select commit Hold shift + click to select a range
4b5a6db
feat: publish the native Iceberg storage scheme list over JNI
dwsmith1983 Sep 20, 2026
d0a6bde
fix: load the Iceberg gate scheme lists from the native storage factory
dwsmith1983 Sep 20, 2026
7cbce96
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 20, 2026
ab9288b
fix: gate the Iceberg scheme lists on the native list and match schem…
dwsmith1983 Sep 20, 2026
dc53313
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 21, 2026
b8fd751
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 22, 2026
56ea017
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 22, 2026
3e4f309
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 22, 2026
1d3f95d
fix: drop the fallback scheme lists and match Iceberg aliases verbatim
dwsmith1983 Sep 22, 2026
458a971
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 22, 2026
5b3b77d
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 23, 2026
aecaec3
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 23, 2026
b428e36
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 24, 2026
a80dbf5
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 24, 2026
3d829e6
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 24, 2026
99d6451
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 24, 2026
b7b401b
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 24, 2026
bbfa4f7
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 24, 2026
07330c4
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 25, 2026
a6e1096
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 25, 2026
f9c8d29
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 25, 2026
9a7c092
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 25, 2026
a634831
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 25, 2026
a4841bc
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 25, 2026
c347457
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 26, 2026
5286dd4
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 26, 2026
7469e44
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 27, 2026
c15100a
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Sep 27, 2026
af0923c
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 27, 2026
ab39d1b
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 28, 2026
879e13a
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 28, 2026
5d94024
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 29, 2026
6a3dd6b
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 29, 2026
01747cf
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 29, 2026
3be88c6
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 29, 2026
fc3b64c
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 30, 2026
05bf9eb
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 30, 2026
3f56a17
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Sep 30, 2026
5913c79
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Oct 1, 2026
5c91b64
Merge remote-tracking branch 'origin/main' into fix/iceberg-scheme-so…
dwsmith1983 Oct 3, 2026
8fcaeaf
Merge branch 'main' into fix/iceberg-scheme-source-of-truth
dwsmith1983 Oct 3, 2026
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
5 changes: 4 additions & 1 deletion docs/source/user-guide/latest/datasources.md
Original file line number Diff line number Diff line change
Expand Up @@ -265,7 +265,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
Expand Down
197 changes: 171 additions & 26 deletions native/core/src/execution/operators/iceberg_common.rs
Original file line number Diff line number Diff line change
Expand Up @@ -47,16 +47,12 @@ const ICEBERG_PROVIDER_CLASS_PROPERTY: &str = "s3.comet.credential.provider.clas
/// `CometS3CredentialBridge` can read whatever the vendor needs.
const STORAGE_PROPERTY_PREFIXES: &[&str] = &["s3.", "gcs.", "adls.", "client."];

/// 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.
Expand All @@ -66,24 +62,22 @@ pub(crate) fn storage_factory_for(
catalog_name: &str,
access_mode: AccessMode,
) -> Result<(Arc<dyn StorageFactory>, 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
Expand All @@ -92,7 +86,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<dyn StorageFactory> = match access {
Expand All @@ -110,6 +104,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}"
))),
Expand Down Expand Up @@ -302,6 +297,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<String, String>,
reference_path: &str,
Expand Down Expand Up @@ -582,7 +588,19 @@ 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<String, String>) -> 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<String, String>) -> bool {
!scheme.bytes().any(|b| b.is_ascii_uppercase())
&& is_s3_compliant_alias_scheme(scheme, catalog_properties)
}

#[cfg(test)]
Expand Down Expand Up @@ -807,21 +825,111 @@ 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]
fn common_schemes_resolve_for_both_modes() {
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();
Expand Down Expand Up @@ -853,6 +961,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
Expand Down
25 changes: 25 additions & 0 deletions native/core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ extern crate datafusion_comet_jni_bridge;

use jni::{
objects::{JClass, JString},
sys::{jboolean, jstring, JNI_FALSE},
EnvUnowned,
};
use log::info;
Expand All @@ -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 execution;
Expand Down Expand Up @@ -247,6 +251,27 @@ pub extern "system" fn Java_org_apache_comet_NativeBase_getTzdataVersion(
})
}

/// 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<Config> {
let console_append = ConsoleAppender::builder()
Expand Down
6 changes: 4 additions & 2 deletions native/core/src/parquet/objectstore/s3_blob_fs_support.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, String>,
Expand Down
12 changes: 12 additions & 0 deletions spark/src/main/java/org/apache/comet/NativeBase.java
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Loading