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


##########
crates/integrations/datafusion/src/catalog.rs:
##########
@@ -783,12 +1068,19 @@ impl SchemaProvider for PaimonSchemaProvider {
                             declared,
                         ))
                         .await?;
-                    Ok(resolved.map(|inner| {
-                        Arc::new(ReadOnlyTableProvider {
+                    Ok(Some(match resolved {

Review Comment:
   `TableEngineResolver` documents `Ok(None)` as not found, but this converts 
it into `Some(UnavailableEngineTableProvider)`. When an engine reports a 
deleted, absent, or unowned object, DataFusion now treats the table as 
existing, planning can succeed, and execution fails later. Please preserve 
`Ok(None)` for a registered resolver; the metadata-only fallback can remain for 
the no-resolver case.
   



##########
crates/integrations/datafusion/src/catalog.rs:
##########
@@ -924,101 +1208,25 @@ impl SchemaProvider for PaimonSchemaProvider {
                 return false;
             }
         };
-        if let Some(system_name) = object.system_table() {
-            if !system_tables::is_registered(system_name) {
-                return false;
-            }
+        if object
+            .system_table()
+            .is_some_and(|system_name| 
!system_tables::is_registered(system_name))
+        {
+            return false;
+        }
+        if object.system_table().is_some() {

Review Comment:
   Returning false for every recognized system table makes 
`table_exist("t$snapshots")` disagree with `table("t$snapshots")` and with 
successful Paimon system-table queries. This regresses the previous behavior 
and violates the `SchemaProvider::table_exist` existence contract for direct 
callers. Please derive the answer from the snapshotted base table, or retain 
enough metadata to distinguish supported and unsupported derived tables.
   



##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -473,7 +489,18 @@ impl SQLContext {
             ));
         }
 
-        match &statements[0] {
+        let refresh_metadata_after = 
self.metadata_change_targets(&statements[0], false)?;
+        let metadata_mutation_targets = 
self.metadata_change_targets(&statements[0], true)?;
+        {
+            let _refresh_guard = self.metadata_refresh_gate.lock().await;

Review Comment:
   Could the refresh targets be computed before taking a per-target 
singleflight lock? This context-wide gate is acquired by every SQL statement, 
even `SELECT 1`, and is held across remote catalog I/O. One stalled 
missing-table refresh therefore blocks unrelated existing-table queries and 
queries against other catalogs.
   



##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -751,7 +807,204 @@ impl SQLContext {
                 self.ctx.sql(&expanded.to_string()).await
             }
             _ => self.ctx.sql(sql).await,
+        };
+
+        if result.is_ok() && !refresh_metadata_after.is_empty() {
+            if let Err(error) = 
self.refresh_metadata_targets(refresh_metadata_after).await {

Review Comment:
   This refresh runs after the catalog mutation has committed, but the 
successful DDL result is still withheld while an unbounded metadata listing is 
awaited. If the listing stalls, the client is left uncertain and may retry an 
already-applied DDL. Please update the known snapshot directly, or detach or 
bound this best-effort refresh.
   



##########
crates/integrations/datafusion/src/sql_context.rs:
##########
@@ -751,7 +807,204 @@ impl SQLContext {
                 self.ctx.sql(&expanded.to_string()).await
             }
             _ => self.ctx.sql(sql).await,
+        };
+
+        if result.is_ok() && !refresh_metadata_after.is_empty() {
+            if let Err(error) = 
self.refresh_metadata_targets(refresh_metadata_after).await {
+                log::warn!("catalog metadata refresh after DDL failed: 
{error}");
+            }
+        }
+        result
+    }
+
+    fn metadata_refresh_targets(
+        &self,
+        statement: &Statement,
+    ) -> DFResult<HashSet<MetadataRefreshTarget>> {
+        let mut targets = HashSet::new();
+        if matches!(statement, Statement::ShowTables { .. }) {
+            let state = self.ctx.state();
+            targets.insert(MetadataRefreshTarget::Database {
+                catalog: self.current_catalog_name(),
+                database: 
state.config_options().catalog.default_schema.clone(),
+            });
+            return Ok(targets);
+        }
+        if let Statement::ShowColumns { show_options, .. } = statement {
+            if let Some(show_in) = &show_options.show_in {
+                if show_in.parent_type.is_none() {
+                    if let Some(name) = &show_in.parent_name {
+                        let (_, catalog, identifier) = 
self.resolve_catalog_and_table(name)?;
+                        targets.insert(MetadataRefreshTarget::Database {
+                            catalog,
+                            database: identifier.database().to_string(),
+                        });
+                    }
+                }
+            }
+            return Ok(targets);
+        }
+        if matches!(statement, Statement::ShowFunctions { .. }) {
+            return Ok(targets);
+        }
+
+        let statement = 
datafusion::sql::parser::Statement::Statement(Box::new(statement.clone()));
+        let state = self.ctx.state();
+        let default_catalog = 
state.config_options().catalog.default_catalog.clone();
+        let default_schema = 
state.config_options().catalog.default_schema.clone();
+        for reference in state.resolve_table_references(&statement)? {
+            let schema = reference.schema().unwrap_or(&default_schema);
+            let catalog_name = reference.catalog().unwrap_or(&default_catalog);
+            if !self.catalogs.contains_key(catalog_name) {
+                continue;
+            }
+            let provider = self.ctx.catalog(catalog_name).ok_or_else(|| {
+                DataFusionError::Plan(format!("Unknown catalog 
'{catalog_name}'"))
+            })?;
+            let provider = provider
+                .downcast_ref::<crate::catalog::PaimonCatalogProvider>()
+                .ok_or_else(|| {
+                    DataFusionError::Plan(format!(
+                        "Catalog '{catalog_name}' is not a Paimon catalog"
+                    ))
+                })?;
+            if schema.eq_ignore_ascii_case("information_schema") {
+                
targets.insert(MetadataRefreshTarget::Catalog(catalog_name.to_string()));

Review Comment:
   DataFusion information schema enumerates every catalog in the session, but 
this refreshes only the catalog that qualifies the `information_schema` 
relation. In a multi-catalog session, filtering for rows from another catalog 
can therefore return stale or missing metadata. Please refresh all 
participating Paimon catalogs, or safely derive a literal `table_catalog` 
restriction and refresh that target.
   



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