shyjsarah commented on code in PR #881:
URL: https://github.com/apache/paimon-rust/pull/881#discussion_r4056420887
##########
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:
Agreed. I narrowed the PR scope instead of claiming the remaining bridges
are fixed: the description now uses Refs #872, explicitly documents
register_schema, deregister_schema, and deregister_table as follow-up work, and
leaves #872 open.
##########
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:
Fixed in 4f68a94. Refresh targets are computed before locking, locks are
keyed by catalog/database and weakly retained, and independent targets refresh
concurrently. Metadata-free queries and other databases no longer wait behind a
stalled refresh; same-target misses still share the TTL/single-flight behavior.
##########
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:
Fixed in 4f68a94. A registered resolver returning Ok(None) is now preserved
as None. The UnavailableEngineTableProvider fallback is only used when no
resolver is registered. Added a regression test for the registered-resolver
not-found 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:
Fixed in 4f68a94. table_exist now recognizes supported system-table suffixes
and derives existence from the snapshotted base table; missing bases and
unknown suffixes remain false. Added direct table_exist coverage.
##########
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:
Fixed in 4f68a94. Generation tracking is now per database, and a full
refresh merges each database independently. Newer scoped updates and tombstones
win only for their own database, while non-conflicting results from the older
full refresh are still published. Added concurrent merge and stale-resurrection
regression tests.
##########
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:
Fixed in 4f68a94. Successful DDL now applies the known create/drop/rename
delta directly to the in-memory snapshot and does not await a post-commit
listing. Added a blocked-listing regression test to verify the DDL result
returns immediately.
##########
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:
Fixed in 4f68a94. An information_schema reference now refreshes every
registered Paimon catalog concurrently. Those catalog-wide refreshes are
best-effort so one unavailable catalog preserves its last good snapshot without
hiding strict failures for ordinary missing-object targets. Added
multi-catalog, failure, and concurrency tests.
--
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]