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
6 changes: 3 additions & 3 deletions Cargo.lock

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

35 changes: 18 additions & 17 deletions encodings/parquet-variant/src/array.rs
Original file line number Diff line number Diff line change
Expand Up @@ -35,12 +35,9 @@ use vortex_array::scalar::Scalar;
use vortex_array::validity::Validity;
use vortex_array::vtable::child_to_validity;
use vortex_array::vtable::validity_to_child;
#[expect(
deprecated,
reason = "TODO(aduffy): figure out what to do with Parquet Variant"
)]
use vortex_arrow::ArrowArrayExecutor;
use vortex_arrow::ArrowExportOptions;
use vortex_arrow::ArrowSession;
use vortex_arrow::ArrowSessionExt;
use vortex_arrow::to_arrow_null_buffer;
use vortex_buffer::BitBuffer;
use vortex_error::VortexExpect;
Expand Down Expand Up @@ -436,19 +433,22 @@ pub trait ParquetVariantArrayExt:
}

/// Converts this storage array to Arrow's canonical Parquet Variant extension storage.
#[expect(
deprecated,
reason = "TODO(aduffy): figure out what to do with Parquet Variant"
)]
fn to_arrow(&self, ctx: &mut ExecutionCtx) -> VortexResult<ArrowVariantArray> {
fn to_arrow(
&self,
options: &ArrowExportOptions,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowVariantArray> {
let session = ctx.session().clone();
let arrow = session.arrow();
let exporter = arrow.exporter(options);
let metadata = self.metadata();
let len = metadata.len();
let nulls = to_arrow_null_buffer(self.parquet_variant_validity(), len, ctx)?;

let mut fields = Vec::with_capacity(3);
let mut arrays: Vec<ArrowArrayRef> = Vec::with_capacity(3);

let metadata_arrow = metadata.clone().execute_arrow(None, ctx)?;
let metadata_arrow = exporter.execute_arrow(metadata.clone(), None, ctx)?;
fields.push(Arc::new(Field::new(
"metadata",
metadata_arrow.data_type().clone(),
Expand All @@ -457,7 +457,7 @@ pub trait ParquetVariantArrayExt:
arrays.push(metadata_arrow);

if let Some(value) = self.value() {
let value_arrow = value.clone().execute_arrow(None, ctx)?;
let value_arrow = exporter.execute_arrow(value.clone(), None, ctx)?;
fields.push(Arc::new(Field::new(
"value",
value_arrow.data_type().clone(),
Expand All @@ -467,7 +467,7 @@ pub trait ParquetVariantArrayExt:
}

if let Some(typed_value) = self.typed_value() {
let tv_arrow = typed_value.clone().execute_arrow(None, ctx)?;
let tv_arrow = exporter.execute_arrow(typed_value.clone(), None, ctx)?;
fields.push(Arc::new(Field::new(
"typed_value",
tv_arrow.data_type().clone(),
Expand Down Expand Up @@ -514,6 +514,7 @@ mod tests {
use vortex_array::dtype::DType;
use vortex_array::dtype::Nullability;
use vortex_array::validity::Validity;
use vortex_arrow::ArrowExportOptions;
use vortex_arrow::ArrowSessionExt;
use vortex_buffer::buffer;
use vortex_error::VortexResult;
Expand All @@ -539,7 +540,7 @@ mod tests {
.ok_or_else(|| vortex_err!("expected parquet variant child"))?;

let mut ctx = SESSION.create_execution_ctx();
let roundtripped = inner.to_arrow(&mut ctx)?;
let roundtripped = inner.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
let roundtripped = roundtripped.inner();

assert_eq!(struct_array.len(), roundtripped.len());
Expand Down Expand Up @@ -647,7 +648,7 @@ mod tests {
let pv_array = ParquetVariant::try_new(Validity::NonNullable, metadata, Some(value), None)?;

let mut ctx = SESSION.create_execution_ctx();
let variant_arr = pv_array.to_arrow(&mut ctx)?;
let variant_arr = pv_array.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
let struct_arr = variant_arr.inner();

assert_eq!(struct_arr.num_columns(), 2);
Expand All @@ -669,7 +670,7 @@ mod tests {
)?;

let mut ctx = SESSION.create_execution_ctx();
let variant_arr = pv_array.to_arrow(&mut ctx)?;
let variant_arr = pv_array.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
let struct_arr = variant_arr.inner();

assert_eq!(struct_arr.num_columns(), 3);
Expand Down Expand Up @@ -756,7 +757,7 @@ mod tests {
assert!(parquet_array.typed_value().is_some());

let mut ctx = SESSION.create_execution_ctx();
let roundtripped = parquet_array.to_arrow(&mut ctx)?;
let roundtripped = parquet_array.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
let roundtripped = roundtripped.inner();
assert_eq!(
roundtripped.column_names(),
Expand Down
22 changes: 14 additions & 8 deletions encodings/parquet-variant/src/arrow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ use vortex_array::arrays::Variant;
use vortex_array::arrays::variant::VariantArraySlotsExt;
use vortex_array::dtype::DType;
use vortex_arrow::ArrowExport;
use vortex_arrow::ArrowExportOptions;
use vortex_arrow::ArrowExportVTable;
use vortex_arrow::ArrowImport;
use vortex_arrow::ArrowImportVTable;
Expand Down Expand Up @@ -63,8 +64,12 @@ fn parquet_variant_storage_request(fields: &Fields) -> Option<(bool, bool)> {
pub(crate) fn export_storage_to_target<T: ParquetVariantArrayExt>(
parquet_array: &T,
target_fields: &Fields,
options: &ArrowExportOptions,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowArrayRef> {
let session = ctx.session().clone();
let arrow = session.arrow();
let exporter = arrow.exporter(options);
let mut arrays = Vec::with_capacity(target_fields.len());

for field in target_fields {
Expand All @@ -82,11 +87,7 @@ pub(crate) fn export_storage_to_target<T: ParquetVariantArrayExt>(
);
};

arrays.push(ctx.session().clone().arrow().execute_arrow(
child,
Some(field.as_ref()),
ctx,
)?);
arrays.push(exporter.execute_arrow(child, Some(field.as_ref()), ctx)?);
}

let nulls = to_arrow_null_buffer(
Expand All @@ -104,17 +105,18 @@ pub(crate) fn export_storage_to_target<T: ParquetVariantArrayExt>(
pub(crate) fn export_unshredded_storage_to_target<T: ParquetVariantArrayExt>(
parquet_array: &T,
target_fields: &Fields,
options: &ArrowExportOptions,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowArrayRef> {
let arrow_variant = parquet_array.to_arrow(ctx)?;
let arrow_variant = parquet_array.to_arrow(options, ctx)?;
let unshredded = unshred_variant(&arrow_variant)?;
let unshredded_array = if parquet_array.as_ref().dtype().is_nullable() {
ParquetVariant::from_arrow_variant_nullable(&unshredded, &ctx.session().arrow())?
} else {
ParquetVariant::from_arrow_variant(&unshredded, &ctx.session().arrow())?
};
let unshredded_parquet = unshredded_array.as_::<ParquetVariant>();
export_storage_to_target(&unshredded_parquet, target_fields, ctx)
export_storage_to_target(&unshredded_parquet, target_fields, options, ctx)
}

pub(crate) fn parquet_variant_for_export(
Expand Down Expand Up @@ -178,6 +180,7 @@ impl ArrowExportVTable for ParquetVariant {
&self,
array: ArrayRef,
target: &Field,
options: &ArrowExportOptions,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowExport> {
if target
Expand All @@ -203,6 +206,7 @@ impl ArrowExportVTable for ParquetVariant {
return Ok(ArrowExport::Exported(export_unshredded_storage_to_target(
&parquet_array,
fields,
options,
ctx,
)?));
}
Expand All @@ -220,11 +224,13 @@ impl ArrowExportVTable for ParquetVariant {
return Ok(ArrowExport::Exported(export_storage_to_target(
&parquet_array,
fields,
options,
ctx,
)?));
}

let arrow_variant = Arc::new(parquet_array.to_arrow(ctx)?.into_inner()) as ArrowArrayRef;
let arrow_variant =
Arc::new(parquet_array.to_arrow(options, ctx)?.into_inner()) as ArrowArrayRef;

if arrow_variant.data_type() == target.data_type() {
Ok(ArrowExport::Exported(arrow_variant))
Expand Down
3 changes: 2 additions & 1 deletion encodings/parquet-variant/src/kernel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ use vortex_array::scalar_fn::ScalarFnVTable;
use vortex_array::scalar_fn::fns::variant_get::VariantGet;
use vortex_array::scalar_fn::fns::variant_get::VariantPath;
use vortex_array::scalar_fn::fns::variant_get::VariantPathElement;
use vortex_arrow::ArrowExportOptions;
use vortex_arrow::ArrowSession;
use vortex_arrow::ArrowSessionExt;
use vortex_error::VortexResult;
Expand Down Expand Up @@ -106,7 +107,7 @@ impl ExecuteParentKernel<ParquetVariant> for VariantGetKernel {
return Ok(None);
}

let arrow_variant = array.to_arrow(ctx)?;
let arrow_variant = array.to_arrow(&ArrowExportOptions::default(), ctx)?;
let arrow_input: ArrowArrayRef = Arc::new(arrow_variant.into_inner());
let session = ctx.session().clone();
let as_type = to_arrow_as_type(parent.options.dtype(), &session.arrow())?;
Expand Down
3 changes: 2 additions & 1 deletion encodings/parquet-variant/src/operations.rs
Original file line number Diff line number Diff line change
Expand Up @@ -396,6 +396,7 @@ mod tests {
use vortex_array::dtype::Nullability;
use vortex_array::scalar::Scalar;
use vortex_array::scalar::ScalarValue;
use vortex_arrow::ArrowExportOptions;
use vortex_arrow::ArrowSessionExt;
use vortex_error::VortexResult;
use vortex_session::VortexSession;
Expand Down Expand Up @@ -548,7 +549,7 @@ mod tests {

let inner_pv = vortex_arr.as_opt::<ParquetVariant>().unwrap();
let mut ctx = array_session().create_execution_ctx();
let roundtripped = inner_pv.to_arrow(&mut ctx)?;
let roundtripped = inner_pv.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
assert_eq!(roundtripped.inner().null_count(), 2);

Ok(())
Expand Down
18 changes: 13 additions & 5 deletions encodings/uuid/src/arrow.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ use vortex_array::dtype::extension::ExtDType;
use vortex_array::dtype::extension::ExtVTable;
use vortex_array::validity::Validity;
use vortex_arrow::ArrowExport;
use vortex_arrow::ArrowExportOptions;
use vortex_arrow::ArrowExportVTable;
use vortex_arrow::ArrowImport;
use vortex_arrow::ArrowImportVTable;
Expand Down Expand Up @@ -86,6 +87,7 @@ impl ArrowExportVTable for Uuid {
&self,
array: ArrayRef,
_target: &Field,
options: &ArrowExportOptions,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowExport> {
let is_uuid = array
Expand All @@ -96,7 +98,7 @@ impl ArrowExportVTable for Uuid {
if !is_uuid {
return Ok(ArrowExport::Unsupported(array));
}
Ok(ArrowExport::Exported(try_fsl_to_fsb(array, ctx)?))
Ok(ArrowExport::Exported(try_fsl_to_fsb(array, options, ctx)?))
}
}

Expand Down Expand Up @@ -165,7 +167,11 @@ impl ArrowImportVTable for Uuid {

/// Reinterpret a Vortex UUID extension array's `FixedSizeList<u8; 16>` storage as an Arrow
/// `FixedSizeBinary[16]` array, sharing the underlying byte buffer.
fn try_fsl_to_fsb(array: ArrayRef, ctx: &mut ExecutionCtx) -> VortexResult<ArrowArrayRef> {
fn try_fsl_to_fsb(
array: ArrayRef,
options: &ArrowExportOptions,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowArrayRef> {
let executed = array.execute::<ExtensionArray>(ctx)?;
let storage = executed.storage_array().clone();
let storage_arrow_type = DataType::FixedSizeList(
Expand All @@ -180,9 +186,11 @@ fn try_fsl_to_fsb(array: ArrayRef, ctx: &mut ExecutionCtx) -> VortexResult<Arrow
);

let session = ctx.session().clone();
let arrow_storage = session
.arrow()
.execute_arrow(storage, Some(&storage_field), ctx)?;
let arrow_storage =
session
.arrow()
.exporter(options)
.execute_arrow(storage, Some(&storage_field), ctx)?;

let fsl = arrow_storage.as_fixed_size_list();
let bytes = fsl
Expand Down
9 changes: 8 additions & 1 deletion vortex-arrow/src/executor/byte_view.rs
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ use vortex_array::dtype::Nullability;
use vortex_buffer::Buffer;
use vortex_error::VortexResult;

use crate::ArrowExporter;
use crate::CompactBuffers;
use crate::dtype::from_arrow_data_type;
use crate::null_buffer::to_null_buffer;

Expand Down Expand Up @@ -52,6 +54,7 @@ pub fn execute_varbinview_to_arrow<T: ByteViewType>(

pub(super) fn to_arrow_byte_view<T: ByteViewType>(
array: ArrayRef,
exporter: &ArrowExporter<'_>,
ctx: &mut ExecutionCtx,
) -> VortexResult<ArrowArrayRef> {
// First we cast the array into the desired ByteView type.
Expand All @@ -62,7 +65,11 @@ pub(super) fn to_arrow_byte_view<T: ByteViewType>(

let array = array.execute::<ArrayRef>(ctx)?;
let varbinview = array.execute::<VarBinViewArray>(ctx)?;
execute_varbinview_to_arrow::<T>(&varbinview, ctx)
if exporter.options().get_or_default::<CompactBuffers>().0 {
execute_varbinview_to_arrow::<T>(&varbinview, ctx)
} else {
canonical_varbinview_to_arrow::<T>(&varbinview, ctx)
}
}

#[cfg(test)]
Expand Down
Loading
Loading