mixermt commented on code in PR #6106:
URL: https://github.com/apache/datafusion-comet/pull/6106#discussion_r4115677484


##########
native/core/src/execution/operators/iceberg_common.rs:
##########
@@ -78,36 +82,172 @@ pub(crate) fn storage_factory_for(
         // promotes a HOSTLESS `blob:///bucket/key` into the host at the open 
boundary -- see
         // s3_blob_fs_support for why that never touches the recorded 
delete-matching string.
         s if is_s3_family_scheme(s, catalog_properties) => {
-            let customized_credential_load =
+            let (customized_credential_load, cacheable) =
                 build_s3_credential_loader(path, catalog_properties, 
catalog_name, access_mode)?;
-            if is_s3_compliant_alias_scheme(s, catalog_properties) {
-                Ok(Arc::new(BlobHostPromotingS3StorageFactory::new(
-                    customized_credential_load,
-                )))
-            } else {
-                Ok(Arc::new(OpenDalStorageFactory::S3 {
-                    customized_credential_load,
-                }))
-            }
+            let factory: Arc<dyn StorageFactory> =
+                if is_s3_compliant_alias_scheme(s, catalog_properties) {
+                    Arc::new(BlobHostPromotingS3StorageFactory::new(
+                        customized_credential_load,
+                    ))
+                } else {
+                    Arc::new(OpenDalStorageFactory::S3 {
+                        customized_credential_load,
+                    })
+                };
+            Ok((factory, cacheable))
         }
         _ => Err(DataFusionError::Execution(format!(
             "Unsupported storage scheme: {scheme}"
         ))),
     }
 }
 
+#[derive(Clone, Debug, PartialEq, Eq, Hash)]

Review Comment:
   Replaced the derive with a manual `Debug` that prints `access_mode`, 
`catalog_name` and `reference_path` and omits `properties`. 
`cache_key_debug_does_not_leak_properties` asserts the exact output from a key 
built with vended secrets.



##########
native/core/src/execution/operators/iceberg_common.rs:
##########
@@ -78,36 +82,172 @@ pub(crate) fn storage_factory_for(
         // promotes a HOSTLESS `blob:///bucket/key` into the host at the open 
boundary -- see
         // s3_blob_fs_support for why that never touches the recorded 
delete-matching string.
         s if is_s3_family_scheme(s, catalog_properties) => {
-            let customized_credential_load =
+            let (customized_credential_load, cacheable) =
                 build_s3_credential_loader(path, catalog_properties, 
catalog_name, access_mode)?;
-            if is_s3_compliant_alias_scheme(s, catalog_properties) {
-                Ok(Arc::new(BlobHostPromotingS3StorageFactory::new(
-                    customized_credential_load,
-                )))
-            } else {
-                Ok(Arc::new(OpenDalStorageFactory::S3 {
-                    customized_credential_load,
-                }))
-            }
+            let factory: Arc<dyn StorageFactory> =
+                if is_s3_compliant_alias_scheme(s, catalog_properties) {
+                    Arc::new(BlobHostPromotingS3StorageFactory::new(
+                        customized_credential_load,
+                    ))
+                } else {
+                    Arc::new(OpenDalStorageFactory::S3 {
+                        customized_credential_load,
+                    })
+                };
+            Ok((factory, cacheable))
         }
         _ => Err(DataFusionError::Execution(format!(
             "Unsupported storage scheme: {scheme}"
         ))),
     }
 }
 
