fix: postgres describe for more statements (#8974)

* fix: postgres describe for more statements

Signed-off-by: Ning Sun <sunning@greptime.com>

* fix: cover more show statements

Signed-off-by: Ning Sun <sunning@greptime.com>

* fix: address review comments

- add missing `clippy::too_many_arguments` allow on
  `query_from_information_schema_dataframe` (CI clippy failure)
- take `&ShowKind` in the information-schema dataframe helper so `kind`
  is no longer cloned at every call site; only the WHERE arm (which needs
  an owned expression for `sql_to_expr`) clones internally
- document why re-applying TQL explain formats never overwrites an
  existing value (per-query context state)

Signed-off-by: Ning Sun <sunning@greptime.com>

* chore: trim comments to essentials

Signed-off-by: Ning Sun <sunning@greptime.com>

---------

Signed-off-by: Ning Sun <sunning@greptime.com>
(cherry picked from commit d32cd77505)
This commit is contained in:
Ning Sun
2026-09-04 13:14:41 +08:00
committed by discord9
parent 03eb28f5b6
commit 2527bc09e7
10 changed files with 825 additions and 142 deletions
-1
View File
@@ -47,7 +47,6 @@ impl RecordBatchStreamCursor {
}
}
/// Returns the schema of the underlying record batch stream.
pub fn schema(&self) -> SchemaRef {
self.schema.clone()
}
+101
View File
@@ -59,6 +59,7 @@ use common_recordbatch::error::StreamTimeoutSnafu;
use common_telemetry::logging::SlowQueryOptions;
use common_telemetry::{debug, error, tracing};
use dashmap::DashMap;
use datafusion::dataframe::DataFrame;
use datafusion::physical_plan::ExecutionPlan;
use datafusion_expr::LogicalPlan;
use futures::{Stream, StreamExt, future};
@@ -847,6 +848,22 @@ impl Instance {
query_interceptor.pre_execute(stmt.as_ref(), Some(&plan), query_ctx.clone())?;
// TQL EXPLAIN/ANALYZE formats are consumed from the query context at
// execution time (see `optimize_physical_plan`); re-apply the side
// effect of `plan_tql` that was lost when the plan was built during
// Describe. `explain_format` is per-query state, so this never
// overwrites anything.
if let Some(Statement::Tql(tql)) = &stmt {
let format = match tql {
Tql::Explain(explain) => explain.format.as_ref(),
Tql::Analyze(analyze) => analyze.format.as_ref(),
Tql::Eval(_) => None,
};
if let Some(format) = format {
query_ctx.set_explain_format(format.to_string());
}
}
let query = stmt
.as_ref()
.map(|s| s.to_string())
@@ -928,6 +945,59 @@ impl Instance {
vec![result]
}
/// Builds the [`DataFrame`] for an information-schema-backed `SHOW`
/// statement; `None` for other statements. The future is boxed to keep
/// `do_describe_inner`'s state machine small.
fn show_statement_dataframe<'a>(
&'a self,
stmt: &'a Statement,
query_ctx: &'a QueryContextRef,
) -> Pin<Box<dyn Future<Output = Option<query::error::Result<DataFrame>>> + Send + 'a>> {
Box::pin(async move {
let engine = &self.query_engine;
let catalog_manager = self.catalog_manager();
let ctx = query_ctx.clone();
let dataframe = match stmt {
Statement::ShowDatabases(show) => {
query::sql::show_databases_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowTables(show) => {
query::sql::show_tables_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowViews(show) => {
query::sql::show_views_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowFlows(show) => {
query::sql::show_flows_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowColumns(show) => {
query::sql::show_columns_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowTableStatus(show) => {
query::sql::show_table_status_dataframe(show, engine, catalog_manager, ctx)
.await
}
Statement::ShowCharset(kind) => {
query::sql::show_charsets_dataframe(kind, engine, catalog_manager, ctx).await
}
Statement::ShowCollation(kind) => {
query::sql::show_collations_dataframe(kind, engine, catalog_manager, ctx).await
}
Statement::ShowIndex(show) => {
query::sql::show_index_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowRegion(show) => {
query::sql::show_region_dataframe(show, engine, catalog_manager, ctx).await
}
Statement::ShowProcesslist(show) => {
query::sql::show_processlist_dataframe(show, engine, catalog_manager, ctx).await
}
_ => return None,
};
Some(dataframe)
})
}
async fn do_describe_inner(
&self,
stmt: Statement,
@@ -947,6 +1017,37 @@ impl Instance {
let plannable = is_inner_plannable(&stmt)
|| matches!(&stmt, Statement::Explain(explain) if is_inner_plannable(explain.statement.as_ref()));
if let Statement::Tql(tql) = stmt {
// TQL produces a logical plan; describe it from the plan so the
// extended-protocol RowDescription matches the executed DataRows.
self.check_sql_permission(&Statement::Tql(tql.clone()), &query_ctx)
.await?;
let plan = self.statement_executor.plan_tql(tql, &query_ctx).await?;
return self
.query_engine
.describe(plan, query_ctx)
.await
.map(Some)
.context(error::DescribeStatementSnafu);
}
// Describe SHOW statements from the same projection the executor builds.
if let Some(dataframe) = self
.show_statement_dataframe(&stmt, &query_ctx)
.await
.transpose()
.context(PlanStatementSnafu)?
{
self.check_sql_permission(&stmt, &query_ctx).await?;
let plan = dataframe.into_unoptimized_plan();
return self
.query_engine
.describe(plan, query_ctx)
.await
.map(Some)
.context(error::DescribeStatementSnafu);
}
if plannable {
self.check_sql_permission(&stmt, &query_ctx).await?;
+2
View File
@@ -12,6 +12,8 @@
// See the License for the specific language governing permissions and
// limitations under the License.
#![recursion_limit = "256"]
pub mod error;
pub mod events;
pub mod frontend;
+1 -1
View File
@@ -83,7 +83,7 @@ use table::table_reference::TableReference;
pub use self::admin::{
AdminEventRecorderHandle, AdminFunctionLayer, AdminFunctionLayerRef,
AdminFunctionRecordingLayer, AdminFunctionRequest, AdminFunctionResponse, AdminFunctionService,
AdminFunctionServiceRef,
AdminFunctionServiceRef, admin_output_schema,
};
use self::set::{
set_bytea_output, set_datestyle, set_intervalstyle, set_timezone, validate_client_encoding,
+114 -57
View File
@@ -19,6 +19,7 @@ use std::sync::Arc;
use common_function::function::FunctionContext;
use common_function::function_registry::{FUNCTION_REGISTRY, get_admin_function};
use common_function::state::FunctionState;
use common_query::Output;
use common_recordbatch::{RecordBatch, RecordBatches};
use common_sql::convert::sql_value_to_value;
@@ -77,6 +78,103 @@ struct CoreAdminFunctionService {
query_engine: query::QueryEngineRef,
}
/// Parts of an `ADMIN` call needed both for execution and schema derivation.
struct ResolvedAdminFunction {
admin_udf: datafusion_expr::ScalarUDF,
fn_name: String,
args: Vec<VectorRef>,
arg_types: Vec<ArrowDataType>,
ret_type: ArrowDataType,
}
/// Resolves the function, parses its literal arguments and derives its
/// return type, without executing it.
fn resolve_admin_function(
stmt: &Admin,
query_ctx: &QueryContextRef,
state: Arc<FunctionState>,
) -> Result<ResolvedAdminFunction> {
let Admin::Func(func) = stmt;
// the function name should be in lower case.
let func_name = func.name.to_string().to_lowercase();
let factory = get_admin_function(&func_name)
.or_else(|| FUNCTION_REGISTRY.get_function(&func_name))
.context(error::AdminFunctionNotFoundSnafu {
name: func_name.clone(),
})?;
let func_ctx = FunctionContext {
query_ctx: query_ctx.clone(),
state,
};
let admin_udf = factory.provide(func_ctx);
admin_udf
.as_async()
.context(error::AdminFunctionNotFoundSnafu { name: func_name })?;
let fn_name = admin_udf.name().to_string();
let signature = admin_udf.signature();
// Parse function arguments
let FunctionArguments::List(args) = &func.args else {
return error::BuildAdminFunctionArgsSnafu {
msg: format!("unsupported function args {} for {}", func.args, fn_name),
}
.fail();
};
let arg_values = args
.args
.iter()
.map(|arg| {
let FunctionArg::Unnamed(FunctionArgExpr::Expr(Expr::Value(value))) = arg else {
return error::BuildAdminFunctionArgsSnafu {
msg: format!("unsupported function arg {arg} for {}", fn_name),
}
.fail();
};
Ok(&value.value)
})
.collect::<Result<Vec<_>>>()?;
let args = args_to_vector(&signature.type_signature, &arg_values, query_ctx)?;
let arg_types = args
.iter()
.map(|arg| arg.data_type().as_arrow_type())
.collect::<Vec<_>>();
let ret_type =
admin_udf
.return_type(&arg_types)
.map_err(|e| error::Error::BuildAdminFunctionArgs {
msg: format!(
"Failed to get return type of admin function {}: {}",
fn_name, e
),
})?;
Ok(ResolvedAdminFunction {
admin_udf,
fn_name,
args,
arg_types,
ret_type,
})
}
/// Output schema of an `ADMIN` statement, mirroring what
/// [`CoreAdminFunctionService::execute`] produces. `None` if the function
/// or arguments are unresolvable; execution will then surface the error.
pub fn admin_output_schema(stmt: &Admin, query_ctx: &QueryContextRef) -> Option<Schema> {
let resolved =
resolve_admin_function(stmt, query_ctx, Arc::new(FunctionState::default())).ok()?;
Some(Schema::new(vec![ColumnSchema::new(
// Use statement as the result column name
stmt.to_string(),
ConcreteDataType::from_arrow_type(&resolved.ret_type),
false,
)]))
}
impl CoreAdminFunctionService {
fn new(query_engine: query::QueryEngineRef) -> Self {
Self { query_engine }
@@ -88,62 +186,23 @@ impl CoreAdminFunctionService {
query_ctx,
} = request;
let Admin::Func(func) = &stmt;
// the function name should be in lower case.
let func_name = func.name.to_string().to_lowercase();
let factory = get_admin_function(&func_name)
.or_else(|| FUNCTION_REGISTRY.get_function(&func_name))
.context(error::AdminFunctionNotFoundSnafu {
name: func_name.clone(),
})?;
let func_ctx = FunctionContext {
query_ctx: query_ctx.clone(),
state: self.query_engine.engine_state().function_state(),
};
let admin_udf = factory.provide(func_ctx);
let resolved = resolve_admin_function(
&stmt,
&query_ctx,
self.query_engine.engine_state().function_state(),
)?;
let ResolvedAdminFunction {
admin_udf,
fn_name,
args,
arg_types,
ret_type,
} = resolved;
let admin_async_fn = admin_udf
.as_async()
.context(error::AdminFunctionNotFoundSnafu { name: func_name })?;
let fn_name = admin_udf.name();
let signature = admin_udf.signature();
// Parse function arguments
let FunctionArguments::List(args) = &func.args else {
return error::BuildAdminFunctionArgsSnafu {
msg: format!("unsupported function args {} for {}", func.args, fn_name),
}
.fail();
};
let arg_values = args
.args
.iter()
.map(|arg| {
let FunctionArg::Unnamed(FunctionArgExpr::Expr(Expr::Value(value))) = arg else {
return error::BuildAdminFunctionArgsSnafu {
msg: format!("unsupported function arg {arg} for {}", fn_name),
}
.fail();
};
Ok(&value.value)
})
.collect::<Result<Vec<_>>>()?;
let args = args_to_vector(&signature.type_signature, &arg_values, &query_ctx)?;
let arg_types = args
.iter()
.map(|arg| arg.data_type().as_arrow_type())
.collect::<Vec<_>>();
let ret_type = admin_udf.return_type(&arg_types).map_err(|e| {
error::Error::BuildAdminFunctionArgs {
msg: format!(
"Failed to get return type of admin function {}: {}",
fn_name, e
),
}
})?;
.context(error::AdminFunctionNotFoundSnafu {
name: fn_name.clone(),
})?;
// Convert arguments to DataFusion ColumnarValue format
let columnar_args: Vec<datafusion_expr::ColumnarValue> = args
@@ -174,9 +233,7 @@ impl CoreAdminFunctionService {
let result_columnar = admin_async_fn
.invoke_async_with_args(func_args)
.await
.with_context(|_| ExecuteAdminFunctionSnafu {
msg: fn_name.to_string(),
})?;
.with_context(|_| ExecuteAdminFunctionSnafu { msg: fn_name })?;
// Convert result back to VectorRef
let result_columnar: common_query::prelude::ColumnarValue =
+3 -7
View File
@@ -46,13 +46,9 @@ const STAGE: &str = "stage";
const NODE: &str = "node";
const PLAN: &str = "plan";
/// Returns the fixed output schema of [`DistAnalyzeExec`]:
/// (`stage`: UInt32, `node`: UInt32, `plan`: Utf8).
///
/// Exposed so protocol-level `Describe` handlers can advertise the schema
/// clients will actually receive (the physical plan is rewritten in
/// `optimize_physical_plan` and its schema differs from the DataFusion
/// `Analyze` logical plan schema).
/// Fixed output schema of [`DistAnalyzeExec`], for `Describe` handlers:
/// execution rewrites the plan in `optimize_physical_plan`, so this schema
/// differs from the logical `Analyze` plan's.
pub fn dist_analyze_output_schema() -> SchemaRef {
SchemaRef::new(Schema::new(vec![
Field::new(STAGE, DataType::UInt32, true),
+206 -43
View File
@@ -41,6 +41,7 @@ use common_recordbatch::RecordBatches;
use common_recordbatch::adapter::RecordBatchStreamAdapter;
use common_time::Timestamp;
use common_time::timezone::get_timezone;
use datafusion::dataframe::DataFrame;
use datafusion::prelude::SessionContext;
use datafusion_expr::{Expr, SortExpr, col, lit};
use datatypes::prelude::*;
@@ -109,7 +110,7 @@ const INDEX_KEY_NAME_COLUMN: &str = "Key_name";
const INDEX_SEQ_IN_INDEX_COLUMN: &str = "Seq_in_index";
const INDEX_COLUMN_NAME_COLUMN: &str = "Column_name";
static DESCRIBE_TABLE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
pub static DESCRIBE_TABLE_OUTPUT_SCHEMA: Lazy<Arc<Schema>> = Lazy::new(|| {
Arc::new(Schema::new(vec![
ColumnSchema::new(
COLUMN_NAME_COLUMN,
@@ -178,6 +179,18 @@ pub async fn show_databases(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let dataframe =
show_databases_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW DATABASES` without executing it.
pub async fn show_databases_dataframe(
stmt: &ShowDatabases,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let projects = if stmt.full {
vec![
(schemata::SCHEMA_NAME, SCHEMAS_COLUMN),
@@ -191,7 +204,7 @@ pub async fn show_databases(
let like_field = Some(schemata::SCHEMA_NAME);
let sort = vec![col(schemata::SCHEMA_NAME).sort(true, true)];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -201,7 +214,7 @@ pub async fn show_databases(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -249,6 +262,45 @@ async fn query_from_information_schema_table(
sort: Vec<SortExpr>,
kind: ShowKind,
) -> Result<Output> {
let dataframe = query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
table_name,
select,
projects,
filters,
like_field,
sort,
&kind,
)
.await?;
dataframe_to_output(dataframe).await
}
async fn dataframe_to_output(dataframe: DataFrame) -> Result<Output> {
let stream = dataframe.execute_stream().await?;
Ok(Output::new_with_stream(Box::pin(
RecordBatchStreamAdapter::try_new(stream).context(error::CreateRecordBatchSnafu)?,
)))
}
/// Builds the [`DataFrame`] for a `SHOW` statement without executing it,
/// so `Describe` handlers can derive the output schema from the same
/// projection the executor uses.
#[allow(clippy::too_many_arguments)]
async fn query_from_information_schema_dataframe(
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
table_name: &str,
select: Vec<Expr>,
projects: Vec<(&str, &str)>,
filters: Vec<Expr>,
like_field: Option<&str>,
sort: Vec<SortExpr>,
kind: &ShowKind,
) -> Result<DataFrame> {
let table = catalog_manager
.table(
query_ctx.current_catalog(),
@@ -336,7 +388,7 @@ async fn query_from_information_schema_table(
.expect("Must be the datafusion planner");
let filter = planner
.sql_to_expr(filter, dataframe.schema(), false, query_ctx)
.sql_to_expr(filter.clone(), dataframe.schema(), false, query_ctx)
.await?;
// Apply the `where` clause filters
@@ -344,11 +396,7 @@ async fn query_from_information_schema_table(
}
};
let stream = dataframe.execute_stream().await?;
Ok(Output::new_with_stream(Box::pin(
RecordBatchStreamAdapter::try_new(stream).context(error::CreateRecordBatchSnafu)?,
)))
Ok(dataframe)
}
/// Execute `SHOW COLUMNS` statement.
@@ -358,8 +406,19 @@ pub async fn show_columns(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let schema_name = if let Some(database) = stmt.database {
database
let dataframe = show_columns_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW COLUMNS` without executing it.
pub async fn show_columns_dataframe(
stmt: &ShowColumns,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let schema_name = if let Some(database) = &stmt.database {
database.clone()
} else {
query_ctx.current_schema()
};
@@ -397,7 +456,7 @@ pub async fn show_columns(
let like_field = Some(columns::COLUMN_NAME);
let sort = vec![col(columns::COLUMN_NAME).sort(true, true)];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -407,7 +466,7 @@ pub async fn show_columns(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -419,8 +478,19 @@ pub async fn show_index(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let schema_name = if let Some(database) = stmt.database {
database
let dataframe = show_index_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW INDEX` without executing it.
pub async fn show_index_dataframe(
stmt: &ShowIndex,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let schema_name = if let Some(database) = &stmt.database {
database.clone()
} else {
query_ctx.current_schema()
};
@@ -471,7 +541,7 @@ pub async fn show_index(
col(statistics::SEQ_IN_INDEX).sort(true, true),
];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -481,7 +551,7 @@ pub async fn show_index(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -493,8 +563,19 @@ pub async fn show_region(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let schema_name = if let Some(database) = stmt.database {
database
let dataframe = show_region_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW REGION` without executing it.
pub async fn show_region_dataframe(
stmt: &ShowRegion,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let schema_name = if let Some(database) = &stmt.database {
database.clone()
} else {
query_ctx.current_schema()
};
@@ -517,7 +598,7 @@ pub async fn show_region(
col(columns::PEER_ID).sort(true, true),
];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -527,7 +608,7 @@ pub async fn show_region(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -539,8 +620,19 @@ pub async fn show_tables(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let schema_name = if let Some(database) = stmt.database {
database
let dataframe = show_tables_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for [`ShowTables`] without executing it.
pub async fn show_tables_dataframe(
stmt: &ShowTables,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let schema_name = if let Some(database) = &stmt.database {
database.clone()
} else {
query_ctx.current_schema()
};
@@ -564,15 +656,16 @@ pub async fn show_tables(
// Transform the WHERE clause for backward compatibility:
// Replace "Tables" with "Tables_in_{schema}" to support old queries
let kind = match stmt.kind {
ShowKind::Where(mut filter) => {
let rewritten_kind = match &stmt.kind {
ShowKind::Where(filter) => {
let mut filter = filter.clone();
replace_column_in_expr(&mut filter, "Tables", &tables_column);
ShowKind::Where(filter)
}
other => other,
kind => kind.clone(),
};
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -582,7 +675,7 @@ pub async fn show_tables(
filters,
like_field,
sort,
kind,
&rewritten_kind,
)
.await
}
@@ -594,8 +687,20 @@ pub async fn show_table_status(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let schema_name = if let Some(database) = stmt.database {
database
let dataframe =
show_table_status_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for [`ShowTableStatus`] without executing it.
pub async fn show_table_status_dataframe(
stmt: &ShowTableStatus,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let schema_name = if let Some(database) = &stmt.database {
database.clone()
} else {
query_ctx.current_schema()
};
@@ -629,7 +734,7 @@ pub async fn show_table_status(
let like_field = Some(tables::TABLE_NAME);
let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -639,7 +744,7 @@ pub async fn show_table_status(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -651,6 +756,18 @@ pub async fn show_collations(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let dataframe =
show_collations_dataframe(&kind, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW COLLATION` without executing it.
pub async fn show_collations_dataframe(
kind: &ShowKind,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
// Refer to https://dev.mysql.com/doc/refman/8.0/en/show-collation.html
let projects = vec![
("collation_name", "Collation"),
@@ -665,7 +782,7 @@ pub async fn show_collations(
let like_field = Some("collation_name");
let sort = vec![];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -687,6 +804,18 @@ pub async fn show_charsets(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let dataframe =
show_charsets_dataframe(&kind, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW CHARSET` without executing it.
pub async fn show_charsets_dataframe(
kind: &ShowKind,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
// Refer to https://dev.mysql.com/doc/refman/8.0/en/show-character-set.html
let projects = vec![
("character_set_name", "Charset"),
@@ -699,7 +828,7 @@ pub async fn show_charsets(
let like_field = Some("character_set_name");
let sort = vec![];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -916,8 +1045,19 @@ pub async fn show_views(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let schema_name = if let Some(database) = stmt.database {
database
let dataframe = show_views_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for [`ShowViews`] without executing it.
pub async fn show_views_dataframe(
stmt: &ShowViews,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let schema_name = if let Some(database) = &stmt.database {
database.clone()
} else {
query_ctx.current_schema()
};
@@ -930,7 +1070,7 @@ pub async fn show_views(
let like_field = Some(tables::TABLE_NAME);
let sort = vec![col(tables::TABLE_NAME).sort(true, true)];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -940,7 +1080,7 @@ pub async fn show_views(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -952,12 +1092,23 @@ pub async fn show_flows(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let dataframe = show_flows_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for [`ShowFlows`] without executing it.
pub async fn show_flows_dataframe(
stmt: &ShowFlows,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let projects = vec![(flows::FLOW_NAME, FLOWS_COLUMN)];
let filters = vec![col(flows::TABLE_CATALOG).eq(lit(query_ctx.current_catalog()))];
let like_field = Some(flows::FLOW_NAME);
let sort = vec![col(flows::FLOW_NAME).sort(true, true)];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx,
@@ -967,7 +1118,7 @@ pub async fn show_flows(
filters,
like_field,
sort,
stmt.kind,
&stmt.kind,
)
.await
}
@@ -1340,6 +1491,18 @@ pub async fn show_processlist(
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<Output> {
let dataframe =
show_processlist_dataframe(&stmt, query_engine, catalog_manager, query_ctx).await?;
dataframe_to_output(dataframe).await
}
/// Builds the [`DataFrame`] for `SHOW PROCESSLIST` without executing it.
pub async fn show_processlist_dataframe(
stmt: &ShowProcessList,
query_engine: &QueryEngineRef,
catalog_manager: &CatalogManagerRef,
query_ctx: QueryContextRef,
) -> Result<DataFrame> {
let projects = if stmt.full {
vec![
(process_list::ID, "Id"),
@@ -1367,17 +1530,17 @@ pub async fn show_processlist(
};
let like_field = None;
let sort = vec![col("id").sort(true, true)];
query_from_information_schema_table(
query_from_information_schema_dataframe(
query_engine,
catalog_manager,
query_ctx.clone(),
"process_list",
vec![],
projects.clone(),
projects,
filters,
like_field,
sort,
ShowKind::All,
&ShowKind::All,
)
.await
}
+54 -29
View File
@@ -28,6 +28,7 @@ use datafusion_pg_catalog::sql::PostgresCompatibilityParser;
use datatypes::prelude::ConcreteDataType;
use datatypes::schema::{Schema, SchemaRef};
use futures::{Sink, SinkExt, Stream, StreamExt, future, stream};
use operator::statement::admin_output_schema;
use pgwire::api::portal::{Format, Portal};
use pgwire::api::query::{ExtendedQueryHandler, SimpleQueryHandler};
use pgwire::api::results::{
@@ -43,6 +44,7 @@ use pgwire::messages::data::DataRow;
use query::dist_analyze_output_schema;
use query::planner::DfLogicalPlanner;
use query::query_engine::DescribeResult;
use query::sql::DESCRIBE_TABLE_OUTPUT_SCHEMA;
use session::Session;
use session::context::QueryContextRef;
use snafu::ResultExt;
@@ -555,11 +557,8 @@ fn describe_fields(
session: &Arc<Session>,
) -> PgWireResult<Vec<FieldInfo>> {
match sql_plan {
// EXPLAIN ANALYZE: at execution time the physical plan is replaced with
// GreptimeDB's `DistAnalyzeExec` (see `optimize_physical_plan`), whose
// output schema (stage/node/plan) differs from the DataFusion `Analyze`
// logical plan schema (plan_type/plan). Describe with the schema the
// client will actually receive so the DataRow field count matches.
// Execution swaps in DistAnalyzeExec (stage/node/plan), whose schema
// differs from the logical `Analyze` plan's (plan_type/plan).
SqlPlan::Plan(LogicalPlan::Analyze(_), _) => {
let schema: Schema =
Schema::try_from(dist_analyze_output_schema()).map_err(convert_err)?;
@@ -658,17 +657,6 @@ fn describe_fields(
),
]),
// single column show statements
SqlPlan::Statement(
Statement::ShowTables(_) | Statement::ShowFlows(_) | Statement::ShowViews(_),
_,
) => Ok(vec![FieldInfo::new(
"name".to_string(),
None,
None,
Type::TEXT,
format.format_for(0),
)]),
#[cfg(feature = "enterprise")]
SqlPlan::Statement(Statement::ShowTriggers(_), _) => Ok(vec![FieldInfo::new(
"name".to_string(),
@@ -690,22 +678,59 @@ fn describe_fields(
Ok(vec![])
}
}
// FETCH cursor: return the cursor's schema so the RowDescription
// matches the DataRow field count sent during Execute.
// Single column named after the variable (see `query::sql::show_variable`).
SqlPlan::Statement(Statement::ShowVariables(show), _) => Ok(vec![FieldInfo::new(
show.variable.to_string().to_uppercase(),
None,
None,
Type::TEXT,
format.format_for(0),
)]),
// Mirrors `query::sql::show_status` (currently always empty).
SqlPlan::Statement(Statement::ShowStatus(_), _) => Ok(vec![
FieldInfo::new(
"Variable_name".to_string(),
None,
None,
Type::TEXT,
format.format_for(0),
),
FieldInfo::new(
"Value".to_string(),
None,
None,
Type::TEXT,
format.format_for(1),
),
]),
SqlPlan::Statement(Statement::ShowSearchPath(_), _) => Ok(vec![FieldInfo::new(
"search_path".to_string(),
None,
None,
Type::TEXT,
format.format_for(0),
)]),
// Mirrors `query::sql::describe_table`.
SqlPlan::Statement(Statement::DescribeTable(_), _) => {
schema_to_pg(&DESCRIBE_TABLE_OUTPUT_SCHEMA, format, None).map_err(convert_err)
}
// Single column typed with the function's return type (see
// `operator::statement::admin_output_schema`).
SqlPlan::Statement(Statement::Admin(admin), _) => {
let query_ctx = session.new_query_context();
match admin_output_schema(admin, &query_ctx) {
Some(schema) => schema_to_pg(&schema, format, None).map_err(convert_err),
// Unresolvable; execution will surface the error.
None => Ok(vec![]),
}
}
// Describe from the declared cursor's schema.
SqlPlan::Statement(Statement::FetchCursor(fetch), _) => {
let cursor_name = fetch.cursor_name.to_string();
match session.get_cursor(&cursor_name) {
Some(cursor) => {
// `cursor.schema()` is the GreptimeDB `SchemaRef` captured
// when the cursor was declared; `schema_to_pg` accepts it
// directly (no DataFusion -> GreptimeDB conversion needed).
schema_to_pg(&cursor.schema(), format, None).map_err(convert_err)
}
None => {
// Cursor not found (e.g. DECLARE hasn't executed yet in
// an extended-protocol batch). Return NoData as a fallback.
Ok(vec![])
}
Some(cursor) => schema_to_pg(&cursor.schema(), format, None).map_err(convert_err),
// Cursor not declared yet; execution will error.
None => Ok(vec![]),
}
}
_ => {
+1 -4
View File
@@ -111,10 +111,7 @@ impl Session {
.into()
}
/// Returns the cursor with the given name, if it exists.
///
/// Cursors are stored in the session's mutable inner data and shared across
/// all query contexts created from this session.
/// Cursors are shared across query contexts created from this session.
pub fn get_cursor(&self, name: &str) -> Option<Arc<RecordBatchStreamCursor>> {
let guard = self.mutable_inner.read().unwrap();
guard.cursors.get(name).cloned()
+343
View File
@@ -91,6 +91,8 @@ macro_rules! sql_tests {
test_mysql_prepare_stmt_insert_timestamp,
test_mysql_prepare_stmt_timezone,
test_mysql_federated_prepare_stmt,
test_mysql_prepare_tql_and_show,
test_postgres_extended_query_row_returning_statements,
test_declare_fetch_close_cursor,
test_alter_update_on,
);
@@ -1487,6 +1489,280 @@ fn assert_pg_numeric_range_error(error: tokio_postgres::Error) {
assert_eq!("numeric_value_out_of_range", error.message());
}
pub async fn test_postgres_extended_query_row_returning_statements(store_type: StorageType) {
// Regression test for the tokio-postgres >= 0.7.14 DataRow/RowDescription
// mismatch: statements answered with NoData at Describe but emitting
// DataRows at Execute must describe their real output schema.
let (mut guard, fe_pg_server) = setup_pg_server(store_type, "test_pg_extended_row_stmts").await;
let addr = fe_pg_server.bind_addr().unwrap().to_string();
let (client, connection) = tokio_postgres::connect(&format!("postgres://{addr}/public"), NoTls)
.await
.unwrap();
let (tx, rx) = tokio::sync::oneshot::channel();
tokio::spawn(async move {
connection.await.unwrap();
tx.send(()).unwrap();
});
client
.execute(
"CREATE TABLE demo_metrics (ts timestamp time index, val double, host string primary key skipping index)",
&[],
)
.await
.unwrap();
client
.execute(
"INSERT INTO demo_metrics (ts, host, val) VALUES (1000, 'host-a', 1.0), (2000, 'host-b', 2.0)",
&[],
)
.await
.unwrap();
// ---- SHOW DATABASES: single `Database` column ----
let rows = client.query("SHOW DATABASES", &[]).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(1, rows[0].columns().len());
assert_eq!("Database", rows[0].columns()[0].name());
assert!(rows.iter().any(|r| r.get::<_, String>(0) == "public"));
// ---- SHOW FULL DATABASES: `Database` + `Options` columns ----
let rows = client.query("SHOW FULL DATABASES", &[]).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(2, rows[0].columns().len());
assert_eq!("Database", rows[0].columns()[0].name());
assert_eq!("Options", rows[0].columns()[1].name());
assert!(rows.iter().any(|r| r.get::<_, String>(0) == "public"));
// ---- SHOW TABLES: single `Tables_in_<schema>` column ----
let rows = client.query("SHOW TABLES", &[]).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(1, rows[0].columns().len());
assert_eq!("Tables_in_public", rows[0].columns()[0].name());
assert!(rows.iter().any(|r| r.get::<_, String>(0) == "demo_metrics"));
// ---- SHOW FULL TABLES: `Tables_in_<schema>` + `Table_type` columns ----
let rows = client.query("SHOW FULL TABLES", &[]).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(2, rows[0].columns().len());
assert_eq!("Tables_in_public", rows[0].columns()[0].name());
assert_eq!("Table_type", rows[0].columns()[1].name());
// ---- SHOW VIEWS / SHOW FLOWS: empty results, still described with one column ----
let rows = client.query("SHOW VIEWS", &[]).await.unwrap();
assert!(rows.is_empty());
let stmt = client.prepare("SHOW VIEWS").await.unwrap();
assert_eq!(1, stmt.columns().len());
assert_eq!("Views", stmt.columns()[0].name());
let rows = client.query("SHOW FLOWS", &[]).await.unwrap();
assert!(rows.is_empty());
let stmt = client.prepare("SHOW FLOWS").await.unwrap();
assert_eq!(1, stmt.columns().len());
assert_eq!("Flows", stmt.columns()[0].name());
// ---- SHOW TABLE STATUS: fixed eighteen-column schema ----
let rows = client.query("SHOW TABLE STATUS", &[]).await.unwrap();
assert!(!rows.is_empty());
let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect();
assert_eq!(
vec![
"Name",
"Engine",
"Version",
"Row_format",
"Rows",
"Avg_row_length",
"Data_length",
"Max_data_length",
"Index_length",
"Data_free",
"Auto_increment",
"Create_time",
"Update_time",
"Check_time",
"Collation",
"Checksum",
"Create_options",
"Comment",
],
names
);
// ---- SHOW COLUMNS / SHOW FULL COLUMNS ----
let rows = client
.query("SHOW COLUMNS FROM demo_metrics", &[])
.await
.unwrap();
assert_eq!(3, rows.len());
let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect();
assert_eq!(
vec![
"Field",
"Type",
"Null",
"Key",
"Default",
"Extra",
"Greptime_type"
],
names
);
let rows = client
.query("SHOW FULL COLUMNS FROM demo_metrics", &[])
.await
.unwrap();
assert_eq!(10, rows[0].columns().len());
// ---- SHOW CHARSET / SHOW COLLATION ----
let rows = client.query("SHOW CHARSET", &[]).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(4, rows[0].columns().len());
let rows = client.query("SHOW COLLATION", &[]).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(6, rows[0].columns().len());
// ---- SHOW INDEX: fixed fifteen-column schema ----
let rows = client
.query("SHOW INDEX IN demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect();
assert_eq!(
vec![
"Table",
"Non_unique",
"Key_name",
"Seq_in_index",
"Column_name",
"Collation",
"Cardinality",
"Sub_part",
"Packed",
"Null",
"Index_type",
"Comment",
"Index_comment",
"Visible",
"Expression",
],
names
);
// ---- SHOW REGION ----
let rows = client
.query("SHOW REGION IN demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
assert_eq!(4, rows[0].columns().len());
// ---- SHOW SEARCH_PATH: single string column ----
let rows = client.query("SHOW SEARCH_PATH", &[]).await.unwrap();
assert_eq!(1, rows.len());
assert_eq!(1, rows[0].columns().len());
assert_eq!("search_path", rows[0].columns()[0].name());
assert_eq!("public", rows[0].get::<_, String>(0));
// ---- SHOW VARIABLES: single column named after the variable ----
let rows = client.query("SHOW VARIABLES timezone", &[]).await.unwrap();
assert_eq!(1, rows.len());
assert_eq!(1, rows[0].columns().len());
assert_eq!("TIMEZONE", rows[0].columns()[0].name());
let _ = rows[0].get::<_, String>(0);
// ---- DESCRIBE TABLE: fixed six-column string schema ----
let rows = client
.query("DESCRIBE TABLE demo_metrics", &[])
.await
.unwrap();
assert_eq!(3, rows.len());
let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect();
assert_eq!(
vec!["Column", "Type", "Key", "Null", "Default", "Semantic Type"],
names
);
// first column of each row is the column name; ensure values decode as TEXT
let columns: Vec<String> = rows.iter().map(|r| r.get::<_, String>(0)).collect();
assert!(columns.contains(&"ts".to_string()));
assert!(columns.contains(&"val".to_string()));
assert!(columns.contains(&"host".to_string()));
// ---- ADMIN function: single column named after the statement ----
let rows = client
.query("ADMIN flush_table('demo_metrics')", &[])
.await
.unwrap();
assert_eq!(1, rows.len());
assert_eq!(1, rows[0].columns().len());
assert!(rows[0].columns()[0].name().contains("flush_table"));
// ---- TQL EVAL / EXPLAIN / ANALYZE: described from the planned query ----
let rows = client
.query("TQL EVAL (0, 3000, '1s') demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
let names: Vec<&str> = rows[0].columns().iter().map(|c| c.name()).collect();
assert_eq!(vec!["ts", "val", "host"], names);
let rows = client
.query("TQL EXPLAIN (0, 3000, '1s') demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
assert_eq!(2, rows[0].columns().len());
let rows = client
.query("TQL ANALYZE (0, 3000, '1s') demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
assert_eq!(3, rows[0].columns().len());
// FORMAT JSON variants keep working through the plan-based execution path
let rows = client
.query("TQL EXPLAIN FORMAT JSON (0, 3000, '1s') demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
let rows = client
.query("TQL ANALYZE FORMAT JSON (0, 3000, '1s') demo_metrics", &[])
.await
.unwrap();
assert!(!rows.is_empty());
// ---- the same statements also work over the simple query protocol ----
for sql in [
"SHOW DATABASES",
"SHOW FULL TABLES",
"SHOW TABLE STATUS",
"SHOW COLUMNS FROM demo_metrics",
"SHOW CHARSET",
"SHOW COLLATION",
"SHOW INDEX IN demo_metrics",
"SHOW REGION IN demo_metrics",
"SHOW SEARCH_PATH",
"DESCRIBE TABLE demo_metrics",
"ADMIN flush_table('demo_metrics')",
"TQL EVAL (0, 3000, '1s') demo_metrics",
] {
let msgs = client.simple_query(sql).await.unwrap();
assert!(
msgs.iter().any(|m| matches!(m, SimpleQueryMessage::Row(_))),
"simple query {sql} should return rows"
);
}
drop(client);
rx.await.unwrap();
let _ = fe_pg_server.shutdown().await;
guard.remove_all().await;
}
pub async fn test_postgres_explain_bind_parameter(store_type: StorageType) {
// Regression test for #8029: EXPLAIN / EXPLAIN ANALYZE must accept bind
// parameters over the Postgres extended query protocol.
@@ -1890,6 +2166,73 @@ pub async fn test_mysql_federated_prepare_stmt(store_type: StorageType) {
guard.remove_all().await;
}
pub async fn test_mysql_prepare_tql_and_show(store_type: StorageType) {
// `do_describe` now plans TQL and information-schema-backed SHOW
// statements, so MySQL prepared statements derive their column metadata
// from the planned query and execute the plan directly. This exercises
// that path end-to-end.
common_telemetry::init_default_ut_logging();
let (mut guard, fe_mysql_server) =
setup_mysql_server(store_type, "test_mysql_prepare_tql_and_show").await;
let addr = fe_mysql_server.bind_addr().unwrap().to_string();
let pool = MySqlPoolOptions::new()
.max_connections(2)
.connect(&format!("mysql://{addr}/public"))
.await
.unwrap();
sqlx::query(
"CREATE TABLE demo_metrics (ts timestamp time index, val double, host string primary key)",
)
.execute(&pool)
.await
.unwrap();
sqlx::query("INSERT INTO demo_metrics (ts, host, val) VALUES (1000, 'host-a', 1.0), (2000, 'host-b', 2.0)")
.execute(&pool)
.await
.unwrap();
// sqlx::query uses the binary prepared statement protocol
// (COM_STMT_PREPARE + COM_STMT_EXECUTE).
let rows = sqlx::query("TQL EVAL (0, 3000, '1s') demo_metrics")
.fetch_all(&pool)
.await
.unwrap();
assert!(!rows.is_empty());
assert_eq!(3, rows[0].columns().len());
let rows = sqlx::query("TQL ANALYZE (0, 3000, '1s') demo_metrics")
.fetch_all(&pool)
.await
.unwrap();
assert!(!rows.is_empty());
// SHOW statements share the same describe path; prepared SHOW FULL
// TABLES must report both columns.
let rows = sqlx::query("SHOW TABLES").fetch_all(&pool).await.unwrap();
assert!(!rows.is_empty());
assert_eq!(1, rows[0].columns().len());
let rows = sqlx::query("SHOW FULL TABLES")
.fetch_all(&pool)
.await
.unwrap();
assert!(!rows.is_empty());
assert_eq!(2, rows[0].columns().len());
let rows = sqlx::query("SHOW DATABASES")
.fetch_all(&pool)
.await
.unwrap();
assert!(!rows.is_empty());
assert_eq!(1, rows[0].columns().len());
let _ = fe_mysql_server.shutdown().await;
guard.remove_all().await;
}
pub async fn test_postgres_array_types(store_type: StorageType) {
let (mut guard, fe_pg_server) = setup_pg_server(store_type, "test_postgres_array_types").await;
let addr = fe_pg_server.bind_addr().unwrap().to_string();