1use std::{collections::HashMap, sync::Arc};
5
6use arrow_array::{
7 Array, ArrayRef, StructArray, UInt64Array,
8 builder::{PrimitiveBuilder, StringBuilder},
9 cast::AsArray,
10 types::{UInt8Type, UInt32Type, UInt64Type},
11};
12use arrow_buffer::Buffer;
13use arrow_schema::{DataType, Field as ArrowField, Fields};
14use futures::{FutureExt, future::BoxFuture};
15use lance_core::{
16 Error, Result, datatypes::BLOB_V2_DESC_FIELDS, datatypes::Field, error::LanceOptionExt,
17};
18
19use crate::{
20 buffer::LanceBuffer,
21 constants::PACKED_STRUCT_META_KEY,
22 decoder::PageEncoding,
23 encoder::{EncodeTask, EncodedColumn, EncodedPage, FieldEncoder, OutOfLineBuffers},
24 format::ProtobufUtils21,
25 repdef::{DefinitionInterpretation, RepDefBuilder},
26};
27use lance_core::datatypes::BlobKind;
28
29pub struct BlobStructuralEncoder {
35 descriptor_encoder: Box<dyn FieldEncoder>,
37 def_meaning: Option<Arc<[DefinitionInterpretation]>>,
39}
40
41impl BlobStructuralEncoder {
42 pub fn new(
43 field: &Field,
44 make_descriptor_encoder: impl FnOnce(Field) -> Result<Box<dyn FieldEncoder>>,
45 ) -> Result<Self> {
46 let mut descriptor_metadata = HashMap::with_capacity(1);
49 descriptor_metadata.insert(PACKED_STRUCT_META_KEY.to_string(), "true".to_string());
50
51 let descriptor_data_type = DataType::Struct(Fields::from(vec![
52 ArrowField::new("position", DataType::UInt64, false),
53 ArrowField::new("size", DataType::UInt64, false),
54 ]));
55
56 let descriptor_field = Field::try_from(
58 ArrowField::new(&field.name, descriptor_data_type, field.nullable)
59 .with_metadata(descriptor_metadata),
60 )?;
61
62 let descriptor_encoder = make_descriptor_encoder(descriptor_field)?;
64
65 Ok(Self {
66 descriptor_encoder,
67 def_meaning: None,
68 })
69 }
70
71 fn wrap_tasks(
72 tasks: Vec<EncodeTask>,
73 def_meaning: Arc<[DefinitionInterpretation]>,
74 ) -> Vec<EncodeTask> {
75 tasks
76 .into_iter()
77 .map(|task| {
78 let def_meaning = def_meaning.clone();
79 task.then(|encoded_page| async move {
80 let encoded_page = encoded_page?;
81
82 let PageEncoding::Structural(inner_layout) = encoded_page.description else {
83 return Err(Error::internal(
84 "Expected inner encoding to return structural layout".to_string(),
85 ));
86 };
87
88 let wrapped = ProtobufUtils21::blob_layout(inner_layout, &def_meaning);
89 Ok(EncodedPage {
90 column_idx: encoded_page.column_idx,
91 data: encoded_page.data,
92 description: PageEncoding::Structural(wrapped),
93 num_rows: encoded_page.num_rows,
94 row_number: encoded_page.row_number,
95 })
96 })
97 .boxed()
98 })
99 .collect::<Vec<_>>()
100 }
101}
102
103impl FieldEncoder for BlobStructuralEncoder {
104 fn maybe_encode(
105 &mut self,
106 array: ArrayRef,
107 external_buffers: &mut OutOfLineBuffers,
108 mut repdef: RepDefBuilder,
109 row_number: u64,
110 num_rows: u64,
111 ) -> Result<Vec<EncodeTask>> {
112 if let Some(validity) = array.nulls() {
113 repdef.add_validity_bitmap(validity.clone());
114 } else {
115 repdef.add_no_null(array.len());
116 }
117
118 let binary_array = array.as_binary_opt::<i64>().ok_or_else(|| {
120 Error::invalid_input_source(
121 format!("Expected LargeBinary array, got {}", array.data_type()).into(),
122 )
123 })?;
124
125 let repdef = RepDefBuilder::serialize(vec![repdef]);
126
127 let rep = repdef.repetition_levels.as_ref();
128 let def = repdef.definition_levels.as_ref();
129 let def_meaning: Arc<[DefinitionInterpretation]> = repdef.def_meaning.into();
130
131 let mut encode_tasks = match self.def_meaning.as_ref() {
136 Some(existing) if existing != &def_meaning => {
137 let existing = existing.clone();
138 Self::wrap_tasks(self.descriptor_encoder.flush(external_buffers)?, existing)
139 }
140 _ => Vec::new(),
141 };
142 self.def_meaning = Some(def_meaning.clone());
143
144 let mut positions = Vec::with_capacity(binary_array.len());
146 let mut sizes = Vec::with_capacity(binary_array.len());
147
148 for i in 0..binary_array.len() {
149 if binary_array.is_null(i) {
150 let mut repdef = (def.expect_ok()?[i] as u64) << 16;
154 if let Some(rep) = rep {
155 repdef += rep[i] as u64;
156 }
157
158 debug_assert_ne!(repdef, 0);
159 positions.push(repdef);
160 sizes.push(0);
161 } else {
162 let value = binary_array.value(i);
163 if value.is_empty() {
164 positions.push(0);
166 sizes.push(0);
167 } else {
168 let position =
170 external_buffers.add_buffer(LanceBuffer::from(Buffer::from(value)));
171 positions.push(position);
172 sizes.push(value.len() as u64);
173 }
174 }
175 }
176
177 let position_array = Arc::new(UInt64Array::from(positions));
179 let size_array = Arc::new(UInt64Array::from(sizes));
180 let descriptor_array = Arc::new(StructArray::new(
181 Fields::from(vec![
182 ArrowField::new("position", DataType::UInt64, false),
183 ArrowField::new("size", DataType::UInt64, false),
184 ]),
185 vec![position_array as ArrayRef, size_array as ArrayRef],
186 None, ));
188
189 let descriptor_tasks = self.descriptor_encoder.maybe_encode(
191 descriptor_array,
192 external_buffers,
193 RepDefBuilder::default(),
194 row_number,
195 num_rows,
196 )?;
197 encode_tasks.extend(Self::wrap_tasks(descriptor_tasks, def_meaning));
198
199 Ok(encode_tasks)
200 }
201
202 fn flush(&mut self, external_buffers: &mut OutOfLineBuffers) -> Result<Vec<EncodeTask>> {
203 let encode_tasks = self.descriptor_encoder.flush(external_buffers)?;
204
205 let def_meaning = self
208 .def_meaning
209 .clone()
210 .unwrap_or_else(|| Arc::new([DefinitionInterpretation::AllValidItem]));
211
212 Ok(Self::wrap_tasks(encode_tasks, def_meaning))
213 }
214
215 fn finish(
216 &mut self,
217 external_buffers: &mut OutOfLineBuffers,
218 ) -> BoxFuture<'_, Result<Vec<EncodedColumn>>> {
219 self.descriptor_encoder.finish(external_buffers)
220 }
221
222 fn num_columns(&self) -> u32 {
223 self.descriptor_encoder.num_columns()
224 }
225}
226
227pub struct BlobV2StructuralEncoder {
229 descriptor_encoder: Box<dyn FieldEncoder>,
230}
231
232impl BlobV2StructuralEncoder {
233 pub fn new(
234 field: &Field,
235 make_descriptor_encoder: impl FnOnce(Field) -> Result<Box<dyn FieldEncoder>>,
236 ) -> Result<Self> {
237 let mut descriptor_metadata = HashMap::with_capacity(1);
238 descriptor_metadata.insert(PACKED_STRUCT_META_KEY.to_string(), "true".to_string());
239
240 let descriptor_data_type = DataType::Struct(BLOB_V2_DESC_FIELDS.clone());
241
242 let descriptor_field = Field::try_from(
243 ArrowField::new(&field.name, descriptor_data_type, field.nullable)
244 .with_metadata(descriptor_metadata),
245 )?;
246
247 let descriptor_encoder = make_descriptor_encoder(descriptor_field)?;
248
249 Ok(Self { descriptor_encoder })
250 }
251}
252
253impl FieldEncoder for BlobV2StructuralEncoder {
254 fn maybe_encode(
255 &mut self,
256 array: ArrayRef,
257 external_buffers: &mut OutOfLineBuffers,
258 repdef: RepDefBuilder,
259 row_number: u64,
260 num_rows: u64,
261 ) -> Result<Vec<EncodeTask>> {
262 let struct_arr = array.as_struct();
263
264 let kind_col = struct_arr
265 .column_by_name("kind")
266 .ok_or_else(|| {
267 Error::invalid_input_source("Blob v2 struct missing `kind` field".into())
268 })?
269 .as_primitive::<UInt8Type>();
270 let data_col = struct_arr
271 .column_by_name("data")
272 .ok_or_else(|| {
273 Error::invalid_input_source("Blob v2 struct missing `data` field".into())
274 })?
275 .as_binary::<i64>();
276 let uri_col = struct_arr
277 .column_by_name("uri")
278 .ok_or_else(|| {
279 Error::invalid_input_source("Blob v2 struct missing `uri` field".into())
280 })?
281 .as_string::<i32>();
282 let blob_id_col = struct_arr
283 .column_by_name("blob_id")
284 .ok_or_else(|| {
285 Error::invalid_input_source("Blob v2 struct missing `blob_id` field".into())
286 })?
287 .as_primitive::<UInt32Type>();
288 let blob_size_col = struct_arr
289 .column_by_name("blob_size")
290 .ok_or_else(|| {
291 Error::invalid_input_source("Blob v2 struct missing `blob_size` field".into())
292 })?
293 .as_primitive::<UInt64Type>();
294 let packed_position_col = struct_arr
295 .column_by_name("position")
296 .ok_or_else(|| {
297 Error::invalid_input_source("Blob v2 struct missing `position` field".into())
298 })?
299 .as_primitive::<UInt64Type>();
300
301 let row_count = struct_arr.len();
302
303 let mut kind_builder = PrimitiveBuilder::<UInt8Type>::with_capacity(row_count);
304 let mut position_builder = PrimitiveBuilder::<UInt64Type>::with_capacity(row_count);
305 let mut size_builder = PrimitiveBuilder::<UInt64Type>::with_capacity(row_count);
306 let mut blob_id_builder = PrimitiveBuilder::<UInt32Type>::with_capacity(row_count);
307 let mut uri_builder = StringBuilder::with_capacity(row_count, row_count * 16);
308
309 for i in 0..row_count {
310 let (kind_value, position_value, size_value, blob_id_value, uri_value) =
311 if struct_arr.is_null(i) || kind_col.is_null(i) {
312 (BlobKind::Inline as u8, 0, 0, 0, "".to_string())
313 } else {
314 let kind_val = BlobKind::try_from(kind_col.value(i))?;
315 match kind_val {
316 BlobKind::Dedicated => (
317 BlobKind::Dedicated as u8,
318 0,
319 blob_size_col.value(i),
320 blob_id_col.value(i),
321 "".to_string(),
322 ),
323 BlobKind::External => {
324 let uri = uri_col.value(i).to_string();
325 let position = if packed_position_col.is_null(i) {
326 0
327 } else {
328 packed_position_col.value(i)
329 };
330 let size = if blob_size_col.is_null(i) {
331 0
332 } else {
333 blob_size_col.value(i)
334 };
335 let external_base_id = if blob_id_col.is_null(i) {
336 0
337 } else {
338 blob_id_col.value(i)
339 };
340 (
341 BlobKind::External as u8,
342 position,
343 size,
344 external_base_id,
345 uri,
346 )
347 }
348 BlobKind::Packed => (
349 BlobKind::Packed as u8,
350 packed_position_col.value(i),
351 blob_size_col.value(i),
352 blob_id_col.value(i),
353 "".to_string(),
354 ),
355 BlobKind::Inline => {
356 let data_val = data_col.value(i);
357 let blob_len = data_val.len() as u64;
358 let position = external_buffers
359 .add_buffer(LanceBuffer::from(Buffer::from(data_val)));
360
361 (
362 BlobKind::Inline as u8,
363 position,
364 blob_len,
365 0,
366 "".to_string(),
367 )
368 }
369 }
370 };
371
372 kind_builder.append_value(kind_value);
373 position_builder.append_value(position_value);
374 size_builder.append_value(size_value);
375 blob_id_builder.append_value(blob_id_value);
376 uri_builder.append_value(uri_value);
377 }
378 let children: Vec<ArrayRef> = vec![
379 Arc::new(kind_builder.finish()),
380 Arc::new(position_builder.finish()),
381 Arc::new(size_builder.finish()),
382 Arc::new(blob_id_builder.finish()),
383 Arc::new(uri_builder.finish()),
384 ];
385
386 let descriptor_array = Arc::new(StructArray::try_new(
387 BLOB_V2_DESC_FIELDS.clone(),
388 children,
389 struct_arr.nulls().cloned(),
390 )?) as ArrayRef;
391
392 self.descriptor_encoder.maybe_encode(
393 descriptor_array,
394 external_buffers,
395 repdef,
396 row_number,
397 num_rows,
398 )
399 }
400
401 fn flush(&mut self, external_buffers: &mut OutOfLineBuffers) -> Result<Vec<EncodeTask>> {
402 self.descriptor_encoder.flush(external_buffers)
403 }
404
405 fn finish(
406 &mut self,
407 external_buffers: &mut OutOfLineBuffers,
408 ) -> BoxFuture<'_, Result<Vec<EncodedColumn>>> {
409 self.descriptor_encoder.finish(external_buffers)
410 }
411
412 fn num_columns(&self) -> u32 {
413 self.descriptor_encoder.num_columns()
414 }
415}
416
417#[cfg(test)]
418mod tests {
419 use super::*;
420 use crate::{
421 encoder::{ColumnIndexSequence, EncodingOptions},
422 testing::{
423 TestCases, TestEncoding, check_round_trip_encoding_of_data,
424 check_round_trip_encoding_of_data_with_expected, create_test_field_encoder,
425 test_encoding_strategy,
426 },
427 };
428 use arrow_array::{
429 ArrayRef, LargeBinaryArray, StringArray, StructArray, UInt8Array, UInt32Array, UInt64Array,
430 };
431 use arrow_schema::{DataType, Field as ArrowField};
432
433 #[test]
434 fn test_blob_encoder_creation() {
435 let field = Field::try_from(
436 ArrowField::new("blob_field", DataType::LargeBinary, true).with_metadata(
437 HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]),
438 ),
439 )
440 .unwrap();
441 let mut column_index = ColumnIndexSequence::default();
442 let options = EncodingOptions::default();
443 let strategy = test_encoding_strategy(TestEncoding::StructuralU16);
444
445 let encoder =
446 create_test_field_encoder(strategy.as_ref(), &field, &mut column_index, &options);
447
448 assert!(encoder.is_ok());
449 }
450
451 #[tokio::test]
452 async fn test_blob_encoding_simple() {
453 let field = Field::try_from(
454 ArrowField::new("blob_field", DataType::LargeBinary, true).with_metadata(
455 HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]),
456 ),
457 )
458 .unwrap();
459 let mut column_index = ColumnIndexSequence::default();
460 let options = EncodingOptions::default();
461 let strategy = test_encoding_strategy(TestEncoding::StructuralU16);
462
463 let mut encoder =
464 create_test_field_encoder(strategy.as_ref(), &field, &mut column_index, &options)
465 .unwrap();
466
467 let large_data = vec![0u8; 1024 * 100]; let data: Vec<Option<&[u8]>> =
470 vec![Some(b"hello world"), None, Some(&large_data), Some(b"")];
471 let array = Arc::new(LargeBinaryArray::from(data));
472
473 let mut external_buffers = OutOfLineBuffers::new(0, 8);
475 let repdef = RepDefBuilder::default();
476
477 let tasks = encoder
478 .maybe_encode(array, &mut external_buffers, repdef, 0, 4)
479 .unwrap();
480
481 if tasks.is_empty() {
483 let _flush_tasks = encoder.flush(&mut external_buffers).unwrap();
484 }
485
486 assert!(encoder.num_columns() > 0);
489
490 let buffers = external_buffers.take_buffers();
492 assert!(
493 !buffers.is_empty(),
494 "Large blobs should be stored in external buffers"
495 );
496 }
497
498 #[tokio::test]
499 async fn test_blob_round_trip() {
500 let blob_metadata =
502 HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]);
503
504 let val1: &[u8] = &vec![1u8; 1024]; let val2: &[u8] = &vec![2u8; 10240]; let val3: &[u8] = &vec![3u8; 102400]; let array = Arc::new(LargeBinaryArray::from(vec![
509 Some(val1),
510 None,
511 Some(val2),
512 Some(val3),
513 ]));
514
515 check_round_trip_encoding_of_data(
517 vec![array],
518 &TestCases::default().with_array_and_u16_encodings(),
519 blob_metadata,
520 )
521 .await;
522 }
523
524 #[tokio::test]
525 async fn test_blob_round_trip_empty_values() {
526 let blob_metadata =
532 HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]);
533
534 let val1: &[u8] = &vec![1u8; 1024];
535 let val2: &[u8] = &vec![2u8; 10240];
536 let empty: &[u8] = &[];
537 let array = Arc::new(LargeBinaryArray::from(vec![
538 Some(empty),
539 Some(val1),
540 None,
541 Some(empty),
542 Some(val2),
543 None,
544 Some(empty),
545 ]));
546
547 check_round_trip_encoding_of_data(vec![array], &TestCases::default(), blob_metadata).await;
548 }
549
550 #[tokio::test]
551 async fn test_blob_round_trip_varying_chunk_nullability() {
552 let blob_metadata =
553 HashMap::from([(lance_arrow::BLOB_META_KEY.to_string(), "true".to_string())]);
554 let all_valid = Arc::new(LargeBinaryArray::from(vec![Some(b"first".as_ref())]));
555 let with_null = Arc::new(LargeBinaryArray::from(vec![
556 Some(b"second".as_ref()),
557 None,
558 Some(b"".as_ref()),
559 ]));
560 let all_valid_again = Arc::new(LargeBinaryArray::from(vec![Some(b"last".as_ref())]));
561
562 check_round_trip_encoding_of_data(
563 vec![all_valid, with_null, all_valid_again],
564 &TestCases::default().with_encoding(TestEncoding::StructuralU16),
565 blob_metadata,
566 )
567 .await;
568 }
569
570 #[tokio::test]
571 async fn test_blob_v2_external_round_trip() {
572 let blob_metadata = HashMap::from([(
573 lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
574 lance_arrow::BLOB_V2_EXT_NAME.to_string(),
575 )]);
576
577 let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
578 let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
579 let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
580 let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
581 let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
582 let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
583
584 let kind_array = UInt8Array::from(vec![
585 BlobKind::Inline as u8,
586 BlobKind::External as u8,
587 BlobKind::External as u8,
588 ]);
589 let data_array = LargeBinaryArray::from(vec![Some(b"inline".as_ref()), None, None]);
590 let uri_array = StringArray::from(vec![
591 None,
592 Some("file:///tmp/external.bin"),
593 Some("s3://bucket/blob"),
594 ]);
595 let blob_id_array = UInt32Array::from(vec![0, 0, 0]);
596 let blob_size_array = UInt64Array::from(vec![0, 0, 0]);
597 let position_array = UInt64Array::from(vec![0, 0, 0]);
598
599 let struct_array = StructArray::from(vec![
600 (kind_field, Arc::new(kind_array) as ArrayRef),
601 (data_field, Arc::new(data_array) as ArrayRef),
602 (uri_field, Arc::new(uri_array) as ArrayRef),
603 (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
604 (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
605 (position_field, Arc::new(position_array) as ArrayRef),
606 ]);
607
608 let expected_descriptor = StructArray::from(vec![
609 (
610 Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
611 Arc::new(UInt8Array::from(vec![
612 BlobKind::Inline as u8,
613 BlobKind::External as u8,
614 BlobKind::External as u8,
615 ])) as ArrayRef,
616 ),
617 (
618 Arc::new(ArrowField::new("position", DataType::UInt64, false)),
619 Arc::new(UInt64Array::from(vec![0, 0, 0])) as ArrayRef,
620 ),
621 (
622 Arc::new(ArrowField::new("size", DataType::UInt64, false)),
623 Arc::new(UInt64Array::from(vec![6, 0, 0])) as ArrayRef,
624 ),
625 (
626 Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
627 Arc::new(UInt32Array::from(vec![0, 0, 0])) as ArrayRef,
628 ),
629 (
630 Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
631 Arc::new(StringArray::from(vec![
632 "",
633 "file:///tmp/external.bin",
634 "s3://bucket/blob",
635 ])) as ArrayRef,
636 ),
637 ]);
638
639 check_round_trip_encoding_of_data_with_expected(
640 vec![Arc::new(struct_array)],
641 Some(Arc::new(expected_descriptor)),
642 &TestCases::default().with_u32_structural_encodings(),
643 blob_metadata,
644 )
645 .await;
646 }
647
648 #[tokio::test]
649 async fn test_blob_v2_dedicated_round_trip() {
650 let blob_metadata = HashMap::from([(
651 lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
652 lance_arrow::BLOB_V2_EXT_NAME.to_string(),
653 )]);
654
655 let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
656 let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
657 let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
658 let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
659 let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
660 let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
661
662 let kind_array = UInt8Array::from(vec![BlobKind::Dedicated as u8, BlobKind::Inline as u8]);
663 let data_array = LargeBinaryArray::from(vec![None, Some(b"abc".as_ref())]);
664 let uri_array = StringArray::from(vec![Option::<&str>::None, None]);
665 let blob_id_array = UInt32Array::from(vec![42, 0]);
666 let blob_size_array = UInt64Array::from(vec![12, 0]);
667 let position_array = UInt64Array::from(vec![0, 0]);
668
669 let struct_array = StructArray::from(vec![
670 (kind_field, Arc::new(kind_array) as ArrayRef),
671 (data_field, Arc::new(data_array) as ArrayRef),
672 (uri_field, Arc::new(uri_array) as ArrayRef),
673 (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
674 (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
675 (position_field, Arc::new(position_array) as ArrayRef),
676 ]);
677
678 let expected_descriptor = StructArray::from(vec![
679 (
680 Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
681 Arc::new(UInt8Array::from(vec![
682 BlobKind::Dedicated as u8,
683 BlobKind::Inline as u8,
684 ])) as ArrayRef,
685 ),
686 (
687 Arc::new(ArrowField::new("position", DataType::UInt64, false)),
688 Arc::new(UInt64Array::from(vec![0, 0])) as ArrayRef,
689 ),
690 (
691 Arc::new(ArrowField::new("size", DataType::UInt64, false)),
692 Arc::new(UInt64Array::from(vec![12, 3])) as ArrayRef,
693 ),
694 (
695 Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
696 Arc::new(UInt32Array::from(vec![42, 0])) as ArrayRef,
697 ),
698 (
699 Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
700 Arc::new(StringArray::from(vec!["", ""])) as ArrayRef,
701 ),
702 ]);
703
704 check_round_trip_encoding_of_data_with_expected(
705 vec![Arc::new(struct_array)],
706 Some(Arc::new(expected_descriptor)),
707 &TestCases::default().with_u32_structural_encodings(),
708 blob_metadata,
709 )
710 .await;
711 }
712
713 #[tokio::test]
714 async fn test_blob_v2_external_with_range_round_trip() {
715 let blob_metadata = HashMap::from([(
716 lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
717 lance_arrow::BLOB_V2_EXT_NAME.to_string(),
718 )]);
719
720 let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
721 let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
722 let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
723 let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
724 let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
725 let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
726
727 let kind_array = UInt8Array::from(vec![BlobKind::External as u8]);
728 let data_array = LargeBinaryArray::from(vec![None::<&[u8]>]);
729 let uri_array = StringArray::from(vec![Some("memory://container.pack")]);
730 let blob_id_array = UInt32Array::from(vec![0]);
731 let blob_size_array = UInt64Array::from(vec![42]);
732 let position_array = UInt64Array::from(vec![7]);
733
734 let struct_array = StructArray::from(vec![
735 (kind_field, Arc::new(kind_array) as ArrayRef),
736 (data_field, Arc::new(data_array) as ArrayRef),
737 (uri_field, Arc::new(uri_array) as ArrayRef),
738 (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
739 (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
740 (position_field, Arc::new(position_array) as ArrayRef),
741 ]);
742
743 let expected_descriptor = StructArray::from(vec![
744 (
745 Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
746 Arc::new(UInt8Array::from(vec![BlobKind::External as u8])) as ArrayRef,
747 ),
748 (
749 Arc::new(ArrowField::new("position", DataType::UInt64, false)),
750 Arc::new(UInt64Array::from(vec![7])) as ArrayRef,
751 ),
752 (
753 Arc::new(ArrowField::new("size", DataType::UInt64, false)),
754 Arc::new(UInt64Array::from(vec![42])) as ArrayRef,
755 ),
756 (
757 Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
758 Arc::new(UInt32Array::from(vec![0])) as ArrayRef,
759 ),
760 (
761 Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
762 Arc::new(StringArray::from(vec!["memory://container.pack"])) as ArrayRef,
763 ),
764 ]);
765
766 check_round_trip_encoding_of_data_with_expected(
767 vec![Arc::new(struct_array)],
768 Some(Arc::new(expected_descriptor)),
769 &TestCases::default().with_u32_structural_encodings(),
770 blob_metadata,
771 )
772 .await;
773 }
774
775 #[tokio::test]
776 async fn test_blob_v2_packed_round_trip() {
777 let blob_metadata = HashMap::from([(
778 lance_arrow::ARROW_EXT_NAME_KEY.to_string(),
779 lance_arrow::BLOB_V2_EXT_NAME.to_string(),
780 )]);
781
782 let kind_field = Arc::new(ArrowField::new("kind", DataType::UInt8, true));
783 let data_field = Arc::new(ArrowField::new("data", DataType::LargeBinary, true));
784 let uri_field = Arc::new(ArrowField::new("uri", DataType::Utf8, true));
785 let blob_id_field = Arc::new(ArrowField::new("blob_id", DataType::UInt32, true));
786 let blob_size_field = Arc::new(ArrowField::new("blob_size", DataType::UInt64, true));
787 let position_field = Arc::new(ArrowField::new("position", DataType::UInt64, true));
788
789 let kind_array = UInt8Array::from(vec![BlobKind::Packed as u8]);
790 let data_array = LargeBinaryArray::from(vec![None::<&[u8]>]);
791 let uri_array = StringArray::from(vec![None::<&str>]);
792 let blob_id_array = UInt32Array::from(vec![7]);
793 let blob_size_array = UInt64Array::from(vec![5]);
794 let position_array = UInt64Array::from(vec![10]);
795
796 let struct_array = StructArray::from(vec![
797 (kind_field, Arc::new(kind_array) as ArrayRef),
798 (data_field, Arc::new(data_array) as ArrayRef),
799 (uri_field, Arc::new(uri_array) as ArrayRef),
800 (blob_id_field, Arc::new(blob_id_array) as ArrayRef),
801 (blob_size_field, Arc::new(blob_size_array) as ArrayRef),
802 (position_field, Arc::new(position_array) as ArrayRef),
803 ]);
804
805 let expected_descriptor = StructArray::from(vec![
806 (
807 Arc::new(ArrowField::new("kind", DataType::UInt8, false)),
808 Arc::new(UInt8Array::from(vec![BlobKind::Packed as u8])) as ArrayRef,
809 ),
810 (
811 Arc::new(ArrowField::new("position", DataType::UInt64, false)),
812 Arc::new(UInt64Array::from(vec![10])) as ArrayRef,
813 ),
814 (
815 Arc::new(ArrowField::new("size", DataType::UInt64, false)),
816 Arc::new(UInt64Array::from(vec![5])) as ArrayRef,
817 ),
818 (
819 Arc::new(ArrowField::new("blob_id", DataType::UInt32, false)),
820 Arc::new(UInt32Array::from(vec![7])) as ArrayRef,
821 ),
822 (
823 Arc::new(ArrowField::new("blob_uri", DataType::Utf8, false)),
824 Arc::new(StringArray::from(vec![""])) as ArrayRef,
825 ),
826 ]);
827
828 check_round_trip_encoding_of_data_with_expected(
829 vec![Arc::new(struct_array)],
830 Some(Arc::new(expected_descriptor)),
831 &TestCases::default().with_u32_structural_encodings(),
832 blob_metadata,
833 )
834 .await;
835 }
836}