+#[derive(Clone, Debug, PartialEq, Eq, Hash)]
+struct FileIoCacheKey {
+    access_mode: u8,
+    catalog_name: String,
+    /// The full path: the S3 access bridge is scoped to the exact path it was 
built for.
+    reference_path: String,
+    properties: Vec<(String, String)>,
+}
+
+impl FileIoCacheKey {
+    /// `None` for `memory:///`, whose namespace must stay private to its task.
+    fn new(
+        catalog_properties: &HashMap<String, String>,
+        reference_path: &str,
+        catalog_name: &str,
+        access_mode: AccessMode,
+    ) -> Option<Self> {
+        if scheme_of(reference_path) == "memory" {
+            return None;
+        }
+        let mut properties: Vec<(String, String)> = catalog_properties
+            .iter()
+            .map(|(k, v)| (k.clone(), v.clone()))
+            .collect();
+        properties.sort();
+        Some(Self {
+            access_mode: access_mode as u8,
+            catalog_name: catalog_name.to_string(),
+            reference_path: reference_path.to_string(),
+            properties,
+        })
+    }
+}
+
+const FILE_IO_CACHE_CAPACITY: usize = 64;
+
+/// Least recently used entries are evicted first.
+struct FileIoCache {

Review Comment:
   Resolving: `closeAll` has only the shutdown hook and the dispatcher test as 
callers, and that test runs in a JVM that never populates the cache.



##########
native/core/src/execution/operators/iceberg_common.rs:
##########
@@ -78,36 +82,172 @@ pub(crate) fn storage_factory_for(
         // promotes a HOSTLESS `blob:///bucket/key` into the host at the open 
boundary -- see
         // s3_blob_fs_support for why that never touches the recorded 
delete-matching string.
         s if is_s3_family_scheme(s, catalog_properties) => {
-            let customized_credential_load =
+            let (customized_credential_load, cacheable) =
                 build_s3_credential_loader(path, catalog_properties, 
catalog_name, access_mode)?;
-            if is_s3_compliant_alias_scheme(s, catalog_properties) {
-                Ok(Arc::new(BlobHostPromotingS3StorageFactory::new(
-                    customized_credential_load,
-                )))
-            } else {
-                Ok(Arc::new(OpenDalStorageFactory::S3 {
-                    customized_credential_load,
-                }))
-            }
+            let factory: Arc<dyn StorageFactory> =
+                if is_s3_compliant_alias_scheme(s, catalog_properties) {
+                    Arc::new(BlobHostPromotingS3StorageFactory::new(
+                        customized_credential_load,
+                    ))
+                } else {
+                    Arc::new(OpenDalStorageFactory::S3 {
+                        customized_credential_load,
+                    })
+                };
+            Ok((factory, cacheable))
         }
         _ => Err(DataFusionError::Execution(format!(
             "Unsupported storage scheme: {scheme}"
         ))),
     }
 }
 
+#[derive(Clone, Debug, PartialEq, Eq, Hash)]
+struct FileIoCacheKey {
+    access_mode: u8,
+    catalog_name: String,
+    /// The full path: the S3 access bridge is scoped to the exact path it was 
built for.
+    reference_path: String,
+    properties: Vec<(String, String)>,
+}
+
+impl FileIoCacheKey {
+    /// `None` for `memory:///`, whose namespace must stay private to its task.
+    fn new(
+        catalog_properties: &HashMap<String, String>,
+        reference_path: &str,
+        catalog_name: &str,
+        access_mode: AccessMode,
+    ) -> Option<Self> {
+        if scheme_of(reference_path) == "memory" {
+            return None;
+        }
+        let mut properties: Vec<(String, String)> = catalog_properties
+            .iter()
+            .map(|(k, v)| (k.clone(), v.clone()))
+            .collect();
+        properties.sort();
+        Some(Self {
+            access_mode: access_mode as u8,
+            catalog_name: catalog_name.to_string(),
+            reference_path: reference_path.to_string(),
+            properties,
+        })
+    }
+}
+
+const FILE_IO_CACHE_CAPACITY: usize = 64;
+
+/// Least recently used entries are evicted first.
+struct FileIoCache {
+    entries: HashMap<FileIoCacheKey, FileIO>,
+    order: VecDeque<FileIoCacheKey>,
+    capacity: usize,
+}
+
+impl FileIoCache {
+    fn new(capacity: usize) -> Self {
+        Self {
+            entries: HashMap::new(),
+            order: VecDeque::new(),
+            capacity,
+        }
+    }
+
+    fn get(&mut self, key: &FileIoCacheKey) -> Option<FileIO> {
+        let file_io = self.entries.get(key)?.clone();
+        if let Some(pos) = self.order.iter().position(|k| k == key) {
+            let recent = self.order.remove(pos)?;
+            self.order.push_back(recent);
+        }
+        Some(file_io)
+    }
+
+    /// Returns the replaced or evicted `FileIO` so the caller can drop it 
outside the lock.
+    fn insert(&mut self, key: FileIoCacheKey, file_io: FileIO) -> 
Option<FileIO> {
+        if let Some(previous) = self.entries.insert(key.clone(), file_io) {
+            return Some(previous);
+        }
+        self.order.push_back(key);
+        if self.entries.len() > self.capacity {
+            let oldest = self.order.pop_front()?;
+            return self.entries.remove(&oldest);
+        }
+        None
+    }
+}
+
+/// Shared per executor so tasks reuse one FileIO: its factory, parsed config 
and access bridge,
+/// plus the storage client where the backend caches operators.
+static FILE_IO_CACHE: LazyLock<Mutex<FileIoCache>> =

Review Comment:
   Tracking in #5898. For that rebase: the storage factory arms now return a 
`cacheable` flag, so the HDFS arm should choose it deliberately; `false` keeps 
the client task-scoped.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


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

Reply via email to