Skip to content
Open
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
70 changes: 69 additions & 1 deletion vortex-array/src/aggregate_fn/accumulator_grouped.rs
Original file line number Diff line number Diff line change
Expand Up @@ -273,7 +273,10 @@ impl<V: AggregateFnVTable> DynGroupedAccumulator for GroupedAccumulator<V> {
}

fn flush(&mut self) -> VortexResult<ArrayRef> {
let states = std::mem::take(&mut self.partials);
let mut states = std::mem::take(&mut self.partials);
if states.len() == 1 {
return Ok(states.pop().vortex_expect("checked one partial"));
}
Ok(ChunkedArray::try_new(states, self.partial_dtype.clone())?.into_array())
}

Expand Down Expand Up @@ -418,3 +421,68 @@ fn fixed_size_list_group_ranges(groups: &FixedSizeListArray) -> GroupRanges {
size: groups.list_size() as usize,
}
}

#[cfg(test)]
mod tests {
use vortex_error::VortexResult;

use crate::ArrayRef;
use crate::IntoArray;
use crate::aggregate_fn::DynGroupedAccumulator;
use crate::aggregate_fn::GroupedAccumulator;
use crate::aggregate_fn::NumericalAggregateOpts;
use crate::aggregate_fn::fns::count::Count;
use crate::arrays::Chunked;
use crate::arrays::PrimitiveArray;
use crate::dtype::DType;
use crate::dtype::Nullability::NonNullable;
use crate::dtype::PType;

fn accumulator() -> VortexResult<GroupedAccumulator<Count>> {
GroupedAccumulator::try_new(
Count,
NumericalAggregateOpts::default(),
DType::Primitive(PType::I32, NonNullable),
)
}

fn state(values: impl IntoIterator<Item = u64>) -> ArrayRef {
PrimitiveArray::from_iter(values).into_array()
}

#[test]
fn test_flush_single_partial_returns_original_array() -> VortexResult<()> {
let mut accumulator = accumulator()?;
let state = state([1, 2, 3]);
accumulator.push_result(state.clone())?;

let flushed = accumulator.flush()?;

assert!(ArrayRef::ptr_eq(&flushed, &state));
Ok(())
}

#[test]
fn test_flush_multiple_partials_returns_chunked_array() -> VortexResult<()> {
let mut accumulator = accumulator()?;
accumulator.push_result(state([1, 2]))?;
accumulator.push_result(state([3, 4]))?;

let flushed = accumulator.flush()?;

assert!(flushed.is::<Chunked>());
assert_eq!(flushed.len(), 4);
Ok(())
}

#[test]
fn test_flush_without_partials_returns_empty_chunked_array() -> VortexResult<()> {
let mut accumulator = accumulator()?;

let flushed = accumulator.flush()?;

assert!(flushed.is::<Chunked>());
assert!(flushed.is_empty());
Ok(())
}
}
Loading