Skip to content

Commit d14fa26

Browse files
break(arrow): add export options and allow skipping view buffer compaction (#10220)
Co-authored-by: Connor Tsui <connor.tsui20@gmail.com> Signed-off-by: Lorenz Hübschle <lorenz@firebolt.io>
1 parent dc0e885 commit d14fa26

28 files changed

Lines changed: 782 additions & 187 deletions

File tree

‎encodings/parquet-variant/src/array.rs‎

Lines changed: 18 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -35,12 +35,9 @@ use vortex_array::scalar::Scalar;
3535
use vortex_array::validity::Validity;
3636
use vortex_array::vtable::child_to_validity;
3737
use vortex_array::vtable::validity_to_child;
38-
#[expect(
39-
deprecated,
40-
reason = "TODO(aduffy): figure out what to do with Parquet Variant"
41-
)]
42-
use vortex_arrow::ArrowArrayExecutor;
38+
use vortex_arrow::ArrowExportOptions;
4339
use vortex_arrow::ArrowSession;
40+
use vortex_arrow::ArrowSessionExt;
4441
use vortex_arrow::to_arrow_null_buffer;
4542
use vortex_buffer::BitBuffer;
4643
use vortex_error::VortexExpect;
@@ -440,19 +437,22 @@ pub trait ParquetVariantArrayExt:
440437
}
441438

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

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

