Skip to content
Draft
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -281,6 +281,7 @@ tpchgen-arrow = { version = "2.0.2", git = "https://github.com/clflushopt/tpchge
tracing = { version = "0.1.41", default-features = false }
tracing-perfetto = "0.1.5"
tracing-subscriber = "0.3"
twox-hash = "2.1.2"
url = "2.5.7"
uuid = { version = "1.23", features = ["js"] }
wasm-bindgen-futures = "0.4.58"
Expand Down
1 change: 1 addition & 0 deletions vortex-layout/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@ sketches-ddsketch = { workspace = true }
termtree = { workspace = true }
tokio = { workspace = true, features = ["rt"], optional = true }
tracing = { workspace = true }
twox-hash = { workspace = true }
vortex-array = { workspace = true }
vortex-arrow = { workspace = true }
vortex-btrblocks = { workspace = true }
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

use vortex_array::ExecutionCtx;
use vortex_array::arrays::BoolArray;
use vortex_array::arrays::bool::BoolArrayExt;
use vortex_error::VortexResult;
use vortex_mask::AllOr;

use crate::layouts::zoned::aggregates::bloom_filter::BloomPartial;

/// Similar to [vortex_array::aggregate_fn::fns::min_max::accumulate_bool]
pub(super) fn accumulate_bool(
array: &BoolArray,
partial: &mut BloomPartial,
ctx: &mut ExecutionCtx,
) -> VortexResult<()> {
let mask = array.validity()?.execute_mask(array.len(), ctx)?;
let bits = array.bit_buffer_view();

let (true_count, valid_count) = match mask.bit_buffer() {
AllOr::None => return Ok(()),
AllOr::All => (bits.true_count() as u64, array.len() as u64),
AllOr::Some(validity) => {
let masked = bits.to_bit_buffer() & validity;
(masked.true_count() as u64, validity.true_count() as u64)
}
};

if true_count > 0 {
// True present
partial.insert([0x1]);
}
if true_count < valid_count {
// False present
partial.insert([0x0]);
}

Ok(())
}

#[cfg(test)]
mod tests {
use rstest::rstest;
use vortex_array::IntoArray;
use vortex_array::arrays::BoolArray;
use vortex_array::dtype::DType;
use vortex_array::dtype::Nullability;
use vortex_array::scalar::Scalar;
use vortex_error::VortexResult;
use vortex_error::vortex_err;

use crate::layouts::zoned::aggregates::bloom_filter::test_utils::build_filter;
use crate::layouts::zoned::aggregates::bloom_filter::test_utils::setup;

#[rstest]
#[case::inserts_each_valid_boolean_value(
&[Some(true), Some(false)],
Nullability::NonNullable,
true,
true
)]
#[case::inserts_only_false(
&[Some(false), Some(false)],
Nullability::NonNullable,
false,
true
)]
#[case::inserts_only_true(
&[Some(true), Some(true)],
Nullability::NonNullable,
true,
false
)]
#[case::ignores_null_boolean_values(
&[Some(true), None, Some(true)],
Nullability::Nullable,
true,
false
)]
#[case::all_null_booleans_leave_the_filter_empty(
&[None, None],
Nullability::Nullable,
false,
false
)]
fn membership(
#[case] values: &[Option<bool>],
#[case] nullability: Nullability,
#[case] expect_true: bool,
#[case] expect_false: bool,
) -> VortexResult<()> {
let ctx = setup()?;
let array = match nullability {
Nullability::NonNullable => BoolArray::from_iter(
values
.iter()
.copied()
.collect::<Option<Vec<_>>>()
.ok_or_else(|| vortex_err!("non-null test case contains a null"))?,
),
Nullability::Nullable => BoolArray::from_iter(values.iter().copied()),
};
let bloom_filter = build_filter(array.into_array(), DType::Bool(nullability), ctx)?;

assert_eq!(
bloom_filter.contains_valid_scalar(&Scalar::bool(true, nullability))?,
expect_true
);
assert_eq!(
bloom_filter.contains_valid_scalar(&Scalar::bool(false, nullability))?,
expect_false
);

Ok(())
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

use vortex_array::Array;
use vortex_array::ExecutionCtx;
use vortex_array::arrays::Decimal;
use vortex_array::match_each_decimal_value_type;
use vortex_error::VortexResult;
use vortex_mask::Mask;

use crate::layouts::zoned::aggregates::bloom_filter::BloomPartial;

pub(super) fn accumulate_decimal(
array: &Array<Decimal>,
partial: &mut BloomPartial,
ctx: &mut ExecutionCtx,
) -> VortexResult<()> {
match_each_decimal_value_type!(array.values_type(), |D| {
match array.validity()?.execute_mask(array.len(), ctx)? {
Mask::AllTrue(_) => {
array
.buffer::<D>()
.iter()
.for_each(|value| partial.insert(value.to_le_bytes()));
}
Mask::AllFalse(_) => {}
Mask::Values(v) => {
array
.buffer::<D>()
.iter()
.zip(v.bit_buffer().iter())
.for_each(|(value, valid)| {
if valid {
partial.insert(value.to_le_bytes())
}
});
}
}
});

Ok(())
}

#[cfg(test)]
mod tests {
use rstest::rstest;
use vortex_array::IntoArray;
use vortex_array::arrays::DecimalArray;
use vortex_array::dtype::DType;
use vortex_array::dtype::DecimalDType;
use vortex_array::dtype::NativeDecimalType;
use vortex_array::dtype::Nullability;
use vortex_array::scalar::DecimalValue;
use vortex_array::scalar::Scalar;
use vortex_error::VortexResult;

use crate::layouts::zoned::aggregates::bloom_filter::test_utils::build_filter;
use crate::layouts::zoned::aggregates::bloom_filter::test_utils::setup;

#[rstest]
#[case(3u8, 0i8, &[1i8, 2, 3], 99i8)]
#[case(5u8, 1i8, &[10i16, 20, 30], 99i16)]
#[case(9u8, 2i8, &[1000i32, 2000, 3000], 99999i32)]
#[case(18u8, 2i8, &[1000i64, 2000, 3000], 99999i64)]
#[case(10u8, 2i8, &[1000i128, 2000, 3000], 99999i128)]
fn membership<T>(
#[case] precision: u8,
#[case] scale: i8,
#[case] present: &[T],
#[case] absent: T,
) -> VortexResult<()>
where
T: Copy + Into<DecimalValue> + NativeDecimalType,
{
let ctx = setup()?;
let decimal_dtype = DecimalDType::new(precision, scale);
let dtype = DType::Decimal(decimal_dtype, Nullability::NonNullable);
let values: DecimalArray = DecimalArray::from_iter(present.iter().copied(), decimal_dtype);
let bloom_filter = build_filter(values.into_array(), dtype, ctx)?;

for &v in present {
let scalar = Scalar::decimal(v.into(), decimal_dtype, Nullability::NonNullable);
assert!(bloom_filter.contains_valid_scalar(&scalar)?);
}

let absent_scalar = Scalar::decimal(absent.into(), decimal_dtype, Nullability::NonNullable);
assert!(!bloom_filter.contains_valid_scalar(&absent_scalar)?);

Ok(())
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

use vortex_array::Canonical;
use vortex_array::ExecutionCtx;
use vortex_array::arrays::ExtensionArray;
use vortex_array::arrays::extension::ExtensionArrayExt;
use vortex_error::VortexResult;

use crate::layouts::zoned::aggregates::bloom_filter::BloomPartial;

pub(super) fn accumulate_extension(
array: &ExtensionArray,
partial: &mut BloomPartial,
ctx: &mut ExecutionCtx,
) -> VortexResult<()> {
let storage = array.storage_array().clone();
let canonical = storage.execute::<Canonical>(ctx)?;

super::accumulate_canonical(&canonical, partial, ctx)
}

#[cfg(test)]
mod tests {
use vortex_array::IntoArray;
use vortex_array::arrays::ExtensionArray;
use vortex_array::arrays::PrimitiveArray;
use vortex_array::dtype::DType;
use vortex_array::dtype::Nullability;
use vortex_array::extension::datetime::TimeUnit;
use vortex_array::extension::datetime::Timestamp;
use vortex_array::scalar::Scalar;
use vortex_error::VortexResult;

use crate::layouts::zoned::aggregates::bloom_filter::test_utils::build_filter;
use crate::layouts::zoned::aggregates::bloom_filter::test_utils::setup;

#[test]
fn hashes_extension_values_through_storage() -> VortexResult<()> {
let ctx = setup()?;
let ext_dtype = Timestamp::new(TimeUnit::Milliseconds, Nullability::NonNullable).erased();
let bloom_filter = build_filter(
ExtensionArray::new(
ext_dtype.clone(),
PrimitiveArray::from_iter([1_000i64, 2_000, 3_000]).into_array(),
)
.into_array(),
DType::Extension(ext_dtype.clone()),
ctx,
)?;

for value in [1_000i64, 2_000, 3_000] {
let scalar = Scalar::extension_ref(
ext_dtype.clone(),
Scalar::primitive(value, Nullability::NonNullable),
);
assert!(bloom_filter.contains_valid_scalar(&scalar)?);
}

let absent = Scalar::extension_ref(
ext_dtype,
Scalar::primitive(4_000i64, Nullability::NonNullable),
);
assert!(!bloom_filter.contains_valid_scalar(&absent)?);
Ok(())
}

#[test]
fn ignores_null_extension_values() -> VortexResult<()> {
let ctx = setup()?;
let ext_dtype = Timestamp::new(TimeUnit::Milliseconds, Nullability::Nullable).erased();
let bloom_filter = build_filter(
ExtensionArray::new(
ext_dtype.clone(),
PrimitiveArray::from_option_iter([Some(1_000i64), None, Some(3_000)]).into_array(),
)
.into_array(),
DType::Extension(ext_dtype.clone()),
ctx,
)?;

let present = Scalar::extension_ref(
ext_dtype.clone(),
Scalar::primitive(1_000i64, Nullability::Nullable),
);
let null_slot_value =
Scalar::extension_ref(ext_dtype, Scalar::primitive(0i64, Nullability::Nullable));
assert!(bloom_filter.contains_valid_scalar(&present)?);
assert!(!bloom_filter.contains_valid_scalar(&null_slot_value)?);
Ok(())
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

use vortex_array::Canonical;
use vortex_array::ExecutionCtx;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;

mod bool;
mod decimal;
mod extension;
mod primitive;
mod varbin;

use crate::layouts::zoned::aggregates::bloom_filter::BloomPartial;
use crate::layouts::zoned::aggregates::bloom_filter::canonical::bool::accumulate_bool;
use crate::layouts::zoned::aggregates::bloom_filter::canonical::decimal::accumulate_decimal;
use crate::layouts::zoned::aggregates::bloom_filter::canonical::extension::accumulate_extension;
use crate::layouts::zoned::aggregates::bloom_filter::canonical::primitive::accumulate_primitive;
use crate::layouts::zoned::aggregates::bloom_filter::canonical::varbin::accumulate_varbin;

pub(super) fn accumulate_canonical(
canonical: &Canonical,
partial: &mut BloomPartial,
ctx: &mut ExecutionCtx,
) -> VortexResult<()> {
match canonical {
Canonical::Bool(array) => accumulate_bool(array, partial, ctx)?,
Canonical::Primitive(array) => accumulate_primitive(array, partial, ctx)?,
Canonical::Decimal(array) => accumulate_decimal(array, partial, ctx)?,
Canonical::VarBinView(array) => accumulate_varbin(array, partial, ctx)?,
Canonical::Extension(array) => accumulate_extension(array, partial, ctx)?,

// Nulls are skipped and are not included in any Bloom filter.
Canonical::Null(_) => {}

// TODO (joacoc): pending canonical
Canonical::Struct(_)
| Canonical::List(_)
| Canonical::FixedSizeList(_)
| Canonical::Variant(_)
| Canonical::Union(_)
| Canonical::Map(_) => {
vortex_bail!(
"Unsupported canonical type for bloom filter: {}",
canonical.dtype()
)
}
}

Ok(())
}
Loading