mirror of
https://github.com/rustfs/rustfs.git
synced 2026-07-27 16:48:58 +00:00
@@ -23,8 +23,6 @@ use tracing::debug;
|
||||
|
||||
use crate::sql::analyzer::DefaultAnalyzer;
|
||||
|
||||
const PUSH_DOWN_PROJECTION_INDEX: usize = 24;
|
||||
|
||||
pub trait LogicalOptimizer: Send + Sync {
|
||||
fn optimize(&self, plan: &QueryPlan, session: &SessionCtx) -> QueryResult<LogicalPlan>;
|
||||
|
||||
@@ -93,7 +91,7 @@ impl Default for DefaultLogicalOptimizer {
|
||||
|
||||
impl LogicalOptimizer for DefaultLogicalOptimizer {
|
||||
fn optimize(&self, plan: &QueryPlan, session: &SessionCtx) -> QueryResult<LogicalPlan> {
|
||||
let analyzed_plan = { self.analyzer.analyze(&plan.df_plan, session).map(|p| p).map_err(|e| e)? };
|
||||
let analyzed_plan = { self.analyzer.analyze(&plan.df_plan, session)? };
|
||||
|
||||
debug!("Analyzed logical plan:\n{}\n", plan.df_plan.display_indent_schema(),);
|
||||
|
||||
@@ -101,9 +99,7 @@ impl LogicalOptimizer for DefaultLogicalOptimizer {
|
||||
SessionStateBuilder::new_from_existing(session.inner().clone())
|
||||
.with_optimizer_rules(self.rules.clone())
|
||||
.build()
|
||||
.optimize(&analyzed_plan)
|
||||
.map(|p| p)
|
||||
.map_err(|e| e)?
|
||||
.optimize(&analyzed_plan)?
|
||||
};
|
||||
|
||||
Ok(optimizeed_plan)
|
||||
|
||||
@@ -31,18 +31,14 @@ impl Optimizer for CascadeOptimizer {
|
||||
let physical_plan = {
|
||||
self.physical_planner
|
||||
.create_physical_plan(&optimized_logical_plan, session)
|
||||
.await
|
||||
.map(|p| p)
|
||||
.map_err(|err| err)?
|
||||
.await?
|
||||
};
|
||||
|
||||
debug!("Original physical plan:\n{}\n", displayable(physical_plan.as_ref()).indent(false));
|
||||
|
||||
let optimized_physical_plan = {
|
||||
self.physical_optimizer
|
||||
.optimize(physical_plan, session)
|
||||
.map(|p| p)
|
||||
.map_err(|err| err)?
|
||||
.optimize(physical_plan, session)?
|
||||
};
|
||||
|
||||
Ok(optimized_physical_plan)
|
||||
|
||||
@@ -82,12 +82,7 @@ impl<'a> ExtParser<'a> {
|
||||
|
||||
/// Parse a new expression
|
||||
fn parse_statement(&mut self) -> Result<ExtStatement> {
|
||||
match self.parser.peek_token().token {
|
||||
Token::Word(w) => match w.keyword {
|
||||
_ => Ok(ExtStatement::SqlStatement(Box::new(self.parser.parse_statement()?))),
|
||||
},
|
||||
_ => Ok(ExtStatement::SqlStatement(Box::new(self.parser.parse_statement()?))),
|
||||
}
|
||||
Ok(ExtStatement::SqlStatement(Box::new(self.parser.parse_statement()?)))
|
||||
}
|
||||
|
||||
// Report unexpected token
|
||||
|
||||
@@ -13,14 +13,14 @@ use datafusion::sql::{planner::SqlToRel, sqlparser::ast::Statement};
|
||||
use crate::metadata::ContextProviderExtension;
|
||||
|
||||
pub struct SqlPlanner<'a, S: ContextProviderExtension> {
|
||||
schema_provider: &'a S,
|
||||
_schema_provider: &'a S,
|
||||
df_planner: SqlToRel<'a, S>,
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl<'a, S: ContextProviderExtension + Send + Sync> LogicalPlanner for SqlPlanner<'a, S> {
|
||||
impl<S: ContextProviderExtension + Send + Sync> LogicalPlanner for SqlPlanner<'_, S> {
|
||||
async fn create_logical_plan(&self, statement: ExtStatement, session: &SessionCtx) -> QueryResult<Plan> {
|
||||
let plan = { self.statement_to_plan(statement, session).await.map_err(|err| err)? };
|
||||
let plan = { self.statement_to_plan(statement, session).await? };
|
||||
|
||||
Ok(plan)
|
||||
}
|
||||
@@ -30,7 +30,7 @@ impl<'a, S: ContextProviderExtension + Send + Sync + 'a> SqlPlanner<'a, S> {
|
||||
/// Create a new query planner
|
||||
pub fn new(schema_provider: &'a S) -> Self {
|
||||
SqlPlanner {
|
||||
schema_provider,
|
||||
_schema_provider: schema_provider,
|
||||
df_planner: SqlToRel::new(schema_provider),
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user