Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 2 additions & 3 deletions src/query/ee/src/materialized_view/refresh.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,6 @@ use databend_common_meta_app::schema::UpsertTableOptionReq;
use databend_common_sql::MaterializedViewChecker;
use databend_common_sql::Planner;
use databend_common_sql::parse_materialized_view_query;
use databend_common_sql::plans::Mutation;
use databend_common_sql::plans::Plan;
use databend_common_sql::validate_materialized_view_source;
use databend_common_storages_fuse::FuseTable;
Expand Down Expand Up @@ -366,10 +365,10 @@ impl<'a> MaterializedViewRefresh<'a> {
self.mv_table.get_id(),
)?,
Plan::DataMutation { s_expr, schema, .. } => {
let mutation: Mutation = s_expr.plan().clone().try_into()?;
let mutation = s_expr.mutation()?.clone();
Arc::new(MutationInterpreter::try_create_materialized_view_refresh(
self.ctx.clone(),
*s_expr,
(*s_expr).into_planned()?,
schema,
mutation.metadata,
self.mv_table.get_id(),
Expand Down
6 changes: 2 additions & 4 deletions src/query/service/src/interpreters/access/privilege_access.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,6 @@ use databend_common_sql::plans::AlterSharePlanAction;
use databend_common_sql::plans::InsertInputSource;
use databend_common_sql::plans::MaintenanceTarget;
use databend_common_sql::plans::ModifyColumnAction;
use databend_common_sql::plans::Mutation;
use databend_common_sql::plans::OptimizeCompactBlock;
use databend_common_sql::plans::PresignAction;
use databend_common_sql::plans::RewriteKind;
Expand Down Expand Up @@ -2245,10 +2244,9 @@ impl AccessChecker for PrivilegeAccess {
self.validate_insert_source(ctx, &plan.source).await?;
}
Plan::DataMutation { s_expr, .. } => {
let plan: Mutation = s_expr.plan().clone().try_into()?;
let plan = s_expr.mutation()?;
if enable_experimental_rbac_check {
let s_expr = s_expr.child(0)?;
match s_expr.get_udfs() {
match s_expr.input_udfs() {
Ok(udfs) => {
if !udfs.is_empty() {
self.validate_udf_access(udfs).await?;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ impl CopyIntoLocationInterpreter {
let select_interpreter = SelectInterpreter::try_create(
self.ctx.clone(),
*(bind_context.clone()),
*s_expr.clone(),
s_expr.planned()?.clone(),
metadata.clone(),
formatted_ast.clone(),
false,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -221,7 +221,7 @@ impl CopyIntoTableInterpreter {
let select_interpreter = SelectInterpreter::try_create(
self.ctx.clone(),
*(bind_context.clone()),
*s_expr.clone(),
s_expr.planned()?.clone(),
metadata.clone(),
formatted_ast.clone(),
false,
Expand Down
48 changes: 29 additions & 19 deletions src/query/service/src/interpreters/interpreter_explain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,8 @@ use databend_common_sql::ColumnSet;
use databend_common_sql::FormatOptions;
use databend_common_sql::MetadataRef;
use databend_common_sql::binder::ExplainConfig;
use databend_common_sql::optimizer::ir::PExpr;
use databend_common_sql::optimizer::ir::PlannedQuery;
use databend_common_sql::optimizer::ir::StatContext;
use databend_common_sql::plans::Mutation;
use databend_common_storages_basic::ResultCacheReader;
Expand Down Expand Up @@ -71,7 +73,6 @@ use crate::sessions::TableContextQueryIdentity;
use crate::sessions::TableContextQueryProfile;
use crate::sessions::TableContextRuntimeFilter;
use crate::sessions::TableContextSettings;
use crate::sql::optimizer::ir::SExpr;
use crate::sql::plans::Plan;

pub struct ExplainInterpreter {
Expand Down Expand Up @@ -116,7 +117,7 @@ impl Interpreter for ExplainInterpreter {
formatted_ast,
..
} => {
self.explain_query(s_expr, metadata, bind_context, formatted_ast)
self.explain_query(s_expr.planned()?, metadata, bind_context, formatted_ast)
.await?
}
Plan::Insert(insert_plan) => {
Expand All @@ -142,8 +143,13 @@ impl Interpreter for ExplainInterpreter {
vec!["CreateTableAsSelect:", ""],
)])];
res.extend(
self.explain_query(s_expr, metadata, bind_context, formatted_ast)
.await?,
self.explain_query(
s_expr.planned()?,
metadata,
bind_context,
formatted_ast,
)
.await?,
);
vec![DataBlock::concat(&res)?]
}
Expand All @@ -164,10 +170,10 @@ impl Interpreter for ExplainInterpreter {
schema,
metadata,
} => {
let mutation: Mutation = s_expr.plan().clone().try_into()?;
let mutation: Mutation = s_expr.mutation()?.clone();
let interpreter = MutationInterpreter::try_create(
self.ctx.clone(),
*s_expr.clone(),
s_expr.planned()?.clone(),
schema.clone(),
metadata.clone(),
)?;
Expand All @@ -187,7 +193,9 @@ impl Interpreter for ExplainInterpreter {
} => {
let ctx = self.ctx.clone();
let mut builder = PhysicalPlanBuilder::new(metadata.clone(), ctx, true);
let plan = builder.build(s_expr, bind_context.column_set()).await?;
let plan = builder
.build_query(s_expr.planned()?, bind_context.column_set())
.await?;

let metadata = metadata.read();
let mut context = FormatContext {
Expand Down Expand Up @@ -218,7 +226,7 @@ impl Interpreter for ExplainInterpreter {
..
} => {
self.explain_analyze(
s_expr,
s_expr.planned()?.expr(),
metadata,
bind_context.column_set(),
None,
Expand All @@ -227,11 +235,11 @@ impl Interpreter for ExplainInterpreter {
.await?
}
Plan::DataMutation { s_expr, .. } => {
let plan: Mutation = s_expr.plan().clone().try_into()?;
let plan: Mutation = s_expr.mutation()?.clone();
let mutation_build_info =
build_mutation_info(self.ctx.clone(), &plan, true, None).await?;
self.explain_analyze(
s_expr.child(0)?,
s_expr.planned()?.child(0)?,
&plan.metadata,
*plan.required_columns.clone(),
Some(mutation_build_info),
Expand Down Expand Up @@ -289,14 +297,14 @@ impl Interpreter for ExplainInterpreter {
..
} => {
self.explain_fragments(
*s_expr.clone(),
s_expr.planned()?.clone(),
metadata.clone(),
bind_context.column_set(),
)
.await?
}
Plan::DataMutation { s_expr, schema, .. } => {
self.explain_merge_fragments(*s_expr.clone(), schema.clone())
self.explain_merge_fragments(s_expr.planned()?.clone(), schema.clone())
.await?
}
Plan::InsertMultiTable(plan) => {
Expand Down Expand Up @@ -449,13 +457,13 @@ impl ExplainInterpreter {
#[async_backtrace::framed]
async fn explain_fragments(
&self,
s_expr: SExpr,
query: PlannedQuery,
metadata: MetadataRef,
required: ColumnSet,
) -> Result<Vec<DataBlock>> {
let ctx = self.ctx.clone();
let plan = PhysicalPlanBuilder::new(metadata.clone(), self.ctx.clone(), true)
.build(&s_expr, required)
.build_query(&query, required)
.await?;

let fragments = Fragmenter::try_create(ctx.clone())?.build_fragment(&plan)?;
Expand Down Expand Up @@ -502,7 +510,7 @@ impl ExplainInterpreter {
#[async_backtrace::framed]
async fn explain_analyze_graphical(
&self,
s_expr: &SExpr,
s_expr: &PExpr,
metadata: &MetadataRef,
required: ColumnSet,
ignore_result: bool,
Expand All @@ -525,7 +533,7 @@ impl ExplainInterpreter {
#[async_backtrace::framed]
async fn explain_analyze(
&self,
s_expr: &SExpr,
s_expr: &PExpr,
metadata: &MetadataRef,
required: ColumnSet,
mutation_build_info: Option<MutationBuildInfo>,
Expand Down Expand Up @@ -625,7 +633,7 @@ impl ExplainInterpreter {

async fn explain_query(
&self,
s_expr: &SExpr,
query: &PlannedQuery,
metadata: &MetadataRef,
bind_context: &BindContext,
formatted_ast: &Option<String>,
Expand All @@ -636,15 +644,17 @@ impl ExplainInterpreter {
// we should not use `dry_run` mode to build the physical plan.
// It's because we need to get the same partitions as the original selecting plan.
let mut builder = PhysicalPlanBuilder::new(metadata.clone(), ctx, formatted_ast.is_none());
let mut plan = builder.build(s_expr, bind_context.column_set()).await?;
let mut plan = builder
.build_query(query, bind_context.column_set())
.await?;
self.inject_pruned_partitions_stats(&mut plan, metadata)?;
self.explain_physical_plan(&plan, metadata, formatted_ast)
.await
}

async fn explain_merge_fragments(
&self,
s_expr: SExpr,
s_expr: PExpr,
schema: DataSchemaRef,
) -> Result<Vec<DataBlock>> {
let mutation: Mutation = s_expr.plan().clone().try_into()?;
Expand Down
7 changes: 3 additions & 4 deletions src/query/service/src/interpreters/interpreter_factory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,6 @@ use databend_common_config::GlobalConfig;
use databend_common_exception::ErrorCode;
use databend_common_exception::Result;
use databend_common_sql::binder::ExplainConfig;
use databend_common_sql::plans::Mutation;
use log::error;
use log::warn;

Expand Down Expand Up @@ -282,7 +281,7 @@ impl InterpreterFactory {
} => Ok(Arc::new(SelectInterpreter::try_create(
ctx,
*bind_context.clone(),
*s_expr.clone(),
s_expr.planned()?.clone(),
metadata.clone(),
formatted_ast.clone(),
*ignore_result,
Expand Down Expand Up @@ -654,10 +653,10 @@ impl InterpreterFactory {

Plan::Replace(replace) => ReplaceInterpreter::try_create(ctx, *replace.clone()),
Plan::DataMutation { s_expr, schema, .. } => {
let mutation: Mutation = s_expr.plan().clone().try_into()?;
let mutation = s_expr.mutation()?;
Ok(Arc::new(MutationInterpreter::try_create(
ctx,
*s_expr.clone(),
s_expr.planned()?.clone(),
schema.clone(),
mutation.metadata.clone(),
)?))
Expand Down
4 changes: 3 additions & 1 deletion src/query/service/src/interpreters/interpreter_insert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -572,7 +572,9 @@ impl Interpreter for InsertInterpreter {
let mut builder1 =
PhysicalPlanBuilder::new(metadata.clone(), self.ctx.clone(), false);
(
builder1.build(s_expr, bind_context.column_set()).await?,
builder1
.build_query(s_expr.planned()?, bind_context.column_set())
.await?,
bind_context.columns.clone(),
metadata,
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -325,7 +325,7 @@ impl InsertMultiTableInterpreter {
ordered_source_projection_columns(&result_columns, &source_required);
let mut required = bind_context.column_set();
required.extend(source_required);
let input_source = builder1.build(s_expr, required).await?;
let input_source = builder1.build_query(s_expr.planned()?, required).await?;

// Lazy materialization (triggered by WHERE + LIMIT) may reorder physical output
// columns and inject internal columns like _row_id. Add a reorder EvalScalar
Expand Down
8 changes: 4 additions & 4 deletions src/query/service/src/interpreters/interpreter_mutation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ use databend_common_pipeline::sinks::EmptySink;
use databend_common_sql::binder::MutationStrategy;
use databend_common_sql::binder::MutationType;
use databend_common_sql::executor::physical_plans::MutationKind;
use databend_common_sql::optimizer::ir::SExpr;
use databend_common_sql::optimizer::ir::PExpr;
use databend_common_sql::planner::MetadataRef;
use databend_common_sql::plans;
use databend_common_sql::plans::Mutation;
Expand Down Expand Up @@ -60,7 +60,7 @@ use crate::stream::DataBlockStream;

pub struct MutationInterpreter {
ctx: Arc<QueryContext>,
s_expr: SExpr,
s_expr: PExpr,
schema: DataSchemaRef,
metadata: MetadataRef,
materialized_view_refresh_target: Option<u64>,
Expand All @@ -69,7 +69,7 @@ pub struct MutationInterpreter {
impl MutationInterpreter {
pub fn try_create(
ctx: Arc<QueryContext>,
s_expr: SExpr,
s_expr: PExpr,
schema: DataSchemaRef,
metadata: MetadataRef,
) -> Result<MutationInterpreter> {
Expand All @@ -84,7 +84,7 @@ impl MutationInterpreter {

pub fn try_create_materialized_view_refresh(
ctx: Arc<QueryContext>,
s_expr: SExpr,
s_expr: PExpr,
schema: DataSchemaRef,
metadata: MetadataRef,
target_table_id: u64,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,13 @@ impl Interpreter for OptimizeCompactBlockInterpreter {
let mut build_res = PipelineBuildResult::create();
let mut builder =
PhysicalPlanBuilder::new(MetadataRef::default(), self.ctx.clone(), false);
match builder.build(&self.s_expr, ColumnSet::new()).await {
match builder
.build(
&databend_common_sql::optimizer::ir::PExpr::from(self.s_expr.clone()),
ColumnSet::new(),
)
.await
{
Ok(physical_plan) => {
build_res =
build_query_pipeline_without_render_result_set(&self.ctx, &physical_plan)
Expand Down
2 changes: 1 addition & 1 deletion src/query/service/src/interpreters/interpreter_replace.rs
Original file line number Diff line number Diff line change
Expand Up @@ -486,7 +486,7 @@ impl ReplaceInterpreter {
let select_interpreter = SelectInterpreter::try_create(
ctx.clone(),
*(bind_context.clone()),
*s_expr.clone(),
s_expr.planned()?.clone(),
metadata.clone(),
formatted_ast.clone(),
false,
Expand Down
10 changes: 5 additions & 5 deletions src/query/service/src/interpreters/interpreter_select.rs
Original file line number Diff line number Diff line change
Expand Up @@ -62,12 +62,12 @@ use crate::sessions::TableContextSettings;
use crate::sessions::TableContextTableAccess;
use crate::sessions::TableContextTelemetry;
use crate::sql::BindContext;
use crate::sql::optimizer::ir::SExpr;
use crate::sql::optimizer::ir::PlannedQuery;

/// Interpret SQL query with new SQL planner
pub struct SelectInterpreter {
ctx: Arc<QueryContext>,
s_expr: SExpr,
query: PlannedQuery,
bind_context: BindContext,
metadata: MetadataRef,
formatted_ast: Option<String>,
Expand All @@ -78,14 +78,14 @@ impl SelectInterpreter {
pub fn try_create(
ctx: Arc<QueryContext>,
bind_context: BindContext,
s_expr: SExpr,
query: PlannedQuery,
metadata: MetadataRef,
formatted_ast: Option<String>,
ignore_result: bool,
) -> Result<Self> {
Ok(SelectInterpreter {
ctx,
s_expr,
query,
bind_context,
metadata,
formatted_ast,
Expand Down Expand Up @@ -128,7 +128,7 @@ impl SelectInterpreter {
let mut builder = PhysicalPlanBuilder::new(self.metadata.clone(), self.ctx.clone(), false);
self.ctx.set_status_info("Building physical plan");
builder
.build(&self.s_expr, self.bind_context.column_set())
.build_query(&self.query, self.bind_context.column_set())
.await
}

Expand Down
2 changes: 1 addition & 1 deletion src/query/service/src/interpreters/interpreter_set.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@ impl Interpreter for SetInterpreter {
let select_interpreter = SelectInterpreter::try_create(
self.ctx.clone(),
*(bind_context.clone()),
*s_expr.clone(),
s_expr.planned()?.clone(),
metadata.clone(),
formatted_ast.clone(),
false,
Expand Down
Loading
Loading