This is an automated email from the ASF dual-hosted git repository.
agrove pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/arrow.git
The following commit(s) were added to refs/heads/master by this push:
new e2f59c0 ARROW-4748: [Rust] [DataFusion] Optimize GROUP BY aggregate
queries
e2f59c0 is described below
commit e2f59c02c54c15bee77b0a82fa48334f74c3c6b4
Author: Andy Grove <[email protected]>
AuthorDate: Sun Oct 13 09:29:28 2019 -0600
ARROW-4748: [Rust] [DataFusion] Optimize GROUP BY aggregate queries
Various optimizations to GROUP BY aggregate queries, resulting in a 25%
improvement in the current benchmark, and close to a 2x improvement when
running queries against large data sets.
- Re-use vector for grouping keys instead of creating a new vector for each
row (benchmark went from 47 us to 39 us)
- Two pass approach where pass one builds a vector of references to
accumulators, and the second pass processes one column at a time, updating the
accumulators (now 37.6 us)
- Create aggregate schema once when plan is constructed, instead of once
per batch per partition (now 36.5 us)
- Avoid converting final aggregate map into a vector (now 34.8 us)
Closes #5636 from andygrove/ARROW-4748 and squashes the following commits:
d92bf8ea5 <Andy Grove> avoid converting final aggregate map into a vec
3937c24e9 <Andy Grove> create aggregate schema once, not once per batch
18330fd99 <Andy Grove> Minor optimizations
5fdde93b0 <Andy Grove> Create one vec for the grouping key value instead of
creating one per row
Authored-by: Andy Grove <[email protected]>
Signed-off-by: Andy Grove <[email protected]>
---
rust/datafusion/src/execution/context.rs | 7 +-
.../src/execution/physical_plan/hash_aggregate.rs | 426 ++++++++++++---------
2 files changed, 244 insertions(+), 189 deletions(-)
diff --git a/rust/datafusion/src/execution/context.rs
b/rust/datafusion/src/execution/context.rs
index 92a03b6..6ea3f8b 100644
--- a/rust/datafusion/src/execution/context.rs
+++ b/rust/datafusion/src/execution/context.rs
@@ -274,7 +274,7 @@ impl ExecutionContext {
input,
group_expr,
aggr_expr,
- schema,
+ ..
} => {
let input = self.create_physical_plan(input, batch_size)?;
let input_schema = input.as_ref().schema().clone();
@@ -287,10 +287,7 @@ impl ExecutionContext {
.map(|e| self.create_aggregate_expr(e, &input_schema))
.collect::<Result<Vec<_>>>()?;
Ok(Arc::new(HashAggregateExec::try_new(
- group_expr,
- aggr_expr,
- input,
- schema.clone(),
+ group_expr, aggr_expr, input,
)?))
}
LogicalPlan::Selection { input, expr, .. } => {
diff --git a/rust/datafusion/src/execution/physical_plan/hash_aggregate.rs
b/rust/datafusion/src/execution/physical_plan/hash_aggregate.rs
index d048b0b..da35f13 100644
--- a/rust/datafusion/src/execution/physical_plan/hash_aggregate.rs
+++ b/rust/datafusion/src/execution/physical_plan/hash_aggregate.rs
@@ -28,8 +28,8 @@ use crate::execution::physical_plan::{
};
use arrow::array::{
- ArrayRef, BinaryArray, Int16Array, Int32Array, Int64Array, Int8Array,
UInt16Array,
- UInt32Array, UInt64Array, UInt8Array,
+ ArrayRef, BinaryArray, Float32Array, Float64Array, Int16Array, Int32Array,
+ Int64Array, Int8Array, UInt16Array, UInt32Array, UInt64Array, UInt8Array,
};
use arrow::array::{
BinaryBuilder, Float32Builder, Float64Builder, Int16Builder, Int32Builder,
@@ -38,7 +38,6 @@ use arrow::array::{
use arrow::datatypes::{DataType, Field, Schema};
use arrow::record_batch::RecordBatch;
-use crate::execution::physical_plan::common::get_scalar_value;
use crate::execution::physical_plan::expressions::Column;
use crate::execution::physical_plan::merge::MergeExec;
use crate::logicalplan::ScalarValue;
@@ -58,8 +57,20 @@ impl HashAggregateExec {
group_expr: Vec<Arc<dyn PhysicalExpr>>,
aggr_expr: Vec<Arc<dyn AggregateExpr>>,
input: Arc<dyn ExecutionPlan>,
- schema: Arc<Schema>,
) -> Result<Self> {
+ let input_schema = input.schema();
+
+ let mut fields = Vec::with_capacity(group_expr.len() +
aggr_expr.len());
+ for expr in &group_expr {
+ let name = expr.name();
+ fields.push(Field::new(&name, expr.data_type(&input_schema)?,
true))
+ }
+ for expr in &aggr_expr {
+ let name = expr.name();
+ fields.push(Field::new(&name, expr.data_type(&input_schema)?,
true))
+ }
+ let schema = Arc::new(Schema::new(fields));
+
Ok(HashAggregateExec {
group_expr,
aggr_expr,
@@ -184,11 +195,11 @@ impl Partition for HashAggregatePartition {
/// Create array from `key` attribute in map entry (representing a grouping
scalar value)
macro_rules! group_array_from_map_entries {
- ($BUILDER:ident, $TY:ident, $ENTRIES:expr, $COL_INDEX:expr) => {{
- let mut builder = $BUILDER::new($ENTRIES.len());
+ ($BUILDER:ident, $TY:ident, $MAP:expr, $COL_INDEX:expr) => {{
+ let mut builder = $BUILDER::new($MAP.len());
let mut err = false;
- for j in 0..$ENTRIES.len() {
- match $ENTRIES[j].k[$COL_INDEX] {
+ for k in $MAP.keys() {
+ match k[$COL_INDEX] {
GroupByScalar::$TY(n) => builder.append_value(n).unwrap(),
_ => err = true,
}
@@ -207,11 +218,11 @@ macro_rules! group_array_from_map_entries {
/// Create array from `value` attribute in map entry (representing an
aggregate scalar
/// value)
macro_rules! aggr_array_from_map_entries {
- ($BUILDER:ident, $TY:ident, $TY2:ty, $ENTRIES:expr, $COL_INDEX:expr) => {{
- let mut builder = $BUILDER::new($ENTRIES.len());
+ ($BUILDER:ident, $TY:ident, $TY2:ty, $MAP:expr, $COL_INDEX:expr) => {{
+ let mut builder = $BUILDER::new($MAP.len());
let mut err = false;
- for j in 0..$ENTRIES.len() {
- match $ENTRIES[j].v[$COL_INDEX] {
+ for v in $MAP.values() {
+ match v[$COL_INDEX].as_ref().borrow().get_value()? {
Some(ScalarValue::$TY(n)) => builder.append_value(n as
$TY2).unwrap(),
None => builder.append_null().unwrap(),
_ => err = true,
@@ -228,6 +239,27 @@ macro_rules! aggr_array_from_map_entries {
}};
}
+/// Create array from single accumulator value
+macro_rules! aggr_array_from_accumulator {
+ ($BUILDER:ident, $TY:ident, $TY2:ty, $VALUE:expr) => {{
+ let mut builder = $BUILDER::new(1);
+ match $VALUE {
+ Some(ScalarValue::$TY(n)) => {
+ builder.append_value(n as $TY2)?;
+ Ok(Arc::new(builder.finish()) as ArrayRef)
+ }
+ None => {
+ builder.append_null()?;
+ Ok(Arc::new(builder.finish()) as ArrayRef)
+ }
+ _ => Err(ExecutionError::ExecutionError(
+ "unexpected type when creating aggregate array from aggregate
map"
+ .to_string(),
+ )),
+ }
+ }};
+}
+
#[derive(Debug)]
struct MapEntry {
k: Vec<GroupByScalar>,
@@ -260,6 +292,22 @@ impl GroupedHashAggregateIterator {
}
}
+type AccumulatorSet = Vec<Rc<RefCell<dyn Accumulator>>>;
+
+macro_rules! update_accumulators {
+ ($ARRAY:ident, $ARRAY_TY:ident, $SCALAR_TY:expr, $COL:expr, $ACCUM:expr)
=> {{
+ let primitive_array =
$ARRAY.as_any().downcast_ref::<$ARRAY_TY>().unwrap();
+
+ for row in 0..$ARRAY.len() {
+ if $ARRAY.is_valid(row) {
+ let value = Some($SCALAR_TY(primitive_array.value(row)));
+ let mut accum = $ACCUM[row][$COL].borrow_mut();
+ accum.accumulate_scalar(value)?;
+ }
+ }
+ }};
+}
+
impl BatchIterator for GroupedHashAggregateIterator {
fn schema(&self) -> Arc<Schema> {
self.schema.clone()
@@ -273,7 +321,7 @@ impl BatchIterator for GroupedHashAggregateIterator {
self.finished = true;
// create map to store accumulators for each unique grouping key
- let mut map: FnvHashMap<Vec<GroupByScalar>, Vec<Rc<RefCell<dyn
Accumulator>>>> =
+ let mut map: FnvHashMap<Vec<GroupByScalar>, Rc<AccumulatorSet>> =
FnvHashMap::default();
// iterate over all input batches and update the accumulators
@@ -295,68 +343,123 @@ impl BatchIterator for GroupedHashAggregateIterator {
.map(|expr| expr.evaluate_input(&batch))
.collect::<Result<Vec<_>>>()?;
- // iterate over each row in the batch
+ // create vector large enough to hold the grouping key
+ let mut key = Vec::with_capacity(group_values.len());
+ for _ in 0..group_values.len() {
+ key.push(GroupByScalar::UInt32(0));
+ }
+
+ // iterate over each row in the batch and create the accumulators
for each grouping key
+ let mut accumulators: Vec<Rc<AccumulatorSet>> =
+ Vec::with_capacity(batch.num_rows());
+
for row in 0..batch.num_rows() {
// create grouping key for this row
- let key = create_key(&group_values, row)?;
-
- let updated: Result<bool> = match map.get(&key) {
- Some(accumulators) => {
- let _ = accumulators
- .iter()
- .zip(aggr_input_values.iter())
- .map(|(accum, input)| {
- accum
- .borrow_mut()
- .accumulate_scalar(get_scalar_value(input,
row)?)
- })
- .collect::<Result<Vec<_>>>()?;
- Ok(true)
- }
- None => Ok(false),
- };
+ create_key(&group_values, row, &mut key)?;
- if !updated? {
- let accumulators: Vec<Rc<RefCell<dyn Accumulator>>> = self
+ if let Some(accumulator_set) = map.get(&key) {
+ accumulators.push(accumulator_set.clone());
+ } else {
+ let accumulator_set: AccumulatorSet = self
.aggr_expr
.iter()
.map(|expr| expr.create_accumulator())
.collect();
- let _ = accumulators
- .iter()
- .zip(aggr_input_values.iter())
- .map(|(accum, input)| {
- accum
- .borrow_mut()
- .accumulate_scalar(get_scalar_value(input,
row)?)
- })
- .collect::<Result<Vec<_>>>()?;
-
- map.insert(key.clone(), accumulators);
+ let accumulator_set = Rc::new(accumulator_set);
+
+ map.insert(key.clone(), accumulator_set.clone());
+ accumulators.push(accumulator_set);
}
}
+
+ // iterate over each non-grouping column in the batch and update
the accumulator
+ // for each row
+ for col in 0..aggr_input_values.len() {
+ let array = &aggr_input_values[col];
+
+ match array.data_type() {
+ DataType::Int8 => update_accumulators!(
+ array,
+ Int8Array,
+ ScalarValue::Int8,
+ col,
+ accumulators
+ ),
+ DataType::Int16 => update_accumulators!(
+ array,
+ Int16Array,
+ ScalarValue::Int16,
+ col,
+ accumulators
+ ),
+ DataType::Int32 => update_accumulators!(
+ array,
+ Int32Array,
+ ScalarValue::Int32,
+ col,
+ accumulators
+ ),
+ DataType::Int64 => update_accumulators!(
+ array,
+ Int64Array,
+ ScalarValue::Int64,
+ col,
+ accumulators
+ ),
+ DataType::UInt8 => update_accumulators!(
+ array,
+ UInt8Array,
+ ScalarValue::UInt8,
+ col,
+ accumulators
+ ),
+ DataType::UInt16 => update_accumulators!(
+ array,
+ UInt16Array,
+ ScalarValue::UInt16,
+ col,
+ accumulators
+ ),
+ DataType::UInt32 => update_accumulators!(
+ array,
+ UInt32Array,
+ ScalarValue::UInt32,
+ col,
+ accumulators
+ ),
+ DataType::UInt64 => update_accumulators!(
+ array,
+ UInt64Array,
+ ScalarValue::UInt64,
+ col,
+ accumulators
+ ),
+ DataType::Float32 => update_accumulators!(
+ array,
+ Float32Array,
+ ScalarValue::Float32,
+ col,
+ accumulators
+ ),
+ DataType::Float64 => update_accumulators!(
+ array,
+ Float64Array,
+ ScalarValue::Float64,
+ col,
+ accumulators
+ ),
+ other => {
+ return Err(ExecutionError::ExecutionError(format!(
+ "Unsupported data type {:?} for result of
aggregate expression",
+ other
+ )));
+ }
+ };
+ }
}
- let input_schema = input.schema();
- // convert the map to a vec to make it easier to build arrays
- let entries: Vec<MapEntry> = map
- .iter()
- .map(|(k, v)| {
- let aggr_values = v
- .iter()
- .map(|accum| {
- let accum = accum.borrow_mut();
- accum.get_value()
- })
- .collect::<Result<Vec<_>>>()?;
-
- Ok(MapEntry {
- k: k.clone(),
- v: aggr_values,
- })
- })
- .collect::<Result<Vec<_>>>()?;
+ let input_schema = input.schema();
// build the result arrays
let mut result_arrays: Vec<ArrayRef> =
@@ -368,35 +471,39 @@ impl BatchIterator for GroupedHashAggregateIterator {
.data_type(&input_schema)?
{
DataType::UInt8 => {
- group_array_from_map_entries!(UInt8Builder, UInt8,
entries, i)
+ group_array_from_map_entries!(UInt8Builder, UInt8, map, i)
}
DataType::UInt16 => {
- group_array_from_map_entries!(UInt16Builder, UInt16,
entries, i)
+ group_array_from_map_entries!(UInt16Builder, UInt16, map,
i)
}
DataType::UInt32 => {
- group_array_from_map_entries!(UInt32Builder, UInt32,
entries, i)
+ group_array_from_map_entries!(UInt32Builder, UInt32, map,
i)
}
DataType::UInt64 => {
- group_array_from_map_entries!(UInt64Builder, UInt64,
entries, i)
+ group_array_from_map_entries!(UInt64Builder, UInt64, map,
i)
}
DataType::Int8 => {
- group_array_from_map_entries!(Int8Builder, Int8, entries,
i)
+ group_array_from_map_entries!(Int8Builder, Int8, map, i)
}
DataType::Int16 => {
- group_array_from_map_entries!(Int16Builder, Int16,
entries, i)
+ group_array_from_map_entries!(Int16Builder, Int16, map, i)
}
DataType::Int32 => {
- group_array_from_map_entries!(Int32Builder, Int32,
entries, i)
+ group_array_from_map_entries!(Int32Builder, Int32, map, i)
}
DataType::Int64 => {
- group_array_from_map_entries!(Int64Builder, Int64,
entries, i)
+ group_array_from_map_entries!(Int64Builder, Int64, map, i)
}
DataType::Utf8 => {
let mut builder = BinaryBuilder::new(1);
- for j in 0..entries.len() {
- match &entries[j].k[i] {
+ for k in map.keys() {
+ match &k[i] {
GroupByScalar::Utf8(s) =>
builder.append_string(&s).unwrap(),
- _ => {}
+ _ => {
+ return Err(ExecutionError::ExecutionError(
+ "Unexpected value for Utf8 group
column".to_string(),
+ ))
+ }
}
}
Ok(Arc::new(builder.finish()) as ArrayRef)
@@ -411,60 +518,36 @@ impl BatchIterator for GroupedHashAggregateIterator {
for i in 0..self.aggr_expr.len() {
let aggr_data_type =
self.aggr_expr[i].data_type(&input_schema)?;
let array = match aggr_data_type {
- DataType::UInt8 => aggr_array_from_map_entries!(
- UInt64Builder,
- UInt8,
- u64,
- entries,
- i
- ),
- DataType::UInt16 => aggr_array_from_map_entries!(
- UInt64Builder,
- UInt16,
- u64,
- entries,
- i
- ),
- DataType::UInt32 => aggr_array_from_map_entries!(
- UInt64Builder,
- UInt32,
- u64,
- entries,
- i
- ),
- DataType::UInt64 => aggr_array_from_map_entries!(
- UInt64Builder,
- UInt64,
- u64,
- entries,
- i
- ),
+ DataType::UInt8 => {
+ aggr_array_from_map_entries!(UInt64Builder, UInt8,
u64, map, i)
+ }
+ DataType::UInt16 => {
+ aggr_array_from_map_entries!(UInt64Builder, UInt16,
u64, map, i)
+ }
+ DataType::UInt32 => {
+ aggr_array_from_map_entries!(UInt64Builder, UInt32,
u64, map, i)
+ }
+ DataType::UInt64 => {
+ aggr_array_from_map_entries!(UInt64Builder, UInt64,
u64, map, i)
+ }
DataType::Int8 => {
- aggr_array_from_map_entries!(Int64Builder, Int8, i64,
entries, i)
+ aggr_array_from_map_entries!(Int64Builder, Int8, i64,
map, i)
}
DataType::Int16 => {
- aggr_array_from_map_entries!(Int64Builder, Int16, i64,
entries, i)
+ aggr_array_from_map_entries!(Int64Builder, Int16, i64,
map, i)
}
DataType::Int32 => {
- aggr_array_from_map_entries!(Int64Builder, Int32, i64,
entries, i)
+ aggr_array_from_map_entries!(Int64Builder, Int32, i64,
map, i)
}
DataType::Int64 => {
- aggr_array_from_map_entries!(Int64Builder, Int64, i64,
entries, i)
+ aggr_array_from_map_entries!(Int64Builder, Int64, i64,
map, i)
+ }
+ DataType::Float32 => {
+ aggr_array_from_map_entries!(Float32Builder, Float32,
f32, map, i)
+ }
+ DataType::Float64 => {
+ aggr_array_from_map_entries!(Float64Builder, Float64,
f64, map, i)
}
- DataType::Float32 => aggr_array_from_map_entries!(
- Float32Builder,
- Float32,
- f32,
- entries,
- i
- ),
- DataType::Float64 => aggr_array_from_map_entries!(
- Float64Builder,
- Float64,
- f64,
- entries,
- i
- ),
_ => Err(ExecutionError::ExecutionError(
"Unsupported aggregate expr".to_string(),
)),
@@ -473,18 +556,7 @@ impl BatchIterator for GroupedHashAggregateIterator {
}
}
- let mut fields = vec![];
- for expr in &self.group_expr {
- let name = expr.name();
- fields.push(Field::new(&name, expr.data_type(&input_schema)?,
true))
- }
- for expr in &self.aggr_expr {
- let name = expr.name();
- fields.push(Field::new(&name, expr.data_type(&input_schema)?,
true))
- }
- let schema = Schema::new(fields);
-
- let batch = RecordBatch::try_new(Arc::new(schema), result_arrays)?;
+ let batch = RecordBatch::try_new(self.schema.clone(), result_arrays)?;
Ok(Some(batch))
}
}
@@ -555,47 +627,40 @@ impl BatchIterator for HashAggregateIterator {
// build the result arrays
let mut result_arrays: Vec<ArrayRef> =
Vec::with_capacity(self.aggr_expr.len());
- let entries = vec![MapEntry {
- k: vec![],
- v: accumulators
- .iter()
- .map(|accum| accum.borrow_mut().get_value())
- .collect::<Result<Vec<_>>>()?,
- }];
-
// aggregate values
for i in 0..self.aggr_expr.len() {
let aggr_data_type = self.aggr_expr[i].data_type(&input_schema)?;
+ let value = accumulators[i].borrow_mut().get_value()?;
let array = match aggr_data_type {
DataType::UInt8 => {
- aggr_array_from_map_entries!(UInt64Builder, UInt8, u64,
entries, i)
+ aggr_array_from_accumulator!(UInt64Builder, UInt8, u64,
value)
}
DataType::UInt16 => {
- aggr_array_from_map_entries!(UInt64Builder, UInt16, u64,
entries, i)
+ aggr_array_from_accumulator!(UInt64Builder, UInt16, u64,
value)
}
DataType::UInt32 => {
- aggr_array_from_map_entries!(UInt64Builder, UInt32, u64,
entries, i)
+ aggr_array_from_accumulator!(UInt64Builder, UInt32, u64,
value)
}
DataType::UInt64 => {
- aggr_array_from_map_entries!(UInt64Builder, UInt64, u64,
entries, i)
+ aggr_array_from_accumulator!(UInt64Builder, UInt64, u64,
value)
}
DataType::Int8 => {
- aggr_array_from_map_entries!(Int64Builder, Int8, i64,
entries, i)
+ aggr_array_from_accumulator!(Int64Builder, Int8, i64,
value)
}
DataType::Int16 => {
- aggr_array_from_map_entries!(Int64Builder, Int16, i64,
entries, i)
+ aggr_array_from_accumulator!(Int64Builder, Int16, i64,
value)
}
DataType::Int32 => {
- aggr_array_from_map_entries!(Int64Builder, Int32, i64,
entries, i)
+ aggr_array_from_accumulator!(Int64Builder, Int32, i64,
value)
}
DataType::Int64 => {
- aggr_array_from_map_entries!(Int64Builder, Int64, i64,
entries, i)
+ aggr_array_from_accumulator!(Int64Builder, Int64, i64,
value)
}
DataType::Float32 => {
- aggr_array_from_map_entries!(Float32Builder, Float32, f32,
entries, i)
+ aggr_array_from_accumulator!(Float32Builder, Float32, f32,
value)
}
DataType::Float64 => {
- aggr_array_from_map_entries!(Float64Builder, Float64, f64,
entries, i)
+ aggr_array_from_accumulator!(Float64Builder, Float64, f64,
value)
}
_ => Err(ExecutionError::ExecutionError(
"Unsupported aggregate expr".to_string(),
@@ -604,14 +669,7 @@ impl BatchIterator for HashAggregateIterator {
result_arrays.push(array?);
}
- let mut fields = vec![];
- for expr in &self.aggr_expr {
- let name = expr.name();
- fields.push(Field::new(&name, expr.data_type(&input_schema)?,
true))
- }
- let schema = Schema::new(fields);
-
- let batch = RecordBatch::try_new(Arc::new(schema), result_arrays)?;
+ let batch = RecordBatch::try_new(self.schema.clone(), result_arrays)?;
Ok(Some(batch))
}
}
@@ -632,53 +690,60 @@ enum GroupByScalar {
}
/// Create a Vec<GroupByScalar> that can be used as a map key
-fn create_key(group_by_keys: &Vec<ArrayRef>, row: usize) ->
Result<Vec<GroupByScalar>> {
- group_by_keys
- .iter()
- .map(|col| match col.data_type() {
+fn create_key(
+ group_by_keys: &Vec<ArrayRef>,
+ row: usize,
+ vec: &mut Vec<GroupByScalar>,
+) -> Result<()> {
+ for i in 0..group_by_keys.len() {
+ let col = &group_by_keys[i];
+ match col.data_type() {
DataType::UInt8 => {
let array = col.as_any().downcast_ref::<UInt8Array>().unwrap();
- Ok(GroupByScalar::UInt8(array.value(row)))
+ vec[i] = GroupByScalar::UInt8(array.value(row))
}
DataType::UInt16 => {
let array =
col.as_any().downcast_ref::<UInt16Array>().unwrap();
- Ok(GroupByScalar::UInt16(array.value(row)))
+ vec[i] = GroupByScalar::UInt16(array.value(row))
}
DataType::UInt32 => {
let array =
col.as_any().downcast_ref::<UInt32Array>().unwrap();
- Ok(GroupByScalar::UInt32(array.value(row)))
+ vec[i] = GroupByScalar::UInt32(array.value(row))
}
DataType::UInt64 => {
let array =
col.as_any().downcast_ref::<UInt64Array>().unwrap();
- Ok(GroupByScalar::UInt64(array.value(row)))
+ vec[i] = GroupByScalar::UInt64(array.value(row))
}
DataType::Int8 => {
let array = col.as_any().downcast_ref::<Int8Array>().unwrap();
- Ok(GroupByScalar::Int8(array.value(row)))
+ vec[i] = GroupByScalar::Int8(array.value(row))
}
DataType::Int16 => {
let array = col.as_any().downcast_ref::<Int16Array>().unwrap();
- Ok(GroupByScalar::Int16(array.value(row)))
+ vec[i] = GroupByScalar::Int16(array.value(row))
}
DataType::Int32 => {
let array = col.as_any().downcast_ref::<Int32Array>().unwrap();
- Ok(GroupByScalar::Int32(array.value(row)))
+ vec[i] = GroupByScalar::Int32(array.value(row))
}
DataType::Int64 => {
let array = col.as_any().downcast_ref::<Int64Array>().unwrap();
- Ok(GroupByScalar::Int64(array.value(row)))
+ vec[i] = GroupByScalar::Int64(array.value(row))
}
DataType::Utf8 => {
let array =
col.as_any().downcast_ref::<BinaryArray>().unwrap();
- Ok(GroupByScalar::Utf8(String::from(
+ vec[i] = GroupByScalar::Utf8(String::from(
str::from_utf8(array.value(row)).unwrap(),
- )))
+ ))
}
- _ => Err(ExecutionError::ExecutionError(
- "Unsupported GROUP BY data type".to_string(),
- )),
- })
- .collect::<Result<Vec<GroupByScalar>>>()
+ _ => {
+ return Err(ExecutionError::ExecutionError(
+ "Unsupported GROUP BY data type".to_string(),
+ ))
+ }
+ }
+ }
+ Ok(())
}
#[cfg(test)]
@@ -688,7 +753,6 @@ mod tests {
use crate::execution::physical_plan::csv::CsvExec;
use crate::execution::physical_plan::expressions::{col, sum};
use crate::test;
- use arrow::datatypes::Field;
#[test]
fn aggregate() -> Result<()> {
@@ -703,16 +767,10 @@ mod tests {
let aggr_expr: Vec<Arc<dyn AggregateExpr>> = vec![sum(col(3))];
- let schema = Arc::new(Schema::new(vec![
- Field::new("a", DataType::UInt32, true),
- Field::new("b", DataType::Int64, true),
- ]));
-
let partition_aggregate = HashAggregateExec::try_new(
group_expr.clone(),
aggr_expr.clone(),
Arc::new(csv),
- schema.clone(),
)?;
let result = test::execute(&partition_aggregate)?;