Skip to content
Open
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
2 changes: 1 addition & 1 deletion vortex-array/benches/list_contains_set.rs
Original file line number Diff line number Diff line change
Expand Up @@ -60,7 +60,7 @@ fn random_i64(len: usize) -> (Vec<i64>, Vec<i64>) {

fn bench_in_set(bencher: Bencher, set: Scalar, needles: ArrayRef) {
let session = vortex_array::array_session();
// Optimized as a scan optimizes it, so the set arrives normalized.
// Optimized as a scan optimizes it.
let expr = list_contains(lit(set), root())
.bind(needles.dtype())
.unwrap()
Expand Down
7 changes: 7 additions & 0 deletions vortex-array/src/arrays/chunked/compute/kernel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,8 @@ use crate::arrays::filter::FilterExecuteAdaptor;
use crate::arrays::slice::SliceExecuteAdaptor;
use crate::optimizer::kernels::ArrayKernelsExt;
use crate::scalar_fn::ScalarFnVTable;
use crate::scalar_fn::fns::list_contains::ListContains;
use crate::scalar_fn::fns::list_contains::ListContainsElementExecuteAdaptor;
use crate::scalar_fn::fns::mask::Mask;
use crate::scalar_fn::fns::mask::MaskExecuteAdaptor;
use crate::scalar_fn::fns::zip::Zip;
Expand All @@ -25,4 +27,9 @@ pub(crate) fn initialize(session: &VortexSession) {
kernels.register_execute_parent_kernel(Slice.id(), Chunked, SliceExecuteAdaptor(Chunked));
kernels.register_execute_parent_kernel(Dict.id(), Chunked, TakeExecuteAdaptor(Chunked));
kernels.register_execute_parent_kernel(Zip.id(), Chunked, ZipExecuteAdaptor(Chunked));
kernels.register_execute_parent_kernel(
ListContains.id(),
Chunked,
ListContainsElementExecuteAdaptor(Chunked),
);
}
129 changes: 129 additions & 0 deletions vortex-array/src/arrays/chunked/compute/list_contains.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright the Vortex contributors

use vortex_error::VortexResult;

use crate::ArrayRef;
use crate::ExecutionCtx;
use crate::IntoArray;
use crate::array::ArrayView;
use crate::arrays::Chunked;
use crate::arrays::ChunkedArray;
use crate::arrays::chunked::ChunkedArrayExt;
use crate::dtype::DType;
use crate::scalar_fn::fns::list_contains::ListContains;
use crate::scalar_fn::fns::list_contains::ListContainsElementKernel;
use crate::scalar_fn::fns::list_contains::ListContainsOptions;
use crate::scalar_fn::fns::list_contains::PreparedSet;

/// Probes each chunk of the needles against one prepared set.
///
/// A constant list is prepared first, by the execution of [`ListContains`]. Each chunk then gets a
/// lazy [`ListContains`] over a slice of the prepared set, which shares the one probe, so that the
/// kernels of the chunk's own encoding can probe it.
impl ListContainsElementKernel for Chunked {
fn list_contains(
list: &ArrayRef,
needles: ArrayView<'_, Chunked>,
options: &ListContainsOptions,
_ctx: &mut ExecutionCtx,
) -> VortexResult<Option<ArrayRef>> {
if !list.is::<PreparedSet>() {
return Ok(None);
}

let mut offset = 0;
let chunks = needles
.iter_chunks()
.map(|chunk| {
let set = list.slice(offset..offset + chunk.len())?;
offset += chunk.len();
Ok(ListContains::try_new_opts(set, chunk.clone(), *options)?.into_array())
})
.collect::<VortexResult<Vec<_>>>()?;

let dtype = DType::Bool(options.result_nullability(list.dtype(), needles.dtype()));

// SAFETY: each chunk is `list_contains` of the prepared set and a needle chunk of one dtype,
// so every chunk has the dtype of the whole result.
Ok(Some(
unsafe { ChunkedArray::new_unchecked(chunks, dtype) }.into_array(),
))
}
}

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

use crate::ArrayRef;
use crate::IntoArray;
use crate::VortexSessionExecute;
use crate::array_session;
use crate::arrays::BoolArray;
use crate::arrays::Chunked;
use crate::arrays::ChunkedArray;
use crate::arrays::ConstantArray;
use crate::arrays::PrimitiveArray;
use crate::assert_arrays_eq;
use crate::dtype::DType;
use crate::dtype::Nullability;
use crate::dtype::PType;
use crate::optimizer::ArrayOptimizer;
use crate::scalar::Scalar;
use crate::scalar_fn::fns::list_contains::ListContains;
use crate::scalar_fn::fns::list_contains::ListContainsOptions;

/// The set `{2, null}` of nullable `i32`.
fn set_with_null() -> Scalar {
let element = DType::Primitive(PType::I32, Nullability::Nullable);
Scalar::list(
element.clone(),
vec![
Scalar::primitive(2i32, Nullability::Nullable),
Scalar::null(element),
],
Nullability::NonNullable,
)
}

fn chunked_needles() -> VortexResult<ArrayRef> {
Ok(ChunkedArray::try_new(
vec![
PrimitiveArray::from_option_iter([Some(1i32), Some(2)]).into_array(),
PrimitiveArray::from_option_iter::<i32, _>([]).into_array(),
PrimitiveArray::from_option_iter([None, Some(3), Some(2)]).into_array(),
],
DType::Primitive(PType::I32, Nullability::Nullable),
)?
.into_array())
}

#[rstest]
#[case::default(
ListContainsOptions::default(),
[Some(false), Some(true), None, Some(false), Some(true)],
)]
#[case::sql(
ListContainsOptions { sql_null_semantics: true },
[None, Some(true), None, None, Some(true)],
)]
fn test_constant_list_over_chunked_needles(
#[case] options: ListContainsOptions,
#[case] expected: [Option<bool>; 5],
) -> VortexResult<()> {
let mut ctx = array_session().create_execution_ctx();
let needles = chunked_needles()?;
let list = ConstantArray::new(set_with_null(), needles.len()).into_array();

// The node stays whole, so that execution prepares the list once for every chunk.
let array = ListContains::try_new_opts(list, needles, options)?
.into_array()
.optimize()?;
assert!(!array.is::<Chunked>());

assert_arrays_eq!(array, BoolArray::from_iter(expected), &mut ctx);
Ok(())
}
}
1 change: 1 addition & 0 deletions vortex-array/src/arrays/chunked/compute/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ mod cast;
mod fill_null;
mod filter;
pub(crate) mod kernel;
mod list_contains;
mod mask;
pub(crate) mod rules;
mod slice;
Expand Down
9 changes: 9 additions & 0 deletions vortex-array/src/arrays/chunked/compute/rules.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ use crate::optimizer::rules::ArrayParentReduceRule;
use crate::optimizer::rules::ParentRuleSet;
use crate::scalar_fn::fns::cast::CastReduceAdaptor;
use crate::scalar_fn::fns::fill_null::FillNullReduceAdaptor;
use crate::scalar_fn::fns::list_contains::ListContains;

