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]