fix: invalidate comment DDL cache and lock by object ID (#8390)

* fix: invalidate comment ddl cache locally

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix: fix typos

Signed-off-by: WenyXu <wenymedia@gmail.com>

* chore: apply suggestions

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-07-01 10:02:57 +00:00
committed by GitHub
parent 8a11ee9462
commit 3aa82567c4
4 changed files with 206 additions and 56 deletions
+40 -7
View File
@@ -14,6 +14,7 @@
use api::v1::CommentOnExpr;
use common_error::ext::BoxedError;
use common_meta::cache_invalidator::Context;
use common_meta::procedure_executor::ExecutorContext;
use common_meta::rpc::ddl::{CommentObjectType, CommentOnTask, DdlTask, SubmitDdlTaskRequest};
use common_query::Output;
@@ -23,7 +24,9 @@ use snafu::ResultExt;
use sql::ast::ObjectNamePartExt;
use sql::statements::comment::{Comment, CommentObject};
use crate::error::{ExecuteDdlSnafu, ExternalSnafu, InvalidSqlSnafu, Result};
use crate::error::{
self, ExecuteDdlSnafu, ExternalSnafu, InvalidSqlSnafu, Result, TableMetadataManagerSnafu,
};
use crate::statement::StatementExecutor;
use crate::utils::to_meta_query_context;
@@ -39,7 +42,15 @@ impl StatementExecutor {
///
/// A `Result` containing the `Output` of the operation, or an error if the operation fails.
pub async fn comment(&self, stmt: Comment, query_ctx: QueryContextRef) -> Result<Output> {
let comment_on_task = self.create_comment_on_task_from_stmt(stmt, &query_ctx)?;
let mut comment_on_task = self.create_comment_on_task_from_stmt(stmt, &query_ctx)?;
comment_on_task
.enrich_object_id(
self.table_metadata_manager.table_name_manager(),
self.flow_metadata_manager.flow_name_manager(),
)
.await
.context(TableMetadataManagerSnafu)?;
let cache_idents = comment_on_task.cache_idents();
let request = SubmitDdlTaskRequest::new(
to_meta_query_context(query_ctx),
@@ -49,8 +60,15 @@ impl StatementExecutor {
self.procedure_executor
.submit_ddl_task(&ExecutorContext::default(), request)
.await
.context(ExecuteDdlSnafu)
.map(|_| Output::new_with_affected_rows(0))
.context(ExecuteDdlSnafu)?;
// Invalidates local cache ASAP.
self.cache_invalidator
.invalidate(&Context::default(), &cache_idents)
.await
.context(error::InvalidateTableCacheSnafu)?;
Ok(Output::new_with_affected_rows(0))
}
pub async fn comment_by_expr(
@@ -58,7 +76,15 @@ impl StatementExecutor {
expr: CommentOnExpr,
query_ctx: QueryContextRef,
) -> Result<Output> {
let comment_on_task = self.create_comment_on_task_from_expr(expr)?;
let mut comment_on_task = self.create_comment_on_task_from_expr(expr)?;
comment_on_task
.enrich_object_id(
self.table_metadata_manager.table_name_manager(),
self.flow_metadata_manager.flow_name_manager(),
)
.await
.context(TableMetadataManagerSnafu)?;
let cache_idents = comment_on_task.cache_idents();
let request = SubmitDdlTaskRequest::new(
to_meta_query_context(query_ctx),
@@ -68,8 +94,15 @@ impl StatementExecutor {
self.procedure_executor
.submit_ddl_task(&ExecutorContext::default(), request)
.await
.context(ExecuteDdlSnafu)
.map(|_| Output::new_with_affected_rows(0))
.context(ExecuteDdlSnafu)?;
// Invalidates local cache ASAP.
self.cache_invalidator
.invalidate(&Context::default(), &cache_idents)
.await
.context(error::InvalidateTableCacheSnafu)?;
Ok(Output::new_with_affected_rows(0))
}
fn create_comment_on_task_from_expr(&self, expr: CommentOnExpr) -> Result<CommentOnTask> {