This is an automated email from the ASF dual-hosted git repository.

sunchao pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/datafusion-comet.git


The following commit(s) were added to refs/heads/main by this push:
     new bd1aa66fed fix: isolate object-store registration by backend and 
configuration (#5503)
bd1aa66fed is described below

commit bd1aa66fed75d103b03f62a7de1fcd03479e83c6
Author: Chao Sun <[email protected]>
AuthorDate: Thu Sep 10 11:17:44 2026 -0700

    fix: isolate object-store registration by backend and configuration (#5503)
    
    * fix: isolate object-store registration by backend and configuration
    
    * docs: clarify object-store cache and registration identity
    
    * test: align S3 routing checks with isolated store URLs
---
 native/core/src/parquet/parquet_exec.rs    |  71 ++++++-
 native/core/src/parquet/parquet_support.rs | 287 +++++++++++++++++++++++++++--
 2 files changed, 340 insertions(+), 18 deletions(-)

diff --git a/native/core/src/parquet/parquet_exec.rs 
b/native/core/src/parquet/parquet_exec.rs
index 8796cb2324..96c762f8dd 100644
--- a/native/core/src/parquet/parquet_exec.rs
+++ b/native/core/src/parquet/parquet_exec.rs
@@ -222,6 +222,25 @@ pub(crate) fn init_datasource_exec(
     Ok(data_source_exec)
 }
 
+// Registration URLs use a reserved suffix to distinguish backend/configuration
+// identities. Encryption must use the physical URI that Spark registered. 
Match
+// the complete suffix, from the right, so custom Hadoop schemes are preserved.
+fn physical_object_store_scheme(object_store_url: &ObjectStoreUrl) -> &str {
+    let store_url: &url::Url = object_store_url.as_ref();
+    let scheme = store_url.scheme();
+    if let Some((physical_scheme, identity)) = scheme.rsplit_once("+comet-") {
+        if let Some((hash, backend)) = identity.split_once('-') {
+            if hash.len() == 16
+                && hash.bytes().all(|byte| byte.is_ascii_hexdigit())
+                && matches!(backend, "native" | "hdfs")
+            {
+                return physical_scheme;
+            }
+        }
+    }
+    scheme
+}
+
 #[allow(clippy::too_many_arguments)]
 fn get_options(
     session_timezone: &str,
@@ -282,10 +301,15 @@ fn get_options(
     spark_parquet_options.allow_timestamp_ltz_to_ntz = 
allow_timestamp_ltz_to_ntz;
 
     if encryption_enabled {
+        let store_url: &url::Url = object_store_url.as_ref();
         table_parquet_options.crypto.configure_factory(
             ENCRYPTION_FACTORY_ID,
             &CometEncryptionConfig {
-                uri_base: object_store_url.to_string(),
+                uri_base: format!(
+                    "{}://{}/",
+                    physical_object_store_scheme(object_store_url),
+                    
&store_url[url::Position::BeforeHost..url::Position::AfterPort],
+                ),
             },
         );
     }
@@ -307,6 +331,51 @@ mod tests {
     use parquet::file::properties::{EnabledStatistics, WriterProperties};
     use std::fs::File;
 
+    #[test]
+    fn preserves_physical_uri_for_isolated_encrypted_object_stores() {
+        fn encryption_uri(url: &str) -> String {
+            let object_store_url = ObjectStoreUrl::parse(url).unwrap();
+            let (options, _) = get_options(
+                "UTC",
+                true,
+                false,
+                false,
+                false,
+                &object_store_url,
+                true,
+                &ParquetOptions::default(),
+            );
+            let encryption: CometEncryptionConfig = options
+                .crypto
+                .factory_options
+                .to_extension_options()
+                .unwrap();
+            encryption.uri_base
+        }
+
+        // Native s3a is canonicalized to s3 before registration; 
Hadoop-selected
+        // s3a keeps its physical spelling. Both must strip only the identity 
suffix.
+        for scheme in ["s3", "s3a", "abfss", "custom+comet-existing"] {
+            let normal = format!("{scheme}://bucket:9000/");
+            for backend in ["native", "hdfs"] {
+                let isolated = 
format!("{scheme}+comet-0123456789abcdef-{backend}://bucket:9000/");
+                assert_eq!(encryption_uri(&isolated), encryption_uri(&normal));
+                assert_eq!(encryption_uri(&isolated), normal);
+            }
+        }
+        assert_eq!(encryption_uri("file:///"), "file:///");
+        assert_eq!(
+            encryption_uri("file+comet-0123456789abcdef-hdfs:///"),
+            "file:///"
+        );
+        // A physical custom scheme containing a similar, incomplete suffix is 
not
+        // itself a synthetic registration URL.
+        for scheme in ["custom+comet-name", "custom+comet-1234-native"] {
+            let normal = format!("{scheme}://bucket/");
+            assert_eq!(encryption_uri(&normal), normal);
+        }
+    }
+
     // Regression test for #4990: a fresh `TableParquetOptions::new()` ignored 
session-level
     // `datafusion.execution.parquet.*` settings entirely, so 
`spark.comet.datafusion.
     // execution.parquet.*` (behind `respectDataFusionConfigs`) and 
`spark.comet.parquet.
diff --git a/native/core/src/parquet/parquet_support.rs 
b/native/core/src/parquet/parquet_support.rs
index dd7eb40266..fb92965def 100644
--- a/native/core/src/parquet/parquet_support.rs
+++ b/native/core/src/parquet/parquet_support.rs
@@ -609,9 +609,16 @@ fn create_hdfs_object_store(
     })
 }
 
-type ObjectStoreCache = RwLock<HashMap<(String, u64), Arc<dyn ObjectStore>>>;
-
-/// Process-wide cache of object stores, keyed by `(scheme://host:port, 
config_hash)`.
+/// Cache identity: `(scheme://host:port, config_hash, hdfs_backend)`.
+/// Native `s3a` is normalized to `s3`; Hadoop-selected schemes keep their 
spelling.
+/// The hash covers the object-store configuration. The boolean is `true` for 
the
+/// Hadoop backend (including custom schemes routed through Hadoop), `false` 
for native.
+type ObjectStoreCacheKey = (String, u64, bool);
+type ObjectStoreCache = RwLock<HashMap<ObjectStoreCacheKey, Arc<dyn 
ObjectStore>>>;
+
+/// Process-wide cache keyed by `(physical_scheme://host:port, config_hash, 
hdfs_backend)`.
+/// Backend identity is separate from the normalized URL: a configuration can 
route `s3`
+/// through Hadoop while native `s3a` is normalized to the same `s3` scheme.
 ///
 /// ## Why static / process lifetime?
 ///
@@ -627,7 +634,7 @@ type ObjectStoreCache = RwLock<HashMap<(String, u64), 
Arc<dyn ObjectStore>>>;
 ///
 /// ## Unbounded size
 ///
-/// Cache entries are indexed by `(scheme://host:port, hash-of-configs)`.  A 
typical Spark
+/// Cache entries include the physical URL, configuration hash and backend. A 
typical Spark
 /// job accesses a small, fixed set of buckets with a stable configuration, so 
the number of
 /// distinct keys is O(buckets × credential-configs) and remains small 
throughout the job.
 /// Entries are cheap relative to the cost of creating a new object store (new 
HTTP
@@ -682,7 +689,7 @@ pub(crate) fn prepare_object_store_with_configs(
     );
 
     let config_hash = hash_object_store_configs(object_store_configs);
-    let cache_key = (url_key.clone(), config_hash);
+    let cache_key = (url_key.clone(), config_hash, is_hdfs_scheme);
 
     // Check the cache first to reuse existing object store instances.
     // This enables HTTP connection pooling and avoids redundant DNS lookups.
@@ -720,28 +727,266 @@ pub(crate) fn prepare_object_store_with_configs(
             (store, path)
         };
 
-    let object_store_url = ObjectStoreUrl::parse(url_key.clone())?;
-    runtime_env.register_object_store(&url, object_store);
+    // A RuntimeEnv can plan multiple scans with different backends or 
credentials
+    // for the same bucket. Use the same identity as the cache, even for the 
first
+    // registration, so neither later registration nor planning order changes 
the
+    // store used by an existing scan. Native s3/s3a share the normalized s3 
scheme;
+    // a Hadoop-selected scheme retains its physical spelling.
+    //
+    // Native LocalFileSystem ignores these Hadoop options and keeps file:// 
for
+    // compatibility. An explicitly Hadoop-routed file scheme is still 
isolated.
+    let object_store_url = if scheme == "file" && !is_hdfs_scheme {
+        ObjectStoreUrl::parse(url_key)?
+    } else {
+        let backend = if is_hdfs_scheme { "hdfs" } else { "native" };
+        // DataFusion keys stores only by scheme and authority, so put 
configuration
+        // and backend identity in the scheme while preserving the physical 
authority.
+        // `+comet-` marks our internal registration suffix; encryption lookup 
strips
+        // the complete suffix to recover the physical URI.
+        ObjectStoreUrl::parse(format!(
+            "{scheme}+comet-{config_hash:016x}-{backend}://{}",
+            &url[url::Position::BeforeHost..url::Position::AfterPort],
+        ))?
+    };
+    runtime_env.register_object_store(object_store_url.as_ref(), object_store);
     Ok((object_store_url, object_store_path))
 }
 
 #[cfg(test)]
 mod tests {
-    #[cfg(not(feature = "hdfs-opendal"))]
+    use super::{hash_object_store_configs, object_store_cache, 
prepare_object_store_with_configs};
+    use bytes::Bytes;
     use datafusion::execution::object_store::ObjectStoreUrl;
-    #[cfg(not(feature = "hdfs-opendal"))]
     use datafusion::execution::runtime_env::RuntimeEnv;
-    #[cfg(not(feature = "hdfs-opendal"))]
+    use object_store::memory::InMemory;
     use object_store::path::Path;
-    #[cfg(not(feature = "hdfs-opendal"))]
+    use object_store::{ObjectStore, ObjectStoreExt};
+    use std::collections::HashMap;
     use std::sync::Arc;
     #[cfg(not(feature = "hdfs-opendal"))]
     use url::Url;
 
     #[cfg(not(feature = "hdfs-opendal"))]
     use crate::execution::operators::ExecutionError;
-    #[cfg(not(feature = "hdfs-opendal"))]
-    use std::collections::HashMap;
+
+    struct StoreConfig {
+        input_scheme: &'static str,
+        physical_scheme: &'static str,
+        hdfs_backend: bool,
+        options: HashMap<String, String>,
+    }
+
+    // Seed distinct stores so these tests exercise cache lookup, registration 
and
+    // actual reads without needing cloud credentials or a running Hadoop 
cluster.
+    async fn check_isolated_stores(bucket: &str, cases: [StoreConfig; 2]) {
+        let path = Path::from("directory/part one.parquet");
+        let stores: [Arc<dyn ObjectStore>; 2] =
+            [Arc::new(InMemory::new()), Arc::new(InMemory::new())];
+        let keys = cases.each_ref().map(|case| {
+            (
+                format!("{}://{bucket}", case.physical_scheme),
+                hash_object_store_configs(&case.options),
+                case.hdfs_backend,
+            )
+        });
+        for (index, store) in stores.iter().enumerate() {
+            store
+                .put(&path, Bytes::from(format!("store-{index}")).into())
+                .await
+                .unwrap();
+        }
+        {
+            let mut cache = object_store_cache().write().unwrap();
+            for (key, store) in keys.iter().zip(&stores) {
+                cache.insert(key.clone(), Arc::clone(store));
+            }
+        }
+
+        let mut previous_urls = None;
+        for order in [[0, 1], [1, 0]] {
+            let runtime = Arc::new(RuntimeEnv::default());
+            let mut prepared = [None, None];
+            for index in order {
+                let case = &cases[index];
+                prepared[index] = Some(
+                    prepare_object_store_with_configs(
+                        Arc::clone(&runtime),
+                        format!(
+                            "{}://{bucket}/directory/part%20one.parquet",
+                            case.input_scheme
+                        ),
+                        &case.options,
+                    )
+                    .unwrap(),
+                );
+            }
+            let prepared = prepared.map(Option::unwrap);
+            let urls = prepared.each_ref().map(|(url, _)| url.clone());
+            assert_ne!(urls[0], urls[1]);
+            if let Some(previous) = &previous_urls {
+                assert_eq!(
+                    &urls, previous,
+                    "registration must not depend on planning order"
+                );
+            }
+            previous_urls = Some(urls);
+            for (index, (url, actual_path)) in prepared.iter().enumerate() {
+                assert_eq!(actual_path, &path);
+                let store = runtime.object_store(url).unwrap();
+                assert!(Arc::ptr_eq(&store, &stores[index]));
+                assert_eq!(
+                    
store.get(actual_path).await.unwrap().bytes().await.unwrap(),
+                    Bytes::from(format!("store-{index}")),
+                );
+                assert!(url.as_str().starts_with(&format!(
+                    "{}+comet-{:016x}-{}://",
+                    cases[index].physical_scheme,
+                    keys[index].1,
+                    if cases[index].hdfs_backend {
+                        "hdfs"
+                    } else {
+                        "native"
+                    },
+                )));
+            }
+        }
+        let mut cache = object_store_cache().write().unwrap();
+        for key in keys {
+            cache.remove(&key);
+        }
+    }
+
+    #[tokio::test]
+    async fn isolates_backends_even_when_s3_alias_and_configs_match() {
+        let options = HashMap::from([("fs.comet.libhdfs.schemes".into(), 
"s3".into())]);
+        check_isolated_stores(
+            "comet-isolation-backend-alias",
+            [
+                StoreConfig {
+                    input_scheme: "s3a",
+                    physical_scheme: "s3",
+                    hdfs_backend: false,
+                    options: options.clone(),
+                },
+                StoreConfig {
+                    input_scheme: "s3",
+                    physical_scheme: "s3",
+                    hdfs_backend: true,
+                    options,
+                },
+            ],
+        )
+        .await;
+    }
+
+    #[tokio::test]
+    async fn isolates_native_stores_with_different_configurations() {
+        check_isolated_stores(
+            "comet-isolation-configurations",
+            ["first", "second"].map(|endpoint| StoreConfig {
+                input_scheme: "s3a",
+                physical_scheme: "s3",
+                hdfs_backend: false,
+                options: HashMap::from([("fs.s3a.endpoint".into(), 
endpoint.into())]),
+            }),
+        )
+        .await;
+    }
+
+    #[tokio::test]
+    async fn preserves_custom_hadoop_scheme_when_routing_changes() {
+        check_isolated_stores(
+            "comet-isolation-custom-hadoop",
+            [
+                StoreConfig {
+                    input_scheme: "s3a",
+                    physical_scheme: "s3",
+                    hdfs_backend: false,
+                    options: HashMap::new(),
+                },
+                StoreConfig {
+                    input_scheme: "s3a",
+                    physical_scheme: "s3a",
+                    hdfs_backend: true,
+                    options: 
HashMap::from([("fs.comet.libhdfs.schemes".into(), "s3a".into())]),
+                },
+            ],
+        )
+        .await;
+    }
+
+    #[test]
+    fn native_s3_aliases_share_cache_and_registration_identity() {
+        let options = HashMap::new();
+        let key = (
+            "s3://comet-isolation-native-aliases".to_string(),
+            hash_object_store_configs(&options),
+            false,
+        );
+        let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
+        object_store_cache()
+            .write()
+            .unwrap()
+            .insert(key.clone(), Arc::clone(&store));
+        let mut previous = None;
+        for schemes in [["s3", "s3a"], ["s3a", "s3"]] {
+            let runtime = Arc::new(RuntimeEnv::default());
+            for scheme in schemes {
+                let (url, _) = prepare_object_store_with_configs(
+                    Arc::clone(&runtime),
+                    
format!("{scheme}://comet-isolation-native-aliases/file.parquet"),
+                    &options,
+                )
+                .unwrap();
+                assert!(Arc::ptr_eq(&runtime.object_store(&url).unwrap(), 
&store));
+                assert!(url.as_str().starts_with("s3+comet-"));
+                if let Some(previous) = &previous {
+                    assert_eq!(&url, previous);
+                }
+                previous = Some(url);
+            }
+        }
+        object_store_cache().write().unwrap().remove(&key);
+    }
+
+    #[test]
+    fn keeps_native_file_url_separate_from_explicit_hadoop_file_routing() {
+        let runtime = Arc::new(RuntimeEnv::default());
+        let options = HashMap::from([("fs.comet.libhdfs.schemes".into(), 
"file".into())]);
+        let key = (
+            "file://".to_string(),
+            hash_object_store_configs(&options),
+            true,
+        );
+        let hdfs_store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
+        object_store_cache()
+            .write()
+            .unwrap()
+            .insert(key.clone(), Arc::clone(&hdfs_store));
+        let (hdfs_url, _) = prepare_object_store_with_configs(
+            Arc::clone(&runtime),
+            "file:///comet-isolation-file-routing.parquet".into(),
+            &options,
+        )
+        .unwrap();
+        let (native_url, _) = prepare_object_store_with_configs(
+            Arc::clone(&runtime),
+            "file:///comet-isolation-file-routing.parquet".into(),
+            &HashMap::new(),
+        )
+        .unwrap();
+        assert_eq!(native_url, ObjectStoreUrl::local_filesystem());
+        assert_ne!(native_url, hdfs_url);
+        assert!(Arc::ptr_eq(
+            &runtime.object_store(&hdfs_url).unwrap(),
+            &hdfs_store
+        ));
+        assert!(!Arc::ptr_eq(
+            &runtime.object_store(&native_url).unwrap(),
+            &hdfs_store
+        ));
+        object_store_cache().write().unwrap().remove(&key);
+    }
 
     /// Parses the url, registers the object store, and returns a tuple of the 
object store url and object store path
     #[cfg(not(feature = "hdfs-opendal"))]
@@ -1129,12 +1374,12 @@ mod tests {
         for (input, expected_bucket, expected_path) in [
             (
                 
"blob://test_bucket/comet/spark-warehouse/part-00000.snappy.parquet",
-                "s3://test_bucket",
+                "test_bucket",
                 "/comet/spark-warehouse/part-00000.snappy.parquet",
             ),
             (
                 "blob:///mybucket/warehouse/data/part-0.snappy.parquet",
-                "s3://mybucket",
+                "mybucket",
                 "warehouse/data/part-0.snappy.parquet",
             ),
         ] {
@@ -1146,7 +1391,11 @@ mod tests {
             .unwrap_or_else(|e| panic!("{input} should normalize to s3://: 
{e}"));
             assert_eq!(
                 object_store_url,
-                ObjectStoreUrl::parse(expected_bucket).unwrap()
+                ObjectStoreUrl::parse(format!(
+                    "s3+comet-{:016x}-native://{expected_bucket}",
+                    hash_object_store_configs(&configs),
+                ))
+                .unwrap()
             );
             assert_eq!(path, Path::from(expected_path));
         }
@@ -1190,7 +1439,11 @@ mod tests {
             .unwrap_or_else(|e| panic!("{input} must build an S3 store, not 
libhdfs: {e}"));
             assert_eq!(
                 object_store_url,
-                ObjectStoreUrl::parse("s3://test_bucket").unwrap()
+                ObjectStoreUrl::parse(format!(
+                    "s3+comet-{:016x}-native://test_bucket",
+                    hash_object_store_configs(&configs),
+                ))
+                .unwrap()
             );
             assert_eq!(path, Path::from("/comet/part-00000.snappy.parquet"));
         }


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to