use std::sync::Arc;
use super::super::Planner;
use super::super::util::{check_forbidden_group_by_params, get_effective_limit_literal};
use crate::err::Error;
use crate::exec::expression_registry::{ComputePoint, ExpressionRegistry, resolve_order_by_alias};
use crate::exec::field_path::FieldPath;
use crate::exec::operators::{
Aggregate, Compute, Filter, Limit, RandomShuffle, Sort, SortByKey, SortDirection, SortKey,
SortTopK, SortTopKByKey, Split,
};
#[cfg(all(storage, not(target_family = "wasm")))]
use crate::exec::operators::{ExternalSort, ExternalSortByKey};
use crate::exec::topk_pushdown::{TopKPushdownHandle, TopKPushdownReason};
use crate::exec::{ExecOperator, OperatorMetrics};
use crate::expr::field::Fields;
use crate::expr::{Cond, Expr, Idiom};
#[derive(Default)]
pub(crate) enum WhereClauseState {
#[default]
None,
Original(crate::expr::cond::Cond),
Precompiled(Arc<dyn crate::exec::PhysicalExpr>),
}
#[derive(Default)]
pub(crate) struct SelectPipelineConfig {
pub where_clause: WhereClauseState,
pub split: Option<crate::expr::split::Splits>,
pub group: Option<crate::expr::group::Groups>,
pub order: Option<crate::expr::order::Ordering>,
pub limit: Option<crate::expr::limit::Limit>,
pub start: Option<crate::expr::start::Start>,
pub omit: Vec<Expr>,
pub tempfiles: bool,
pub topk_pushdown: Option<TopKPushdownHandle>,
}
pub(crate) enum FilterAction {
UseOriginal,
FullyConsumed,
Residual(Cond),
}
pub(crate) struct PlannedSource {
pub(crate) operator: Arc<dyn ExecOperator>,
pub(crate) filter_action: FilterAction,
pub(crate) limit_pushed: bool,
pub(crate) topk_pushdown: Option<TopKPushdownHandle>,
}
pub(crate) fn filter_action_for_predicate(
scan_predicate: &Option<Arc<dyn crate::exec::PhysicalExpr>>,
) -> FilterAction {
if scan_predicate.is_some() {
FilterAction::FullyConsumed
} else {
FilterAction::UseOriginal
}
}
pub(crate) struct TopKFirstKeySpec {
pub(crate) first_key: SortKey,
pub(crate) key_count: usize,
}
pub(crate) enum TopKPushdownRequest {
NotApplicable,
Ineligible(TopKPushdownReason),
Eligible(TopKFirstKeySpec),
}
pub(crate) fn compute_topk_pushdown_request(
order: Option<&crate::expr::order::Ordering>,
start: &Option<crate::expr::start::Start>,
limit: &Option<crate::expr::limit::Limit>,
fields: &Fields,
tempfiles: bool,
disqualified: bool,
max_priority_queue_size: usize,
) -> TopKPushdownRequest {
use crate::expr::order::Ordering;
use crate::expr::part::Part;
if disqualified {
return TopKPushdownRequest::NotApplicable;
}
let Some(Ordering::Order(order_list)) = order else {
return TopKPushdownRequest::NotApplicable;
};
let Some(first) = order_list.0.first() else {
return TopKPushdownRequest::NotApplicable;
};
if limit.is_none() {
return TopKPushdownRequest::NotApplicable;
}
match get_effective_limit_literal(start, limit) {
Some(effective_limit) if effective_limit <= max_priority_queue_size => {}
_ => return TopKPushdownRequest::Ineligible(TopKPushdownReason::LimitTooLarge),
}
if tempfiles {
return TopKPushdownRequest::Ineligible(TopKPushdownReason::Tempfiles);
}
if first.collate || first.numeric {
return TopKPushdownRequest::Ineligible(TopKPushdownReason::UnsupportedOrder);
}
let idiom = &first.value;
let field_path = if let Some((resolved_expr, _alias)) = resolve_order_by_alias(idiom, fields) {
match &resolved_expr {
Expr::Idiom(inner_idiom)
if inner_idiom.len() == 1
&& !inner_idiom.0.iter().any(|p| matches!(p, Part::Lookup(_))) =>
{
match FieldPath::try_from(inner_idiom) {
Ok(path) => path,
Err(_) => {
return TopKPushdownRequest::Ineligible(
TopKPushdownReason::UnsupportedOrder,
);
}
}
}
_ => return TopKPushdownRequest::Ineligible(TopKPushdownReason::UnsupportedOrder),
}
} else {
match FieldPath::try_from(idiom) {
Ok(path) => path,
Err(_) => return TopKPushdownRequest::Ineligible(TopKPushdownReason::UnsupportedOrder),
}
};
use crate::exec::field_path::FieldPathPart;
if matches!(field_path.0.first(), Some(FieldPathPart::Field(f)) if f == "id") {
return TopKPushdownRequest::NotApplicable;
}
if crate::exec::topk_pushdown::field_path_wire_segments(&field_path).is_none() {
return TopKPushdownRequest::Ineligible(TopKPushdownReason::UnsupportedOrder);
}
let mut first_key = SortKey::new(field_path);
first_key.direction = if first.direction {
SortDirection::Asc
} else {
SortDirection::Desc
};
TopKPushdownRequest::Eligible(TopKFirstKeySpec {
first_key,
key_count: order_list.0.len(),
})
}
impl<'ctx> Planner<'ctx> {
pub(crate) async fn plan_pipeline(
&self,
source: Arc<dyn ExecOperator>,
fields: Option<Fields>,
config: SelectPipelineConfig,
) -> Result<Arc<dyn ExecOperator>, Error> {
let SelectPipelineConfig {
where_clause,
split,
group,
order,
limit,
start,
omit,
tempfiles,
topk_pushdown,
} = config;
let topk_pushdown = if split.is_some() || group.is_some() {
None
} else {
topk_pushdown
};
let filtered = match where_clause {
WhereClauseState::None => source,
WhereClauseState::Precompiled(predicate) => {
Arc::new(Filter::new(source, predicate)) as Arc<dyn ExecOperator>
}
WhereClauseState::Original(cond) => {
let predicate = self.physical_expr(cond.0).await?;
Arc::new(Filter::new(source, predicate)) as Arc<dyn ExecOperator>
}
};
let split_op = if let Some(splits) = split {
let idioms: Vec<_> = splits.into_iter().map(|s| s.0).collect();
Arc::new(Split {
input: filtered,
idioms,
metrics: Arc::new(OperatorMetrics::new()),
}) as Arc<dyn ExecOperator>
} else {
filtered
};
let fields = fields.unwrap_or_else(Fields::all);
let (grouped, skip_projections) = if let Some(groups) = group {
let group_by: Vec<_> = groups.0.into_iter().map(|g| g.0).collect();
check_forbidden_group_by_params(&fields)?;
let (aggregates, group_by_exprs) = self.plan_aggregation(&fields, &group_by).await?;
(
Arc::new(Aggregate::new(split_op, group_by, group_by_exprs, aggregates))
as Arc<dyn ExecOperator>,
true,
)
} else {
(split_op, false)
};
let mut registry = ExpressionRegistry::with_reserved_and_protected_names(
super::collect_field_names(&fields),
super::collect_simple_source_field_names(&fields),
);
let (sorted, sort_only_omits) = if let Some(order) = order {
if self.can_eliminate_sort(&grouped, &order) {
(grouped, vec![])
} else if skip_projections {
(self.plan_sort(grouped, order, &start, &limit, tempfiles).await?, vec![])
} else {
self.plan_sort_consolidated(
grouped,
order,
&fields,
&start,
&limit,
tempfiles,
&mut registry,
topk_pushdown,
)
.await?
}
} else {
(grouped, vec![])
};
let limited = if limit.is_some() || start.is_some() {
let limit_expr = match limit {
Some(l) => Some(self.physical_expr(l.0).await?),
None => None,
};
let offset_expr = match start {
Some(s) => Some(self.physical_expr(s.0).await?),
None => None,
};
Arc::new(Limit::new(sorted, limit_expr, offset_expr)) as Arc<dyn ExecOperator>
} else {
sorted
};
let mut all_omit = omit;
for field_name in sort_only_omits {
all_omit.push(Expr::Idiom(Idiom::field(field_name)));
}
let projected = if skip_projections {
if !all_omit.is_empty() {
let omit_fields = self.plan_omit(all_omit).await?;
Arc::new(crate::exec::operators::Project::new(limited, vec![], omit_fields, true))
as Arc<dyn ExecOperator>
} else {
limited
}
} else {
self.plan_projections_fast(fields, all_omit, limited, &mut registry).await?
};
Ok(projected)
}
pub(crate) fn can_eliminate_sort(
&self,
input: &Arc<dyn ExecOperator>,
order: &crate::expr::order::Ordering,
) -> bool {
use crate::exec::ordering::SortProperty;
use crate::expr::order::Ordering;
let Ordering::Order(order_list) = order else {
return false; };
let required: Vec<SortProperty> = order_list
.iter()
.filter_map(|field| {
crate::exec::field_path::FieldPath::try_from(&field.value).ok().map(|path| {
let direction = if field.direction {
SortDirection::Asc
} else {
SortDirection::Desc
};
SortProperty {
path,
direction,
collate: field.collate,
numeric: field.numeric,
}
})
})
.collect();
if required.len() != order_list.len() {
return false;
}
let constant_fields = input.constant_output_fields();
let required: Vec<SortProperty> =
required.into_iter().skip_while(|prop| constant_fields.contains(&prop.path)).collect();
if required.is_empty() {
return true;
}
input.output_ordering().satisfies(&required)
}
pub(crate) async fn plan_sort(
&self,
input: Arc<dyn ExecOperator>,
order: crate::expr::order::Ordering,
start: &Option<crate::expr::start::Start>,
limit: &Option<crate::expr::limit::Limit>,
#[allow(unused)] tempfiles: bool,
) -> Result<Arc<dyn ExecOperator>, Error> {
use crate::expr::order::Ordering;
match order {
Ordering::Random => {
let effective_limit = get_effective_limit_literal(start, limit);
Ok(Arc::new(RandomShuffle::new(input, effective_limit)) as Arc<dyn ExecOperator>)
}
Ordering::Order(order_list) => {
let order_by = self.convert_order_list(order_list).await?;
#[cfg(all(storage, not(target_family = "wasm")))]
if tempfiles && let Some(temp_dir) = self.ctx.temporary_directory() {
return Ok(
Arc::new(ExternalSort::new(input, order_by, temp_dir.to_path_buf()))
as Arc<dyn ExecOperator>,
);
}
if let Some(effective_limit) = get_effective_limit_literal(start, limit)
&& effective_limit
<= self.ctx.config.max_order_limit_priority_queue_size as usize
{
return Ok(Arc::new(SortTopK::new(input, order_by, effective_limit))
as Arc<dyn ExecOperator>);
}
Ok(Arc::new(Sort::new(input, order_by)) as Arc<dyn ExecOperator>)
}
}
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn plan_sort_consolidated(
&self,
input: Arc<dyn ExecOperator>,
order: crate::expr::order::Ordering,
fields: &Fields,
start: &Option<crate::expr::start::Start>,
limit: &Option<crate::expr::limit::Limit>,
#[allow(unused)] tempfiles: bool,
registry: &mut ExpressionRegistry,
topk_pushdown: Option<TopKPushdownHandle>,
) -> Result<(Arc<dyn ExecOperator>, Vec<String>), Error> {
use crate::expr::order::Ordering;
use crate::expr::part::Part;
match order {
Ordering::Random => {
let effective_limit = get_effective_limit_literal(start, limit);
Ok((
Arc::new(RandomShuffle::new(input, effective_limit)) as Arc<dyn ExecOperator>,
vec![],
))
}
Ordering::Order(order_list) => {
let mut sort_keys = Vec::with_capacity(order_list.len());
let mut sort_only_fields: Vec<String> = Vec::new();
for order_field in order_list.iter() {
let idiom = &order_field.value;
let field_path = if let Some((resolved_expr, alias)) =
resolve_order_by_alias(idiom, fields)
{
match &resolved_expr {
Expr::Idiom(inner_idiom) => {
if inner_idiom.len() > 1
|| inner_idiom.0.iter().any(|p| matches!(p, Part::Lookup(_)))
{
let name = registry
.register(
&resolved_expr,
ComputePoint::Sort,
Some(alias.clone()),
self,
)
.await?;
FieldPath::field(name)
} else {
match FieldPath::try_from(inner_idiom) {
Ok(path) => path,
Err(_) => {
let name = registry
.register(
&resolved_expr,
ComputePoint::Sort,
Some(alias.clone()),
self,
)
.await?;
FieldPath::field(name)
}
}
}
}
_ => {
let name = registry
.register(
&resolved_expr,
ComputePoint::Sort,
Some(alias.clone()),
self,
)
.await?;
FieldPath::field(name)
}
}
} else {
match FieldPath::try_from(idiom) {
Ok(path) => path,
Err(_) => {
let expr = Expr::Idiom(idiom.clone());
let name = registry
.register(&expr, ComputePoint::Sort, None, self)
.await?;
sort_only_fields.push(name.clone());
FieldPath::field(name)
}
}
};
let direction = if order_field.direction {
SortDirection::Asc
} else {
SortDirection::Desc
};
let mut key = SortKey::new(field_path);
key.direction = direction;
key.collate = order_field.collate;
key.numeric = order_field.numeric;
sort_keys.push(key);
}
let computed = if registry.has_expressions_for_point(ComputePoint::Sort) {
let compute_fields = registry
.get_expressions_for_point(ComputePoint::Sort)
.into_iter()
.map(|(name, expr)| (crate::val::Strand::new(name), expr))
.collect();
Arc::new(Compute::new(input, compute_fields)) as Arc<dyn ExecOperator>
} else {
input
};
#[cfg(all(storage, not(target_family = "wasm")))]
if tempfiles && let Some(temp_dir) = self.ctx.temporary_directory() {
return Ok((
Arc::new(ExternalSortByKey::new(
computed,
sort_keys,
temp_dir.to_path_buf(),
)) as Arc<dyn ExecOperator>,
sort_only_fields,
));
}
if let Some(effective_limit) = get_effective_limit_literal(start, limit)
&& effective_limit
<= self.ctx.config.max_order_limit_priority_queue_size as usize
{
let mut topk = SortTopKByKey::new(computed, sort_keys, effective_limit);
if let Some(handle) = topk_pushdown
&& topk.sort_keys.first() == Some(&handle.expected_first_key)
&& topk.sort_keys.len() == handle.expected_key_count
{
topk = topk.with_threshold_cell(handle.cell);
}
return Ok((Arc::new(topk) as Arc<dyn ExecOperator>, sort_only_fields));
}
Ok((
Arc::new(SortByKey::new(computed, sort_keys)) as Arc<dyn ExecOperator>,
sort_only_fields,
))
}
}
}
}