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


##########
crates/integrations/datafusion/src/catalog.rs:
##########
@@ -423,15 +646,22 @@ impl CatalogProvider for PaimonCatalogProvider {
         let catalog_name = self.catalog_name.clone();
         let session_state = self.session_state.clone();
         let schema_force_view_types = self.schema_force_view_types;
+        let table_engines = self.table_engines();
+        let metadata = Arc::clone(&self.metadata);
         let name = name.to_string();
         block_on_with_runtime(

Review Comment:
   Could we avoid keeping the blocking bridge in the remaining mutation 
callbacks before closing #872? `register_schema`, `deregister_schema`, and 
`deregister_table` still move the future to the process runtime and 
synchronously join the OS thread. If that future waits for work scheduled on a 
single-worker caller runtime, the caller worker remains blocked and the same 
deadlock class as #872 persists. Moving these catalog mutations to explicit 
async paths, or narrowing the issue scope, would make the claimed supersession 
precise.
   



##########
crates/integrations/datafusion/src/catalog.rs:
##########
@@ -49,6 +53,118 @@ pub(crate) type SessionStateProvider = Arc<dyn Fn() -> 
Option<SessionState> + Se
 /// providers, so registrations stay visible to schemas obtained earlier.
 type TableEngines = Arc<RwLock<HashMap<PaimonTableType, Arc<dyn 
TableEngineResolver>>>>;
 
+#[derive(Clone, Debug, Default)]
+struct CatalogMetadataSnapshot {
+    databases: IndexMap<String, Arc<DatabaseMetadata>>,
+}
+
+#[derive(Clone, Debug, Default)]
+struct DatabaseMetadata {
+    objects: IndexMap<String, TableType>,
+}
+
+#[derive(Debug)]
+struct CatalogMetadataState {
+    snapshot: RwLock<Arc<CatalogMetadataSnapshot>>,
+    next_generation: AtomicU64,
+    published_generation: AtomicU64,
+}
+
+impl Default for CatalogMetadataState {
+    fn default() -> Self {
+        Self {
+            snapshot: 
RwLock::new(Arc::new(CatalogMetadataSnapshot::default())),
+            next_generation: AtomicU64::new(0),
+            published_generation: AtomicU64::new(0),
+        }
+    }
+}
+
+impl Deref for CatalogMetadataState {
+    type Target = RwLock<Arc<CatalogMetadataSnapshot>>;
+
+    fn deref(&self) -> &Self::Target {
+        &self.snapshot
+    }
+}
+
+impl CatalogMetadataState {
+    fn begin_refresh(&self) -> u64 {
+        self.next_generation.fetch_add(1, Ordering::AcqRel) + 1
+    }
+
+    fn publish(&self, generation: u64, next: Arc<CatalogMetadataSnapshot>) {
+        let mut current = self.snapshot.write().unwrap_or_else(|e| 
e.into_inner());
+        if generation >= self.published_generation.load(Ordering::Acquire) {

Review Comment:
   Could we make generation tracking per database, or merge non-conflicting 
results? A database-scoped refresh with generation N+1 causes an 
already-running full refresh at N to discard its entire snapshot, including 
newly loaded metadata for unrelated databases, while still returning `Ok(())`. 
That makes a successful full refresh silently ineffective for unrelated entries.
   



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