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


##########
crates/paimon/src/arrow/format/blob.rs:
##########
@@ -1688,36 +1677,68 @@ async fn read_blob_range(reader: &dyn FileRead, range: 
Range<u64>) -> crate::Res
     Ok(bytes)
 }
 
-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);
+        let key_heap_bytes = cache_key.as_ref().map_or(0, String::capacity);
+        let load = cache
+            .entries
+            .get_or_try_insert_with_admission(
+                cache_key,

Review Comment:
   [P1] Reset process-local BLOB cache state before joining inherited loads
   
   If a parent thread is awaiting a BLOB index read when another thread calls 
`fork()`, the child inherits this key's pending `OnceCell` with its 
initialization permit held. A new child read using the inherited Table/FileIO 
and preplanned splits joins that entry and waits indefinitely: the parent's 
initializer cannot finish in the child. This is reachable from Python because 
TableRead releases the GIL while collecting, and the PID-aware runtime creates 
a fresh runtime without resetting FileIO caches. For matching schema IDs with 
no deletion vector/predicate and local block caching disabled, BLOB readers use 
the manifest's file size without a prior stat; Jindo reader creation is lazy, 
so #999's range-read PID check is never reached either.
   
   I reproduced this with an actual `libc::fork()` and a process-bound test 
FileRead on the current combined heads: with `cache.blob-index.max-size=0`, the 
child immediately returns the fork-specific error; with the default `64 MiB`, 
the same child load times out after 500 ms without reaching the reader. This is 
a synthetic reader test, not a real-SDK or end-to-end Python run. Normal 
direct/wrapped cached failures still preserve fork classification, so the 
earlier error-cloning fix does not prevent this wait. The baseline BLOB cache 
inserts only completed indexes and has no pending entry to inherit. Please 
select a fresh per-process cache or bypass inherited cache contexts before 
touching their OnceLock/entry synchronization, and add a bounded 
fork-during-load regression.



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