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]