455-
let metadata_arrow = metadata.clone().execute_arrow(None, ctx)?;
455+
let metadata_arrow = exporter.execute_arrow(metadata.clone(), None, ctx)?;
456456
fields.push(Arc::new(Field::new(
457457
"metadata",
458458
metadata_arrow.data_type().clone(),
@@ -461,7 +461,7 @@ pub trait ParquetVariantArrayExt:
461461
arrays.push(metadata_arrow);
462462

463463
if let Some(value) = self.value() {
464-
let value_arrow = value.clone().execute_arrow(None, ctx)?;
464+
let value_arrow = exporter.execute_arrow(value.clone(), None, ctx)?;
465465
fields.push(Arc::new(Field::new(
466466
"value",
467467
value_arrow.data_type().clone(),
@@ -471,7 +471,7 @@ pub trait ParquetVariantArrayExt:
471471
}
472472

473473
if let Some(typed_value) = self.typed_value() {
474-
let tv_arrow = typed_value.clone().execute_arrow(None, ctx)?;
474+
let tv_arrow = exporter.execute_arrow(typed_value.clone(), None, ctx)?;
475475
fields.push(Arc::new(Field::new(
476476
"typed_value",
477477
tv_arrow.data_type().clone(),
@@ -518,6 +518,7 @@ mod tests {
518518
use vortex_array::dtype::DType;
519519
use vortex_array::dtype::Nullability;
520520
use vortex_array::validity::Validity;
521+
use vortex_arrow::ArrowExportOptions;
521522
use vortex_arrow::ArrowSessionExt;
522523
use vortex_buffer::buffer;
523524
use vortex_error::VortexResult;
@@ -543,7 +544,7 @@ mod tests {
543544
.ok_or_else(|| vortex_err!("expected parquet variant child"))?;
544545

545546
let mut ctx = SESSION.create_execution_ctx();
546-
let roundtripped = inner.to_arrow(&mut ctx)?;
547+
let roundtripped = inner.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
547548
let roundtripped = roundtripped.inner();
548549

549550
assert_eq!(struct_array.len(), roundtripped.len());
@@ -651,7 +652,7 @@ mod tests {
651652
let pv_array = ParquetVariant::try_new(Validity::NonNullable, metadata, Some(value), None)?;
652653

653654
let mut ctx = SESSION.create_execution_ctx();
654-
let variant_arr = pv_array.to_arrow(&mut ctx)?;
655+
let variant_arr = pv_array.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
655656
let struct_arr = variant_arr.inner();
656657

657658
assert_eq!(struct_arr.num_columns(), 2);
@@ -673,7 +674,7 @@ mod tests {
673674
)?;
674675

675676
let mut ctx = SESSION.create_execution_ctx();
676-
let variant_arr = pv_array.to_arrow(&mut ctx)?;
677+
let variant_arr = pv_array.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
677678
let struct_arr = variant_arr.inner();
678679

679680
assert_eq!(struct_arr.num_columns(), 3);
@@ -760,7 +761,7 @@ mod tests {
760761
assert!(parquet_array.typed_value().is_some());
761762

762763
let mut ctx = SESSION.create_execution_ctx();
763-
let roundtripped = parquet_array.to_arrow(&mut ctx)?;
764+
let roundtripped = parquet_array.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
764765
let roundtripped = roundtripped.inner();
765766
assert_eq!(
766767
roundtripped.column_names(),

‎encodings/parquet-variant/src/arrow.rs‎

Lines changed: 14 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,7 @@ use vortex_array::arrays::Variant;
2222
use vortex_array::arrays::variant::VariantArraySlotsExt;
2323
use vortex_array::dtype::DType;
2424
use vortex_arrow::ArrowExport;
25+
use vortex_arrow::ArrowExportOptions;
2526
use vortex_arrow::ArrowExportVTable;
2627
use vortex_arrow::ArrowImport;
2728
use vortex_arrow::ArrowImportVTable;
@@ -63,8 +64,12 @@ fn parquet_variant_storage_request(fields: &Fields) -> Option<(bool, bool)> {
6364
pub(crate) fn export_storage_to_target<T: ParquetVariantArrayExt>(
6465
parquet_array: &T,
6566
target_fields: &Fields,
67+
options: &ArrowExportOptions,
6668
ctx: &mut ExecutionCtx,
6769
) -> VortexResult<ArrowArrayRef> {
70+
let session = ctx.session().clone();
71+
let arrow = session.arrow();
72+
let exporter = arrow.exporter(options);
6873
let mut arrays = Vec::with_capacity(target_fields.len());
6974

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

85-
arrays.push(ctx.session().clone().arrow().execute_arrow(
86-
child,
87-
Some(field.as_ref()),
88-
ctx,
89-
)?);
90+
arrays.push(exporter.execute_arrow(child, Some(field.as_ref()), ctx)?);
9091
}
9192

9293
let nulls = to_arrow_null_buffer(
@@ -104,17 +105,18 @@ pub(crate) fn export_storage_to_target<T: ParquetVariantArrayExt>(
104105
pub(crate) fn export_unshredded_storage_to_target<T: ParquetVariantArrayExt>(
105106
parquet_array: &T,
106107
target_fields: &Fields,
108+
options: &ArrowExportOptions,
107109
ctx: &mut ExecutionCtx,
108110
) -> VortexResult<ArrowArrayRef> {
109-
let arrow_variant = parquet_array.to_arrow(ctx)?;
111+
let arrow_variant = parquet_array.to_arrow(options, ctx)?;
110112
let unshredded = unshred_variant(&arrow_variant)?;
111113
let unshredded_array = if parquet_array.as_ref().dtype().is_nullable() {
112114
ParquetVariant::from_arrow_variant_nullable(&unshredded, &ctx.session().arrow())?
113115
} else {
114116
ParquetVariant::from_arrow_variant(&unshredded, &ctx.session().arrow())?
115117
};
116118
let unshredded_parquet = unshredded_array.as_::<ParquetVariant>();
117-
export_storage_to_target(&unshredded_parquet, target_fields, ctx)
119+
export_storage_to_target(&unshredded_parquet, target_fields, options, ctx)
118120
}
119121

120122
pub(crate) fn parquet_variant_for_export(
@@ -178,6 +180,7 @@ impl ArrowExportVTable for ParquetVariant {
178180
&self,
179181
array: ArrayRef,
180182
target: &Field,
183+
options: &ArrowExportOptions,
181184
ctx: &mut ExecutionCtx,
182185
) -> VortexResult<ArrowExport> {
183186
if target
@@ -203,6 +206,7 @@ impl ArrowExportVTable for ParquetVariant {
203206
return Ok(ArrowExport::Exported(export_unshredded_storage_to_target(
204207
&parquet_array,
205208
fields,
209+
options,
206210
ctx,
207211
)?));
208212
}
@@ -220,11 +224,13 @@ impl ArrowExportVTable for ParquetVariant {
220224
return Ok(ArrowExport::Exported(export_storage_to_target(
221225
&parquet_array,
222226
fields,
227+
options,
223228
ctx,
224229
)?));
225230
}
226231

227-
let arrow_variant = Arc::new(parquet_array.to_arrow(ctx)?.into_inner()) as ArrowArrayRef;
232+
let arrow_variant =
233+
Arc::new(parquet_array.to_arrow(options, ctx)?.into_inner()) as ArrowArrayRef;
228234

229235
if arrow_variant.data_type() == target.data_type() {
230236
Ok(ArrowExport::Exported(arrow_variant))

‎encodings/parquet-variant/src/kernel.rs‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,6 +43,7 @@ use vortex_array::scalar_fn::ScalarFnVTable;
4343
use vortex_array::scalar_fn::fns::variant_get::VariantGet;
4444
use vortex_array::scalar_fn::fns::variant_get::VariantPath;
4545
use vortex_array::scalar_fn::fns::variant_get::VariantPathElement;
46+
use vortex_arrow::ArrowExportOptions;
4647
use vortex_arrow::ArrowSession;
4748
use vortex_arrow::ArrowSessionExt;
4849
use vortex_error::VortexResult;
@@ -106,7 +107,7 @@ impl ExecuteParentKernel<ParquetVariant> for VariantGetKernel {
106107
return Ok(None);
107108
}
108109

109-
let arrow_variant = array.to_arrow(ctx)?;
110+
let arrow_variant = array.to_arrow(&ArrowExportOptions::default(), ctx)?;
110111
let arrow_input: ArrowArrayRef = Arc::new(arrow_variant.into_inner());
111112
let session = ctx.session().clone();
112113
let as_type = to_arrow_as_type(parent.options.dtype(), &session.arrow())?;

‎encodings/parquet-variant/src/operations.rs‎

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -396,6 +396,7 @@ mod tests {
396396
use vortex_array::dtype::Nullability;
397397
use vortex_array::scalar::Scalar;
398398
use vortex_array::scalar::ScalarValue;
399+
use vortex_arrow::ArrowExportOptions;
399400
use vortex_arrow::ArrowSessionExt;
400401
use vortex_error::VortexResult;
401402
use vortex_session::VortexSession;
@@ -548,7 +549,7 @@ mod tests {
548549

549550
let inner_pv = vortex_arr.as_opt::<ParquetVariant>().unwrap();
550551
let mut ctx = array_session().create_execution_ctx();
551-
let roundtripped = inner_pv.to_arrow(&mut ctx)?;
552+
let roundtripped = inner_pv.to_arrow(&ArrowExportOptions::default(), &mut ctx)?;
552553
assert_eq!(roundtripped.inner().null_count(), 2);
553554

554555
Ok(())

‎encodings/uuid/src/arrow.rs‎

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ use vortex_array::dtype::extension::ExtDType;
3333
use vortex_array::dtype::extension::ExtVTable;
3434
use vortex_array::validity::Validity;
3535
use vortex_arrow::ArrowExport;
36+
use vortex_arrow::ArrowExportOptions;
3637
use vortex_arrow::ArrowExportVTable;
3738
use vortex_arrow::ArrowImport;
3839
use vortex_arrow::ArrowImportVTable;
@@ -86,6 +87,7 @@ impl ArrowExportVTable for Uuid {
8687
&self,
8788
array: ArrayRef,
8889
_target: &Field,
90+
options: &ArrowExportOptions,
8991
ctx: &mut ExecutionCtx,
9092
) -> VortexResult<ArrowExport> {
9193
let is_uuid = array
@@ -96,7 +98,7 @@ impl ArrowExportVTable for Uuid {
9698
if !is_uuid {
9799
return Ok(ArrowExport::Unsupported(array));
98100
}
99-
Ok(ArrowExport::Exported(try_fsl_to_fsb(array, ctx)?))
101+
Ok(ArrowExport::Exported(try_fsl_to_fsb(array, options, ctx)?))
100102
}
101103
}
102104

@@ -165,7 +167,11 @@ impl ArrowImportVTable for Uuid {
165167

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

182188
let session = ctx.session().clone();
183-
let arrow_storage = session
184-
.arrow()
185-
.execute_arrow(storage, Some(&storage_field), ctx)?;
189+
let arrow_storage =
190+
session
191+
.arrow()
192+
.exporter(options)
193+
.execute_arrow(storage, Some(&storage_field), ctx)?;
186194

187195
let fsl = arrow_storage.as_fixed_size_list();
188196
let bytes = fsl

‎vortex-arrow/src/executor/byte_view.rs‎

Lines changed: 8 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,8 @@ use vortex_array::dtype::Nullability;
1414
use vortex_buffer::Buffer;
1515
use vortex_error::VortexResult;
1616

17+
use crate::ArrowExporter;
18+
use crate::CompactBuffers;
1719
use crate::dtype::from_arrow_data_type;
1820
use crate::null_buffer::to_null_buffer;
1921

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

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

6366
let array = array.execute::<ArrayRef>(ctx)?;
6467
let varbinview = array.execute::<VarBinViewArray>(ctx)?;
65-
execute_varbinview_to_arrow::<T>(&varbinview, ctx)
68+
if exporter.options().get_or_default::<CompactBuffers>().0 {
69+
execute_varbinview_to_arrow::<T>(&varbinview, ctx)
70+
} else {
71+
canonical_varbinview_to_arrow::<T>(&varbinview, ctx)
72+
}
6673
}
6774

6875
#[cfg(test)]

‎vortex-arrow/src/executor/dictionary.rs‎

Lines changed: 16 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -23,7 +23,8 @@ use vortex_error::VortexError;
2323
use vortex_error::VortexResult;
2424
use vortex_error::vortex_bail;
2525

26-
use crate::ArrowArrayExecutor;
26+
use crate::ArrowExporter;
27+
use crate::executor::execute_arrow_naive;
2728

2829
/// Matches the encodings [`to_arrow_dictionary`] requires for export.
2930
struct ArrowDictExportable;
@@ -40,22 +41,23 @@ pub(super) fn to_arrow_dictionary(
4041
array: ArrayRef,
4142
codes_type: &DataType,
4243
values_type: &DataType,
44+
exporter: &ArrowExporter<'_>,
4345
ctx: &mut ExecutionCtx,
4446
) -> VortexResult<ArrowArrayRef> {
4547
let array = array.execute_until::<ArrowDictExportable>(ctx)?;
4648

4749
let array = match array.try_downcast::<Dict>() {
48-
Ok(dict) => return dict_to_dict(dict, codes_type, values_type, ctx),
50+
Ok(dict) => return dict_to_dict(dict, codes_type, values_type, exporter, ctx),
4951
Err(array) => array,
5052
};
5153
let array = match array.try_downcast::<Constant>() {
52-
Ok(constant) => return constant_to_dict(constant, codes_type, values_type, ctx),
54+
Ok(constant) => return constant_to_dict(constant, codes_type, values_type, exporter, ctx),
5355
Err(array) => array,
5456
};
5557

5658
// Otherwise, we should try and build a dictionary.
5759
// Arrow hides this functionality inside the cast module!
58-
let array = array.execute_arrow(Some(values_type), ctx)?;
60+
let array = execute_arrow_naive(array, Some(values_type), exporter, ctx)?;
5961
arrow_cast::cast(
6062
&array,
6163
&DataType::Dictionary(Box::new(codes_type.clone()), Box::new(values_type.clone())),
@@ -68,6 +70,7 @@ fn constant_to_dict(
6870
array: ConstantArray,
6971
codes_type: &DataType,
7072
values_type: &DataType,
73+
exporter: &ArrowExporter<'_>,
7174
ctx: &mut ExecutionCtx,
7275
) -> VortexResult<ArrowArrayRef> {
7376
let len = array.len();
@@ -78,9 +81,12 @@ fn constant_to_dict(
7881
return Ok(new_null_array(&dict_type, len));
7982
}
8083

81-
let values = ConstantArray::new(scalar.clone(), 1)
82-
.into_array()
83-
.execute_arrow(Some(values_type), ctx)?;
84+
let values = execute_arrow_naive(
85+
ConstantArray::new(scalar.clone(), 1).into_array(),
86+
Some(values_type),
87+
exporter,
88+
ctx,
89+
)?;
8490
let codes = zeroed_codes_array(codes_type, len)?;
8591
make_dict_array(codes_type, codes, values)
8692
}
@@ -90,13 +96,11 @@ fn dict_to_dict(
9096
array: DictArray,
9197
codes_type: &DataType,
9298
values_type: &DataType,
99+
exporter: &ArrowExporter<'_>,
93100
ctx: &mut ExecutionCtx,
94101
) -> VortexResult<ArrowArrayRef> {
95-
let codes = array.codes().clone().execute_arrow(Some(codes_type), ctx)?;
96-
let values = array
97-
.values()
98-
.clone()
99-
.execute_arrow(Some(values_type), ctx)?;
102+
let codes = execute_arrow_naive(array.codes().clone(), Some(codes_type), exporter, ctx)?;
103+
let values = execute_arrow_naive(array.values().clone(), Some(values_type), exporter, ctx)?;
100104
make_dict_array(codes_type, codes, values)
101105
}
102106

0 commit comments

Comments
 (0)