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]