pub(crate) const PARENT_RULES: ParentRuleSet<Chunked> = ParentRuleSet::new(&[
ParentRuleSet::lift(&CastReduceAdaptor(Chunked)),
Expand Down Expand Up @@ -61,6 +62,10 @@ impl ArrayParentReduceRule<Chunked> for ChunkedUnaryScalarFnPushDownRule {
}

/// Push down non-unary scalar functions through chunked arrays where other siblings are constant.
///
/// [`ListContains`] is not pushed down. Its execution prepares a constant list once as a set, and
/// its chunked kernel then probes each chunk against that one set. A push-down would give each
/// chunk its own copy of the list, and so its own set to prepare.
#[derive(Debug)]
struct ChunkedConstantScalarFnPushDownRule;
impl ArrayParentReduceRule<Chunked> for ChunkedConstantScalarFnPushDownRule {
Expand All @@ -72,6 +77,10 @@ impl ArrayParentReduceRule<Chunked> for ChunkedConstantScalarFnPushDownRule {
parent: ArrayView<'_, ScalarFn>,
child_idx: usize,
) -> VortexResult<Option<ArrayRef>> {
if parent.scalar_fn().is::<ListContains>() {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is a hack, we essentially have common setup we would like to share across chunks, if we just dispatch over each chunk there's no place to store the shared state

return Ok(None);
}

for (idx, child) in parent.iter_children().enumerate() {
if idx == child_idx {
continue;
Expand Down
1 change: 1 addition & 0 deletions vortex-array/src/arrays/constant/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,3 +15,4 @@ pub(crate) mod compute;
mod vtable;

pub use vtable::Constant;
pub(crate) use vtable::canonical::list_scalar_elements;
Loading
Loading