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]