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
12 changes: 8 additions & 4 deletions benchmarks/string-bench/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ use vortex::array::arrays::VarBinViewArray;
use vortex::array::arrays::struct_::StructArrayExt;
use vortex::io::session::RuntimeSessionExt;
use vortex::session::VortexSession;
use vortex_arrow::FromArrowArray;
use vortex_arrow::ArrowSessionExt;
use vortex_bench::Format;
use vortex_bench::IdempotentPath;
use vortex_bench::datasets::Dataset;
Expand Down Expand Up @@ -299,9 +299,13 @@ async fn read_parquet_projected(path: PathBuf, column: &str) -> Result<ArrayRef>

let chunks: Vec<ArrayRef> = reader
.map(|batch| {
batch
.map_err(anyhow::Error::from)
.and_then(|rb| ArrayRef::from_arrow(rb, false).map_err(anyhow::Error::from))
batch.map_err(anyhow::Error::from).and_then(|rb| {
let schema = rb.schema();
SESSION
.arrow()
.from_arrow_record_batch(rb, &schema)
.map_err(anyhow::Error::from)
})
})
.try_collect()
.await?;
Expand Down
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(ArrowArrayRef::clone(arrow_variant.metadata_field()), false)?;

let value = arrow_variant
.value_field()
.map(|v| ArrayRef::from_arrow(v as &dyn ArrowArray, value_nullable))
.map(|v| session.from_arrow_array(ArrowArrayRef::clone(v), 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(ArrowArrayRef::clone(tv), 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
37 changes: 24 additions & 13 deletions encodings/parquet-variant/src/kernel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,7 +117,10 @@ 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 {
// Import through the same `as_type` field the cast targeted, so an extension target
// dtype comes back as that extension rather than as its bare storage type.
Expand Down Expand Up @@ -166,9 +169,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 @@ -327,7 +330,6 @@ mod tests {
use vortex_array::scalar_fn::fns::variant_get::VariantPathElement;
use vortex_array::validity::Validity;
use vortex_arrow::ArrowSessionExt;
use vortex_arrow::FromArrowArray;
use vortex_error::VortexResult;
use vortex_error::vortex_bail;
use vortex_error::vortex_ensure;
Expand Down Expand Up @@ -414,7 +416,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 @@ -431,13 +433,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 @@ -837,15 +839,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(ArrowArrayRef::clone(arrow_variant.metadata_field()), 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(ArrowArrayRef::clone(value), 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(ArrowArrayRef::clone(typed_value), typed_value_nullable)
})
.transpose()?;

Ok(
Expand All @@ -856,7 +867,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 @@ -973,7 +984,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
Loading
Loading