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