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
37 changes: 22 additions & 15 deletions encodings/parquet-variant/src/array.rs
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ use vortex_array::vtable::validity_to_child;
reason = "TODO(aduffy): figure out what to do with Parquet Variant"
)]
use vortex_arrow::ArrowArrayExecutor;
use vortex_arrow::FromArrowArray;
use vortex_arrow::ArrowSession;
use vortex_arrow::to_arrow_null_buffer;
use vortex_buffer::BitBuffer;
use vortex_error::VortexExpect;
Expand Down Expand Up @@ -91,20 +91,26 @@ impl ParquetVariant {
)
}

/// Converts an Arrow `parquet_variant_compute::VariantArray` into Parquet Variant storage.
pub fn from_arrow_variant(arrow_variant: &ArrowVariantArray) -> VortexResult<ArrayRef> {
Self::from_arrow_variant_impl(arrow_variant, false)
/// Converts an Arrow `parquet_variant_compute::VariantArray` into Parquet Variant storage,
/// converting the storage children through `session`.
pub fn from_arrow_variant(
arrow_variant: &ArrowVariantArray,
session: &ArrowSession,
) -> VortexResult<ArrayRef> {
Self::from_arrow_variant_impl(arrow_variant, false, session)
}

pub(crate) fn from_arrow_variant_nullable(
arrow_variant: &ArrowVariantArray,
session: &ArrowSession,
) -> VortexResult<ArrayRef> {
Self::from_arrow_variant_impl(arrow_variant, true)
Self::from_arrow_variant_impl(arrow_variant, true, session)
}

fn from_arrow_variant_impl(
arrow_variant: &ArrowVariantArray,
force_nullable: bool,
session: &ArrowSession,
) -> VortexResult<ArrayRef> {
let storage = arrow_variant.inner();
let mut value_nullable = false;
Expand All @@ -130,17 +136,17 @@ impl ParquetVariant {
} else {
Validity::NonNullable
});
let metadata =
ArrayRef::from_arrow(arrow_variant.metadata_field() as &dyn ArrowArray, false)?;
let metadata = session
.from_arrow_array_nullable(arrow_variant.metadata_field() as &dyn ArrowArray, false)?;

let value = arrow_variant
.value_field()
.map(|v| ArrayRef::from_arrow(v as &dyn ArrowArray, value_nullable))
.map(|v| session.from_arrow_array_nullable(v as &dyn ArrowArray, value_nullable))
.transpose()?;

let typed_value = arrow_variant
.typed_value_field()
.map(|tv| ArrayRef::from_arrow(tv.as_ref(), typed_value_nullable))
.map(|tv| session.from_arrow_array_nullable(tv.as_ref(), typed_value_nullable))
.transpose()?;
ParquetVariant::try_new(validity, metadata, value, typed_value).map(IntoArray::into_array)
}
Expand Down Expand Up @@ -508,6 +514,7 @@ mod tests {
use vortex_array::dtype::DType;
use vortex_array::dtype::Nullability;
use vortex_array::validity::Validity;
use vortex_arrow::ArrowSessionExt;
use vortex_buffer::buffer;
use vortex_error::VortexResult;
use vortex_error::vortex_err;
Expand All @@ -526,7 +533,7 @@ mod tests {

fn assert_arrow_variant_storage_roundtrip(struct_array: StructArray) -> VortexResult<()> {
let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;
let inner = vortex_arr
.as_opt::<ParquetVariant>()
.ok_or_else(|| vortex_err!("expected parquet variant child"))?;
Expand Down Expand Up @@ -577,7 +584,7 @@ mod tests {
builder.append_variant(PqVariant::from(true));
let arrow_variant = builder.build();

let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;

assert_eq!(vortex_arr.len(), 3);
assert_eq!(
Expand Down Expand Up @@ -609,7 +616,7 @@ mod tests {

let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;

let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;
assert_eq!(vortex_arr.len(), 3);
assert_eq!(
vortex_arr.dtype(),
Expand Down Expand Up @@ -700,7 +707,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;
let parquet_array = vortex_arr
.as_opt::<ParquetVariant>()
.ok_or_else(|| vortex_err!("expected parquet variant array"))?;
Expand Down Expand Up @@ -737,7 +744,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;
let parquet_array = vortex_arr
.as_opt::<ParquetVariant>()
.ok_or_else(|| vortex_err!("expected parquet variant array"))?;
Expand Down Expand Up @@ -809,7 +816,7 @@ mod tests {
.with_path("a", &DataType::Int32)?
.build();
let shredded = shred_variant(&json_to_variant(&json)?, &shredding)?;
let original = ParquetVariant::from_arrow_variant(&shredded)?;
let original = ParquetVariant::from_arrow_variant(&shredded, &SESSION.arrow())?;
assert!(
original
.as_opt::<ParquetVariant>()
Expand Down
9 changes: 5 additions & 4 deletions encodings/parquet-variant/src/arrow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -109,9 +109,9 @@ pub(crate) fn export_unshredded_storage_to_target<T: ParquetVariantArrayExt>(
let arrow_variant = parquet_array.to_arrow(ctx)?;
let unshredded = unshred_variant(&arrow_variant)?;
let unshredded_array = if parquet_array.as_ref().dtype().is_nullable() {
ParquetVariant::from_arrow_variant_nullable(&unshredded)?
ParquetVariant::from_arrow_variant_nullable(&unshredded, &ctx.session().arrow())?
} else {
ParquetVariant::from_arrow_variant(&unshredded)?
ParquetVariant::from_arrow_variant(&unshredded, &ctx.session().arrow())?
};
let unshredded_parquet = unshredded_array.as_::<ParquetVariant>();
export_storage_to_target(&unshredded_parquet, target_fields, ctx)
Expand Down Expand Up @@ -263,6 +263,7 @@ impl ArrowImportVTable for ParquetVariant {
array: ArrowArrayRef,
field: &Field,
dtype: &DType,
session: &ArrowSession,
) -> VortexResult<ArrowImport> {
if !dtype.is_variant()
|| field
Expand All @@ -276,9 +277,9 @@ impl ArrowImportVTable for ParquetVariant {

let arrow_variant = ArrowVariantArray::try_new(array.as_struct())?;
let imported = if dtype.is_nullable() {
ParquetVariant::from_arrow_variant_nullable(&arrow_variant)?
ParquetVariant::from_arrow_variant_nullable(&arrow_variant, session)?
} else {
ParquetVariant::from_arrow_variant(&arrow_variant)?
ParquetVariant::from_arrow_variant(&arrow_variant, session)?
};
Ok(ArrowImport::Imported(imported.into_array()))
}
Expand Down
43 changes: 28 additions & 15 deletions encodings/parquet-variant/src/kernel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,6 @@ use vortex_array::scalar_fn::fns::variant_get::VariantPath;
use vortex_array::scalar_fn::fns::variant_get::VariantPathElement;
use vortex_arrow::ArrowSession;
use vortex_arrow::ArrowSessionExt;
use vortex_arrow::FromArrowArray;
use vortex_error::VortexResult;
use vortex_error::vortex_ensure_eq;
use vortex_error::vortex_err;
Expand Down Expand Up @@ -120,9 +119,14 @@ impl ExecuteParentKernel<ParquetVariant> for VariantGetKernel {
let arrow_output = arrow_variant_get(&arrow_input, get_options)?;
let output = if parent.options.dtype().is_none_or(DType::is_variant) {
let arrow_variant_output = ArrowVariantArray::try_new(arrow_output.as_ref())?;
ParquetVariant::from_arrow_variant_nullable(&arrow_variant_output)?
ParquetVariant::from_arrow_variant_nullable(
&arrow_variant_output,
&ctx.session().arrow(),
)?
} else {
ArrayRef::from_arrow(arrow_output.as_ref(), true)?
ctx.session()
.arrow()
.from_arrow_array_nullable(arrow_output.as_ref(), true)?
};

vortex_ensure_eq!(
Expand Down Expand Up @@ -164,9 +168,9 @@ fn json_strings_to_variant(
};

if nullable {
ParquetVariant::from_arrow_variant_nullable(&arrow_variant)
ParquetVariant::from_arrow_variant_nullable(&arrow_variant, &session.arrow())
} else {
ParquetVariant::from_arrow_variant(&arrow_variant)
ParquetVariant::from_arrow_variant(&arrow_variant, &session.arrow())
}
}

Expand Down Expand Up @@ -320,7 +324,7 @@ mod tests {
use vortex_array::scalar_fn::fns::variant_get::VariantPath;
use vortex_array::scalar_fn::fns::variant_get::VariantPathElement;
use vortex_array::validity::Validity;
use vortex_arrow::FromArrowArray;
use vortex_arrow::ArrowSessionExt;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
use vortex_error::vortex_ensure;
Expand All @@ -343,7 +347,7 @@ mod tests {
builder.append_variant(PqVariant::from("hello"));
builder.append_variant(PqVariant::from(true));
builder.append_variant(PqVariant::from(99i64));
ParquetVariant::from_arrow_variant(&builder.build())
ParquetVariant::from_arrow_variant(&builder.build(), &SESSION.arrow())
}

fn make_nullable_array() -> VortexResult<ArrayRef> {
Expand All @@ -360,13 +364,13 @@ mod tests {
Some(NullBuffer::from(vec![true, false, true, false])),
)?;
let arrow_variant = ArrowVariantArray::try_new(&null_struct)?;
ParquetVariant::from_arrow_variant(&arrow_variant)
ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())
}

fn make_unshredded_json_array(values: Vec<Option<&str>>) -> VortexResult<ArrayRef> {
let json: ArrowArrayRef = Arc::new(StringArray::from(values));
let arrow_variant = json_to_variant(&json)?;
ParquetVariant::from_arrow_variant(&arrow_variant)
ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())
}

fn parse_path(path: &str) -> VortexResult<VariantPath> {
Expand Down Expand Up @@ -766,15 +770,24 @@ mod tests {
.map(|field| field.is_nullable())
.unwrap_or(false);

let metadata =
ArrayRef::from_arrow(arrow_variant.metadata_field() as &dyn ArrowArray, false)?;
let metadata = SESSION
.arrow()
.from_arrow_array_nullable(arrow_variant.metadata_field() as &dyn ArrowArray, false)?;
let value = arrow_variant
.value_field()
.map(|value| ArrayRef::from_arrow(value as &dyn ArrowArray, value_nullable))
.map(|value| {
SESSION
.arrow()
.from_arrow_array_nullable(value as &dyn ArrowArray, value_nullable)
})
.transpose()?;
let typed_value = arrow_variant
.typed_value_field()
.map(|typed_value| ArrayRef::from_arrow(typed_value.as_ref(), typed_value_nullable))
.map(|typed_value| {
SESSION
.arrow()
.from_arrow_array_nullable(typed_value.as_ref(), typed_value_nullable)
})
.transpose()?;

Ok(
Expand All @@ -785,7 +798,7 @@ mod tests {

fn make_partially_shredded_object_array() -> VortexResult<ArrayRef> {
let arrow_variant = make_partially_shredded_arrow_variant()?;
let parquet_array = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let parquet_array = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;
let mut ctx = SESSION.create_execution_ctx();
let Canonical::Variant(canonical) = parquet_array.execute::<Canonical>(&mut ctx)? else {
return Err(vortex_err!("expected canonical variant array"));
Expand Down Expand Up @@ -902,7 +915,7 @@ mod tests {
None,
)?;
let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;
ParquetVariant::from_arrow_variant(&arrow_variant)
ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())
}

fn assert_typed_value_i32(
Expand Down
23 changes: 16 additions & 7 deletions encodings/parquet-variant/src/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,7 @@ fn parquet_variant_to_scalar(variant: PqVariant<'_, '_>) -> VortexResult<Scalar>
#[cfg(test)]
mod tests {
use std::sync::Arc;
use std::sync::LazyLock;

use arrow_array::Array as _;
use arrow_array::ArrayRef as ArrowArrayRef;
Expand All @@ -393,12 +394,20 @@ mod tests {
use vortex_array::dtype::Nullability;
use vortex_array::scalar::Scalar;
use vortex_array::scalar::ScalarValue;
use vortex_arrow::ArrowSessionExt;
use vortex_error::VortexResult;
use vortex_session::VortexSession;

use crate::ParquetVariant;
use crate::ParquetVariantArrayExt;
use crate::operations::parquet_variant_to_scalar;

static SESSION: LazyLock<VortexSession> = LazyLock::new(|| {
let session = array_session();
crate::initialize(&session);
session
});

fn binary_view_array(values: &[&[u8]]) -> ArrowArrayRef {
let mut builder = BinaryViewBuilder::new();
for value in values {
Expand All @@ -411,7 +420,7 @@ mod tests {
arrow_variant: &ArrowVariantArray,
rows: impl IntoIterator<Item = usize>,
) -> VortexResult<()> {
let vortex_arr = ParquetVariant::from_arrow_variant(arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(arrow_variant, &SESSION.arrow())?;

for index in rows {
let expected_inner = parquet_variant_to_scalar(arrow_variant.try_value(index)?)?;
Expand Down Expand Up @@ -443,7 +452,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&null_struct)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;

assert_eq!(vortex_arr.dtype(), &DType::Variant(Nullability::Nullable));

Expand Down Expand Up @@ -484,7 +493,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&null_struct)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;

let present_variant_null =
vortex_arr.execute_scalar(0, &mut array_session().create_execution_ctx())?;
Expand Down Expand Up @@ -521,7 +530,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&null_struct)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;

assert_eq!(vortex_arr.dtype(), &DType::Variant(Nullability::Nullable));
assert!(
Expand Down Expand Up @@ -550,7 +559,7 @@ mod tests {
builder.append_variant(PqVariant::from(2i32));
let arrow_variant = builder.build();

let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;

assert_eq!(
vortex_arr.dtype(),
Expand Down Expand Up @@ -671,7 +680,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;

let row0 = vortex_arr.execute_scalar(0, &mut array_session().create_execution_ctx())?;
let row0 = row0.as_variant().value().unwrap().as_list();
Expand Down Expand Up @@ -763,7 +772,7 @@ mod tests {
)?;

let arrow_variant = ArrowVariantArray::try_new(&struct_array)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant)?;
let vortex_arr = ParquetVariant::from_arrow_variant(&arrow_variant, &SESSION.arrow())?;
let object = vortex_arr.execute_scalar(0, &mut array_session().create_execution_ctx())?;
let object = object.as_variant().value().unwrap().as_struct();

Expand Down
6 changes: 5 additions & 1 deletion encodings/parquet-variant/src/vtable.rs
Original file line number Diff line number Diff line change
Expand Up @@ -322,6 +322,7 @@ mod tests {
use vortex_array::session::ArraySessionExt;
use vortex_array::stream::ArrayStreamExt;
use vortex_array::validity::Validity;
use vortex_arrow::ArrowSessionExt;
use vortex_buffer::BitBuffer;
use vortex_buffer::ByteBufferMut;
use vortex_buffer::buffer;
Expand Down Expand Up @@ -387,7 +388,10 @@ mod tests {
None,
)?;

ParquetVariant::from_arrow_variant(&ArrowVariantArray::try_new(&arrow_storage)?)
ParquetVariant::from_arrow_variant(
&ArrowVariantArray::try_new(&arrow_storage)?,
&SESSION.arrow(),
)
}

#[fixture]
Expand Down
2 changes: 1 addition & 1 deletion lang/cpp/include/vortex/array.hpp
Original file line number Diff line number Diff line change
Expand Up @@ -135,7 +135,7 @@ class Array {
* Import an Arrow array. Consumes both "array" and "schema", do not use
* or release them afterwards. For a record batch pass nullable = false.
*/
static Array from_arrow(ArrowArray *array, ArrowSchema *schema, bool nullable);
static Array from_arrow(const Session &session, ArrowArray *array, ArrowSchema *schema, bool nullable);

size_t size() const;
bool nullable() const;
Expand Down
Loading
Loading