use std::sync::Arc;
use arrow::datatypes::SchemaRef;
use arrow::record_batch::RecordBatch;
use datafusion_common::Result;
use crate::InputOrderMode;
use crate::aggregates::aggregate_hash_table::FinalMarker;
use crate::aggregates::group_values::GroupByMetrics;
use crate::aggregates::{AggregateExec, AggregateMode};
use super::common_ordered::OrderedAggregateTable;
impl OrderedAggregateTable<FinalMarker> {
pub(in crate::aggregates) fn new_with_input_order(
agg: &AggregateExec,
input_schema: &SchemaRef,
output_schema: SchemaRef,
batch_size: usize,
input_order_mode: &InputOrderMode,
group_by_metrics: GroupByMetrics,
) -> Result<Self> {
Self::new_for_mode(
agg,
input_schema,
output_schema,
Arc::clone(input_schema),
batch_size,
input_order_mode,
&AggregateMode::Final,
vec![None; agg.aggr_expr.len()],
group_by_metrics,
)
}
pub(in crate::aggregates) fn aggregate_batch(
&mut self,
batch: &RecordBatch,
) -> Result<()> {
let evaluated_batch = self.evaluate_batch(batch)?;
debug_assert_eq!(evaluated_batch.grouping_set_args.len(), 1);
self.aggregate_evaluated_batch(&evaluated_batch, true)
}
pub(in crate::aggregates) fn next_output_batch(
&mut self,
) -> Result<Option<RecordBatch>> {
self.next_output_batch_for_mode(true)
}
}