Skip to content
Merged
Show file tree
Hide file tree
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
36 changes: 36 additions & 0 deletions vortex-file/src/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@ use vortex_array::expr::lt_eq;
use vortex_array::expr::or;
use vortex_array::expr::root;
use vortex_array::expr::select;
use vortex_array::expr::stats::Stat;
use vortex_array::extension::datetime::TimeUnit;
use vortex_array::extension::datetime::Timestamp;
use vortex_array::extension::datetime::TimestampOptions;
Expand Down Expand Up @@ -2126,6 +2127,41 @@ async fn test_writer_with_statistics() -> VortexResult<()> {
Ok(())
}

#[tokio::test]
async fn file_sum_is_absent_when_a_chunk_overflows() -> VortexResult<()> {
let dtype = DType::Struct(
StructFields::from_iter([("numbers", DType::from(PType::I64))]),
Nullability::NonNullable,
);
let mut buf = ByteBufferMut::empty();
let mut writer = SESSION
.write_options()
.with_file_statistics(vec![Stat::Sum])
.writer(&mut buf, dtype);

// The first chunk overflows, so the file sum must not be the second chunk's sum of 2.
for chunk in [buffer![i64::MAX, 1], buffer![2i64]] {
writer
.push(StructArray::from_fields(&[("numbers", chunk.into_array())])?.into_array())
.await?;
}

let summary = writer.finish().await?;
let footer_stats = summary
.footer()
.statistics()
.vortex_expect("file statistics were requested");
assert!(footer_stats.stats_sets()[0].get(Stat::Sum).is_absent());

let file = SESSION.open_options().open_buffer(buf)?;
let file_stats = file
.file_stats()
.vortex_expect("file statistics were written");
assert!(file_stats.stats_sets()[0].get(Stat::Sum).is_absent());

Ok(())
}

#[tokio::test]
async fn test_file_metadata_roundtrip() -> VortexResult<()> {
let array =
Expand Down
37 changes: 37 additions & 0 deletions vortex-layout/src/layouts/file_stats.rs
Original file line number Diff line number Diff line change
Expand Up @@ -158,6 +158,13 @@ impl StatsAccumulator {
// upper bound, so aggregating by skipping nulls would be unsound.
continue;
}
Stat::Sum if !values.all_valid(ctx)? => {
// `values` holds each chunk's `Stat::Sum` from `push_chunk`. The legacy `Sum`
// aggregate behind it returns zero for empty input and null only on overflow.
// Summing `values` below would skip those nulls and report a wrong exact
// total, so leave the file sum unset.
continue;
}
Stat::Min | Stat::Max | Stat::Sum => {
if let Some(s) = values.statistics().compute_stat(stat, ctx)?
&& let Some(v) = s.into_value()
Expand Down Expand Up @@ -531,6 +538,7 @@ mod tests {
use vortex_array::IntoArray;
use vortex_array::array_session;
use vortex_array::arrays::BoolArray;
use vortex_array::arrays::PrimitiveArray;
use vortex_array::arrays::bool::BoolArrayExt;
use vortex_array::builders::VarBinViewBuilder;
use vortex_buffer::BitBuffer;
Expand Down Expand Up @@ -636,4 +644,33 @@ mod tests {
&[Stat::Max.name(), Stat::Min.name(), Stat::Sum.name()]
);
}

#[rstest]
#[case::one_chunk_overflows(vec![vec![Some(i64::MAX), Some(1)], vec![Some(2)]], None)]
#[case::all_chunks_overflow(vec![vec![Some(i64::MAX), Some(1)]; 2], None)]
#[case::total_overflows(vec![vec![Some(i64::MAX)], vec![Some(1)]], None)]
#[case::nullable_values(vec![vec![None, Some(3)], vec![Some(4), None]], Some(7))]
#[case::all_null_chunk(vec![vec![None], vec![Some(5)]], Some(5))]
fn combines_chunk_sums(
#[case] chunks: Vec<Vec<Option<i64>>>,
#[case] expected: Option<i64>,
) -> VortexResult<()> {
let mut ctx = array_session().create_execution_ctx();
let dtype = DType::Primitive(PType::I64, Nullability::Nullable);
let mut acc = StatsAccumulator::new(&dtype, &[Stat::Sum], 12);
for chunk in chunks {
acc.push_chunk(
&PrimitiveArray::from_option_iter(chunk).into_array(),
&mut ctx,
)?;
}

let stats = acc.as_stats_set(&[Stat::Sum], &mut ctx)?;
assert_eq!(
stats.get(Stat::Sum),
expected.map_or(Precision::Absent, Precision::exact)
);

Ok(())
}
}
Loading