JingsongLi commented on code in PR #998:
URL: https://github.com/apache/paimon-rust/pull/998#discussion_r4164983034


##########
crates/paimon/src/io/file_io.rs:
##########
@@ -198,15 +277,28 @@ impl FileIO {
     /// subsequently created by [`Self::new_input`] and [`Self::new_output`].
     pub fn with_provider(mut self, provider: Arc<dyn FileIOProvider>) -> Self {
         self.backend = FileIOBackend::Provider(provider);
-        self.cache_namespace = next_file_io_cache_namespace();
         self

Review Comment:
   [P1] Restore cache isolation when replacing a FileIO provider
   
   `with_provider()` no longer renews the namespace, and FileIO clones share 
the new BLOB cache context. For providers implementing the existing `create()` 
API, the default namespace hook returns `None`, so replacing a provider retains 
the same fallback namespace. Standard OpenDAL OSS/S3 operators at different 
endpoints can report the same scheme, bucket name, root, and relative key. I 
reproduced this with two real OSS operators serving equal-size files at the 
same bucket/key: after caching `[NULL, value]` from endpoint A, endpoint B's 
`[value, NULL]` file returned NULL for row 0 instead of `"value"`. Previously 
the fresh numeric namespace prevented this BLOB collision. Please preserve 
provider-attachment isolation in the fallback namespace while keeping the 
catalog budget shared, and cover provider replacement with a regression test.



##########
crates/paimon/src/arrow/format/blob.rs:
##########
@@ -1779,6 +1800,26 @@ impl BlobFileIndex {
     }
 }
 
+fn clone_blob_index_error(error: &Error) -> Error {
+    match error {
+        Error::DataInvalid { message, .. } => Error::DataInvalid {
+            message: message.clone(),
+            source: None,
+        },
+        Error::Unsupported { message } => Error::Unsupported {
+            message: message.clone(),
+        },
+        Error::UnexpectedError { message, .. } => Error::UnexpectedError {
+            message: message.clone(),
+            source: None,
+        },
+        _ => Error::UnexpectedError {
+            message: error.to_string(),
+            source: None,

Review Comment:
   [P1] Preserve fork-safety errors when combining this cache with #999
   
   This is an integration issue with #999: `ProcessForkUnsupported` falls 
through to this arm and becomes `UnexpectedError` with no source. The 
`UnexpectedError` arm above also discards a wrapped fork cause. `load_cached()` 
invokes this helper even for the first failing cold load, so rejecting failed 
entries does not preserve the error for the caller. Python then receives 
`ValueError` instead of `ForkSafetyError`, defeating the dedicated interception 
in apache/paimon#10328 and allowing fallback to inherited Jindo state. A 
targeted interface-combination test using this PR's actual cache code and 
#999's `error.rs` reproduced `is_process_fork_unsupported() == false`; this was 
not a full merged-branch test. Please retain fork classification for both 
direct and wrapped failures and test propagation with the BLOB cache enabled.



##########
crates/paimon/src/arrow/format/blob.rs:
##########
@@ -1673,36 +1662,68 @@ struct BlobFileIndex {
     entries: Vec<BlobEntry>,
 }
 
-impl BlobFileIndex {
-    async fn load_cached(
-        reader: &dyn FileRead,
-        file_size: u64,
-        file_path: &str,
-    ) -> crate::Result<Arc<Self>> {
-        let cache_key = reader
-            .cache_namespace()
-            .filter(|_| !file_path.is_empty())
-            .map(|namespace| BlobIndexCacheKey {
-                namespace,
-                file_path: file_path.to_string(),
-            });
-        if let Some(cache_key) = &cache_key {
-            let mut cache = BLOB_INDEX_CACHE
-                .lock()
-                .unwrap_or_else(|error| error.into_inner());
-            if let Some(index) = cache.get(cache_key) {
-                return Ok(index.clone());
-            }
+type BlobIndexLoadResult = Result<Arc<BlobFileIndex>, Arc<Error>>;
+
+struct BlobIndexCache {
+    entries: FileMetadataCache<String, BlobIndexLoadResult>,
+}
+
+impl BlobIndexCache {
+    fn new(max_bytes: usize) -> Self {
+        Self {
+            entries: FileMetadataCache::new(max_bytes, usize::MAX),
         }
+    }
+}
 
-        let index = Arc::new(Self::load(reader, file_size).await?);
-        if let Some(cache_key) = cache_key {
-            BLOB_INDEX_CACHE
-                .lock()
-                .unwrap_or_else(|error| error.into_inner())
-                .put(cache_key, index.clone());
+impl BlobFileIndex {
+    async fn load_cached(reader: &dyn FileRead, file_size: u64) -> 
crate::Result<Arc<Self>> {
+        let Some(context) = reader
+            .blob_index_cache()
+            .and_then(|cache| cache.downcast_ref::<BlobIndexCacheContext>())
+        else {
+            return Ok(Arc::new(Self::load(reader, file_size).await?));
+        };
+        let cache = context.get_or_init(BlobIndexCache::new);
+        let cache_key = reader.cache_key().map(ToOwned::to_owned);

Review Comment:
   [P1] Include the Azure account authority in the BLOB index cache key
   
   Replacing the previous full-URI key with `reader.cache_key()` makes two 
supported Azure paths collide: a single FileIO reading 
`abfss://[email protected]/data.blob` and the 
corresponding path on `account-b` produces identical keys when `azure.endpoint` 
is unset. The namespace hashes only explicitly configured endpoints, while 
OpenDAL's Azure service identity contains the filesystem name/root but not the 
account host. I reproduced the identical keys with the real Azure operators. 
The second file can therefore reuse the first file's decoded index, giving 
incorrect NULL positions, row counts, or byte ranges. The previous BLOB key 
retained the full URI and distinguished these accounts. Please include the 
effective endpoint/account authority in the metadata identity and add a 
cross-account BLOB regression test.



-- 
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]

Reply via email to