1pub mod session;
7pub mod write;
8
9use std::borrow::Cow;
10use std::collections::{BTreeMap, HashMap, HashSet};
11use std::ops::Range;
12use std::sync::Arc;
13
14use arrow::compute::concat_batches;
15use arrow_array::cast::as_primitive_array;
16use arrow_array::{
17 RecordBatch, RecordBatchReader, StructArray, UInt32Array, UInt64Array, new_null_array,
18};
19use arrow_schema::Schema as ArrowSchema;
20use datafusion::logical_expr::Expr;
21use datafusion::scalar::ScalarValue;
22use futures::future::try_join_all;
23use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt, join, stream};
24use lance_arrow::{RecordBatchExt, SchemaExt};
25use lance_core::datatypes::{OnMissing, OnTypeMismatch, SchemaCompareOptions};
26use lance_core::utils::deletion::DeletionVector;
27use lance_core::utils::tokio::get_num_compute_intensive_cpus;
28use lance_core::{Error, Result, cache::CacheKey, datatypes::Schema};
29use lance_core::{
30 ROW_ADDR, ROW_ADDR_FIELD, ROW_CREATED_AT_VERSION_FIELD, ROW_ID, ROW_ID_FIELD,
31 ROW_LAST_UPDATED_AT_VERSION_FIELD,
32};
33use lance_datafusion::utils::StreamingWriteSource;
34use lance_encoding::decoder::DecoderPlugins;
35use lance_file::previous::reader::{
36 FileReader as PreviousFileReader, read_batch as previous_read_batch,
37};
38use lance_file::reader::{CachedFileMetadata, FileReaderOptions, ReaderProjection};
39use lance_file::version::LanceFileVersion;
40use lance_file::{LanceEncodingsIo, determine_file_version};
41use lance_io::ReadBatchParams;
42use lance_io::scheduler::{FileScheduler, ScanScheduler, SchedulerConfig};
43use lance_io::utils::CachedFileSize;
44use lance_table::format::{DataFile, DeletionFile, Fragment};
45use lance_table::io::deletion::{deletion_file_path, write_deletion_file};
46use lance_table::rowids::RowIdSequence;
47use lance_table::utils::stream::{
48 ReadBatchFutStream, ReadBatchTask, ReadBatchTaskStream, RowIdAndDeletesConfig,
49 wrap_with_row_id_and_delete,
50};
51
52use self::write::FragmentCreateBuilder;
53
54use super::hash_joiner::HashJoiner;
55use super::rowids::load_row_id_sequence;
56use super::scanner::Scanner;
57
58use super::updater::Updater;
59use super::{NewColumnTransform, WriteParams, schema_evolution};
60use crate::dataset::Dataset;
61use crate::dataset::fragment::session::FragmentSession;
62use crate::io::deletion::read_dataset_deletion_file;
63
64#[derive(Debug, Clone)]
68pub struct FileFragment {
69 dataset: Arc<Dataset>,
70
71 pub(super) metadata: Fragment,
72}
73
74const DEFAULT_BATCH_READ_SIZE: u32 = 1024;
75
76#[allow(clippy::len_without_is_empty)]
78pub trait GenericFileReader: std::fmt::Debug + Send + Sync {
79 fn read_range_tasks(
82 &self,
83 range: Range<u64>,
84 batch_size: u32,
85 projection: Arc<lance_core::datatypes::Schema>,
86 ) -> Result<ReadBatchTaskStream>;
87 fn read_ranges_tasks(
89 &self,
90 ranges: Arc<[Range<u64>]>,
91 batch_size: u32,
92 projection: Arc<lance_core::datatypes::Schema>,
93 ) -> Result<ReadBatchTaskStream>;
94 fn read_all_tasks(
96 &self,
97 batch_size: u32,
98 projection: Arc<lance_core::datatypes::Schema>,
99 ) -> Result<ReadBatchTaskStream>;
100 fn take_all_tasks(
102 &self,
103 indices: &[u32],
104 batch_size: u32,
105 projection: Arc<lance_core::datatypes::Schema>,
106 take_priority: Option<u32>,
107 ) -> Result<ReadBatchTaskStream>;
108
109 fn len(&self) -> u32;
111
112 fn projection(&self) -> &Arc<Schema>;
114
115 fn storage_stats(&self) -> Vec<(u32, u64)>;
117
118 fn clone_box(&self) -> Box<dyn GenericFileReader>;
124 fn is_legacy(&self) -> bool;
126 fn as_legacy(&self) -> &PreviousFileReader {
129 self.as_legacy_opt()
130 .expect("legacy function called on v2 file")
131 }
132 fn as_legacy_opt(&self) -> Option<&PreviousFileReader>;
135 fn as_legacy_opt_mut(&mut self) -> Option<&mut PreviousFileReader>;
138}
139
140fn ranges_to_tasks(
141 reader: &PreviousFileReader,
142 ranges: Vec<(i32, Range<usize>)>,
143 projection: Arc<Schema>,
144) -> ReadBatchTaskStream {
145 let reader = reader.clone();
146 stream::iter(ranges)
147 .map(move |(batch_idx, range)| {
148 let num_rows = range.end - range.start;
149 let reader = reader.clone();
150 let projection = projection.clone();
151 let task = tokio::task::spawn(async move {
152 previous_read_batch(
153 &reader,
154 &ReadBatchParams::Range(range.clone()),
155 &projection,
156 batch_idx,
157 )
158 .await
159 })
160 .map(|task_out| task_out.unwrap())
161 .boxed();
162 ReadBatchTask {
163 task,
164 num_rows: num_rows as u32,
165 }
166 })
167 .boxed()
168}
169
170#[derive(Clone, Debug)]
171struct V1Reader {
172 reader: PreviousFileReader,
173 projection: Arc<Schema>,
174}
175
176impl V1Reader {
177 fn new(reader: PreviousFileReader, projection: Arc<Schema>) -> Self {
178 Self { reader, projection }
179 }
180}
181
182impl GenericFileReader for V1Reader {
183 fn read_range_tasks(
185 &self,
186 range: Range<u64>,
187 batch_size: u32,
188 projection: Arc<Schema>,
189 ) -> Result<ReadBatchTaskStream> {
190 let mut to_skip = range.start as u32;
191 let mut remaining = range.end as u32 - to_skip;
192 let mut ranges = Vec::new();
193 let mut batch_idx = 0;
194 while remaining > 0 {
195 let next_batch_len = self.reader.num_rows_in_batch(batch_idx) as u32;
196 let next_batch_idx = batch_idx;
197 batch_idx += 1;
198 if to_skip >= next_batch_len {
199 to_skip -= next_batch_len;
200 continue;
201 }
202 let batch_start = to_skip;
203 to_skip = 0;
204 let batch_end = next_batch_len.min(batch_start + remaining);
205 remaining -= batch_end - batch_start;
206 for chunk_start in (batch_start..batch_end).step_by(batch_size as usize) {
207 let chunk_end = (chunk_start + batch_size).min(batch_end);
208 ranges.push((next_batch_idx, (chunk_start as usize..chunk_end as usize)));
209 }
210 }
211 Ok(ranges_to_tasks(&self.reader, ranges, projection))
212 }
213
214 fn read_all_tasks(
215 &self,
216 batch_size: u32,
217 projection: Arc<Schema>,
218 ) -> Result<ReadBatchTaskStream> {
219 let ranges = (0..self.reader.num_batches())
220 .flat_map(move |batch_idx| {
221 let rows_in_batch = self.reader.num_rows_in_batch(batch_idx as i32);
222 (0..rows_in_batch)
223 .step_by(batch_size as usize)
224 .map(move |start| {
225 let end = (start + batch_size as usize).min(rows_in_batch);
226 (batch_idx as i32, start..end)
227 })
228 })
229 .collect::<Vec<_>>();
230 Ok(ranges_to_tasks(&self.reader, ranges, projection))
231 }
232
233 fn read_ranges_tasks(
234 &self,
235 _ranges: Arc<[Range<u64>]>,
236 _batch_size: u32,
237 _projection: Arc<Schema>,
238 ) -> Result<ReadBatchTaskStream> {
239 Err(Error::internal(
240 "Attempt to perform FilteredRead on v1 files".to_string(),
241 ))
242 }
243
244 fn take_all_tasks(
245 &self,
246 indices: &[u32],
247 _batch_size: u32,
248 projection: Arc<Schema>,
249 _take_priority: Option<u32>,
250 ) -> Result<ReadBatchTaskStream> {
251 let indices_vec = indices.to_vec();
252 let reader = self.reader.clone();
253 let task_fut = async move { reader.take(&indices_vec, projection.as_ref()).await }.boxed();
255 let task = std::future::ready(ReadBatchTask {
256 task: task_fut,
257 num_rows: indices.len() as u32,
258 })
259 .boxed();
260 Ok(futures::stream::once(task).boxed())
261 }
262
263 fn projection(&self) -> &Arc<Schema> {
264 &self.projection
265 }
266
267 fn len(&self) -> u32 {
269 self.reader.len() as u32
270 }
271
272 fn storage_stats(&self) -> Vec<(u32, u64)> {
273 Vec::new()
275 }
276
277 fn clone_box(&self) -> Box<dyn GenericFileReader> {
278 Box::new(self.clone())
279 }
280
281 fn is_legacy(&self) -> bool {
282 true
283 }
284
285 fn as_legacy_opt(&self) -> Option<&PreviousFileReader> {
286 Some(&self.reader)
287 }
288
289 fn as_legacy_opt_mut(&mut self) -> Option<&mut PreviousFileReader> {
290 Some(&mut self.reader)
291 }
292}
293
294mod v2_adapter {
295 use lance_encoding::decoder::FilterExpression;
296
297 use super::*;
298
299 #[derive(Debug, Clone)]
300 pub struct Reader {
301 reader: Arc<lance_file::reader::FileReader>,
302 projection: Arc<Schema>,
303 field_id_to_column_idx: Arc<BTreeMap<u32, u32>>,
304 default_priority: u32,
305 file_scheduler: FileScheduler,
306 }
307
308 impl Reader {
309 pub fn new(
310 reader: Arc<lance_file::reader::FileReader>,
311 projection: Arc<Schema>,
312 field_id_to_column_idx: Arc<BTreeMap<u32, u32>>,
313 default_priority: u32,
314 file_scheduler: FileScheduler,
315 ) -> Self {
316 Self {
317 reader,
318 projection,
319 field_id_to_column_idx,
320 default_priority,
321 file_scheduler,
322 }
323 }
324 }
325
326 impl GenericFileReader for Reader {
327 fn read_range_tasks(
329 &self,
330 range: Range<u64>,
331 batch_size: u32,
332 projection: Arc<Schema>,
333 ) -> Result<ReadBatchTaskStream> {
334 let projection = ReaderProjection::from_field_ids(
335 self.reader.metadata().version(),
336 projection.as_ref(),
337 self.field_id_to_column_idx.as_ref(),
338 )?;
339 Ok(self
340 .reader
341 .read_tasks(
342 ReadBatchParams::Range(range.start as usize..range.end as usize),
343 batch_size,
344 Some(projection),
345 FilterExpression::no_filter(),
346 )?
347 .map(|v2_task| ReadBatchTask {
348 task: v2_task.task.map_err(Error::from).boxed(),
349 num_rows: v2_task.num_rows,
350 })
351 .boxed())
352 }
353
354 fn read_ranges_tasks(
355 &self,
356 ranges: Arc<[Range<u64>]>,
357 batch_size: u32,
358 projection: Arc<Schema>,
359 ) -> Result<ReadBatchTaskStream> {
360 let projection = ReaderProjection::from_field_ids(
361 self.reader.metadata().version(),
362 projection.as_ref(),
363 self.field_id_to_column_idx.as_ref(),
364 )?;
365 Ok(self
366 .reader
367 .read_tasks(
368 ReadBatchParams::Ranges(ranges),
369 batch_size,
370 Some(projection),
371 FilterExpression::no_filter(),
372 )?
373 .map(|v2_task| ReadBatchTask {
374 task: v2_task.task.map_err(Error::from).boxed(),
375 num_rows: v2_task.num_rows,
376 })
377 .boxed())
378 }
379
380 fn read_all_tasks(
381 &self,
382 batch_size: u32,
383 projection: Arc<Schema>,
384 ) -> Result<ReadBatchTaskStream> {
385 let projection = ReaderProjection::from_field_ids(
386 self.reader.metadata().version(),
387 projection.as_ref(),
388 self.field_id_to_column_idx.as_ref(),
389 )?;
390 Ok(self
391 .reader
392 .read_tasks(
393 ReadBatchParams::RangeFull,
394 batch_size,
395 Some(projection),
396 FilterExpression::no_filter(),
397 )?
398 .map(|v2_task| ReadBatchTask {
399 task: v2_task.task.map_err(Error::from).boxed(),
400 num_rows: v2_task.num_rows,
401 })
402 .boxed())
403 }
404
405 fn take_all_tasks(
406 &self,
407 indices: &[u32],
408 batch_size: u32,
409 projection: Arc<Schema>,
410 take_priority: Option<u32>,
411 ) -> Result<ReadBatchTaskStream> {
412 let indices = UInt32Array::from(indices.to_vec());
413 let projection = ReaderProjection::from_field_ids(
414 self.reader.metadata().version(),
415 projection.as_ref(),
416 self.field_id_to_column_idx.as_ref(),
417 )?;
418
419 let reader = if let Some(take_priority) = take_priority {
420 let op_priority = ((take_priority as u64) << 32) | self.default_priority as u64;
421 let scheduler = self.file_scheduler.with_priority(op_priority);
422 Arc::new(
423 self.reader
424 .with_scheduler(Arc::new(LanceEncodingsIo::new(scheduler))),
425 )
426 } else {
427 self.reader.clone()
428 };
429
430 Ok(reader
431 .read_tasks(
432 ReadBatchParams::Indices(indices),
433 batch_size,
434 Some(projection),
435 FilterExpression::no_filter(),
436 )?
437 .map(|v2_task| ReadBatchTask {
438 task: v2_task.task.map_err(Error::from).boxed(),
439 num_rows: v2_task.num_rows,
440 })
441 .boxed())
442 }
443
444 fn storage_stats(&self) -> Vec<(u32, u64)> {
445 let file_statistics = self.reader.file_statistics();
446 let column_idx_to_field_id = self
447 .field_id_to_column_idx
448 .iter()
449 .map(|(field_id, column_idx)| (*column_idx, *field_id))
450 .collect::<HashMap<_, _>>();
451
452 let mut stats = Vec::new();
453 let mut current_field_id = 0;
456 for (column_idx, col_stats) in file_statistics.columns.iter().enumerate() {
457 if let Some(field_id) = column_idx_to_field_id.get(&(column_idx as u32)) {
458 current_field_id = *field_id;
459 }
460 stats.push((current_field_id, col_stats.size_bytes));
461 }
462 stats
463 }
464
465 fn projection(&self) -> &Arc<Schema> {
466 &self.projection
467 }
468
469 fn len(&self) -> u32 {
471 self.reader.metadata().num_rows as u32
472 }
473
474 fn clone_box(&self) -> Box<dyn GenericFileReader> {
475 Box::new(self.clone())
476 }
477
478 fn is_legacy(&self) -> bool {
479 false
480 }
481
482 fn as_legacy_opt(&self) -> Option<&PreviousFileReader> {
483 None
484 }
485
486 fn as_legacy_opt_mut(&mut self) -> Option<&mut PreviousFileReader> {
487 None
488 }
489 }
490}
491
492#[derive(Debug, Clone)]
495struct NullReader {
496 schema: Arc<Schema>,
497 num_rows: u32,
498}
499
500impl NullReader {
501 fn new(schema: Arc<Schema>, num_rows: u32) -> Self {
502 Self { schema, num_rows }
503 }
504
505 fn batch(projection: Arc<ArrowSchema>, num_rows: usize) -> RecordBatch {
506 let columns = projection
507 .fields()
508 .iter()
509 .map(|f| new_null_array(f.data_type(), num_rows))
510 .collect::<Vec<_>>();
511 RecordBatch::try_new(projection, columns).unwrap()
512 }
513}
514
515impl GenericFileReader for NullReader {
516 fn read_range_tasks(
517 &self,
518 range: Range<u64>,
519 batch_size: u32,
520 projection: Arc<Schema>,
521 ) -> Result<ReadBatchTaskStream> {
522 self.read_ranges_tasks(vec![range].into(), batch_size, projection)
523 }
524
525 fn read_ranges_tasks(
526 &self,
527 ranges: Arc<[Range<u64>]>,
528 batch_size: u32,
529 projection: Arc<Schema>,
530 ) -> Result<ReadBatchTaskStream> {
531 let mut remaining_rows = ranges.iter().map(|r| r.end - r.start).sum::<u64>();
532 let projection: Arc<ArrowSchema> = Arc::new(projection.as_ref().into());
533
534 let task_iter = std::iter::from_fn(move || {
535 if remaining_rows == 0 {
536 return None;
537 }
538
539 let num_rows = remaining_rows.min(batch_size as u64) as usize;
540 remaining_rows -= num_rows as u64;
541 let batch = Self::batch(projection.clone(), num_rows);
542 let task = ReadBatchTask {
543 task: futures::future::ready(Ok(batch)).boxed(),
544 num_rows: num_rows as u32,
545 };
546 Some(task)
547 });
548
549 Ok(futures::stream::iter(task_iter).boxed())
550 }
551
552 fn read_all_tasks(
553 &self,
554 batch_size: u32,
555 projection: Arc<Schema>,
556 ) -> Result<ReadBatchTaskStream> {
557 self.read_ranges_tasks(vec![0..self.num_rows as u64].into(), batch_size, projection)
558 }
559
560 fn take_all_tasks(
561 &self,
562 indices: &[u32],
563 batch_size: u32,
564 projection: Arc<Schema>,
565 _take_priority: Option<u32>,
566 ) -> Result<ReadBatchTaskStream> {
567 let num_rows = indices.len() as u64;
568 self.read_ranges_tasks(vec![0..num_rows].into(), batch_size, projection)
569 }
570
571 fn storage_stats(&self) -> Vec<(u32, u64)> {
572 Vec::new()
574 }
575
576 fn projection(&self) -> &Arc<Schema> {
577 &self.schema
578 }
579
580 fn len(&self) -> u32 {
581 self.num_rows
582 }
583
584 fn clone_box(&self) -> Box<dyn GenericFileReader> {
585 Box::new(self.clone())
586 }
587
588 fn is_legacy(&self) -> bool {
589 false
590 }
591
592 fn as_legacy_opt(&self) -> Option<&PreviousFileReader> {
593 None
594 }
595
596 fn as_legacy_opt_mut(&mut self) -> Option<&mut PreviousFileReader> {
597 None
598 }
599}
600
601#[derive(Debug, Default)]
602pub struct FragReadConfig {
603 pub with_row_id: bool,
605 pub with_row_address: bool,
607 pub with_row_last_updated_at_version: bool,
609 pub with_row_created_at_version: bool,
611 pub scan_scheduler: Option<Arc<ScanScheduler>>,
616 pub reader_priority: Option<u32>,
624 pub file_reader_options: Option<FileReaderOptions>,
626}
627
628impl FragReadConfig {
629 pub fn with_row_id(mut self, value: bool) -> Self {
630 self.with_row_id = value;
631 self
632 }
633
634 pub fn with_row_address(mut self, value: bool) -> Self {
635 self.with_row_address = value;
636 self
637 }
638
639 pub fn with_row_last_updated_at_version(mut self, value: bool) -> Self {
640 self.with_row_last_updated_at_version = value;
641 self
642 }
643
644 pub fn with_row_created_at_version(mut self, value: bool) -> Self {
645 self.with_row_created_at_version = value;
646 self
647 }
648
649 pub fn has_system_cols(&self) -> bool {
650 self.with_row_id
651 || self.with_row_address
652 || self.with_row_last_updated_at_version
653 || self.with_row_created_at_version
654 }
655
656 pub fn with_scan_scheduler(mut self, value: Arc<ScanScheduler>) -> Self {
657 self.scan_scheduler = Some(value);
658 self
659 }
660
661 pub fn with_reader_priority(mut self, value: u32) -> Self {
662 self.reader_priority = Some(value);
663 self
664 }
665
666 pub fn with_file_reader_options(mut self, value: FileReaderOptions) -> Self {
667 self.file_reader_options = Some(value);
668 self
669 }
670}
671
672impl FileFragment {
673 pub fn new(dataset: Arc<Dataset>, metadata: Fragment) -> Self {
675 Self { dataset, metadata }
676 }
677
678 pub async fn create(
685 dataset_uri: &str,
686 id: usize,
687 source: impl StreamingWriteSource,
688 params: Option<WriteParams>,
689 ) -> Result<Fragment> {
690 let mut builder = FragmentCreateBuilder::new(dataset_uri);
691
692 if let Some(params) = params.as_ref() {
693 builder = builder.write_params(params);
694 }
695
696 builder.write(source, Some(id as u64)).await
697 }
698
699 pub async fn create_fragments(
701 dataset_uri: &str,
702 source: impl StreamingWriteSource,
703 params: Option<WriteParams>,
704 ) -> Result<Vec<Fragment>> {
705 let mut builder = FragmentCreateBuilder::new(dataset_uri);
706
707 if let Some(params) = params.as_ref() {
708 builder = builder.write_params(params);
709 }
710
711 builder.write_fragments(source).await
712 }
713
714 pub async fn create_from_file(
715 filename: &str,
716 dataset: &Dataset,
717 fragment_id: usize,
718 physical_rows: Option<usize>,
719 ) -> Result<Fragment> {
720 let filepath = dataset.data_dir().child(filename);
721 let file_version =
722 determine_file_version(dataset.object_store.as_ref(), &filepath, None).await?;
723
724 if file_version != dataset.manifest.data_storage_format.lance_file_version()? {
725 return Err(Error::invalid_input(format!(
726 "File version mismatch. Dataset version: {:?} Fragment version: {:?}",
727 dataset.manifest.data_storage_format.lance_file_version()?,
728 file_version
729 )));
730 }
731
732 if file_version == LanceFileVersion::Legacy {
733 let fragment = Fragment::with_file_legacy(
734 fragment_id as u64,
735 filename,
736 dataset.schema(),
737 physical_rows,
738 );
739 Ok(fragment)
740 } else {
741 let mut frag = Fragment::new(fragment_id as u64);
744 let scheduler = ScanScheduler::new(
745 dataset.object_store.clone(),
746 SchedulerConfig::max_bandwidth(&dataset.object_store),
747 );
748 let file_scheduler = scheduler
749 .open_file(&filepath, &CachedFileSize::unknown())
750 .await?;
751 let reader = lance_file::reader::FileReader::try_open(
752 file_scheduler,
753 None,
754 Arc::<DecoderPlugins>::default(),
755 &dataset.metadata_cache.file_metadata_cache(&filepath),
756 dataset.file_reader_options.clone().unwrap_or_default(),
757 )
758 .await?;
759 reader
761 .schema()
762 .check_compatible(dataset.schema(), &SchemaCompareOptions::default())?;
763 let projection = lance_file::reader::ReaderProjection::from_whole_schema(
764 dataset.schema(),
765 reader.metadata().version(),
766 );
767 let physical_rows = reader.metadata().num_rows as usize;
768 frag.physical_rows = Some(physical_rows);
769 frag.id = fragment_id as u64;
770
771 let column_indices = projection
772 .column_indices
773 .into_iter()
774 .map(|c| c as i32)
775 .collect();
776
777 frag.add_file(
778 filename,
779 dataset.schema().field_ids(),
780 column_indices,
781 &file_version,
782 None,
783 );
784 Ok(frag)
785 }
786 }
787
788 pub(crate) async fn storage_stats(
790 &self,
791 dataset_schema: &Schema,
792 scan_scheduler: Arc<ScanScheduler>,
793 ) -> Result<Vec<(u32, u64)>> {
794 let mut stats = Vec::new();
795 for reader in self
796 .open_readers(
797 dataset_schema,
798 &FragReadConfig::default().with_scan_scheduler(scan_scheduler),
799 )
800 .await?
801 {
802 stats.extend(reader.storage_stats());
803 }
804 Ok(stats)
805 }
806
807 pub fn dataset(&self) -> &Dataset {
808 self.dataset.as_ref()
809 }
810
811 pub fn schema(&self) -> &Schema {
812 self.dataset.schema()
813 }
814
815 pub fn metadata(&self) -> &Fragment {
817 &self.metadata
818 }
819
820 pub fn id(&self) -> usize {
822 self.metadata.id as usize
823 }
824
825 pub fn num_data_files(&self) -> usize {
827 self.metadata.files.len()
828 }
829
830 pub fn data_file_for_field(&self, field_id: u32) -> Option<&DataFile> {
832 self.metadata
833 .files
834 .iter()
835 .find(|f| f.fields.contains(&(field_id as i32)))
836 }
837
838 pub async fn open(
853 &self,
854 projection: &Schema,
855 read_config: FragReadConfig,
856 ) -> Result<FragmentReader> {
857 let open_files = self.open_readers(projection, &read_config);
858 let deletion_vec_load = self.get_deletion_vector();
859
860 let row_id_load = if self.dataset.manifest.uses_stable_row_ids() {
861 futures::future::Either::Left(
862 load_row_id_sequence(&self.dataset, &self.metadata).map_ok(Some),
863 )
864 } else {
865 futures::future::Either::Right(futures::future::ready(Ok(None)))
866 };
867
868 let (opened_files, deletion_vec, row_id_sequence) =
869 join!(open_files, deletion_vec_load, row_id_load);
870 let opened_files = opened_files?;
871 let deletion_vec = deletion_vec?;
872 let row_id_sequence = row_id_sequence?;
873
874 if opened_files.is_empty() && !read_config.has_system_cols() {
875 return Err(Error::not_found(format!(
876 "No data files found for schema: {}, fragment_id={}",
877 projection,
878 self.id()
879 )));
880 }
881
882 let num_physical_rows = self.physical_rows().await?;
883
884 let mut reader = FragmentReader::try_new(
885 self.id(),
886 deletion_vec,
887 row_id_sequence,
888 opened_files,
889 ArrowSchema::from(projection),
890 self.count_rows(None).await?,
891 num_physical_rows,
892 Arc::new(self.metadata.clone()),
893 )?;
894
895 if read_config.with_row_id {
896 reader.with_row_id();
897 }
898 if read_config.with_row_address {
899 reader.with_row_address();
900 }
901 if read_config.with_row_last_updated_at_version {
902 reader.with_row_last_updated_at_version();
903 }
904 if read_config.with_row_created_at_version {
905 reader.with_row_created_at_version();
906 }
907
908 Ok(reader)
909 }
910
911 fn get_field_id_offset(data_file: &DataFile) -> u32 {
912 data_file.fields.first().copied().unwrap_or(0) as u32
913 }
914
915 async fn open_reader(
916 &self,
917 data_file: &DataFile,
918 projection: Option<&Schema>,
919 read_config: &FragReadConfig,
920 ) -> Result<Option<Box<dyn GenericFileReader>>> {
921 let full_schema = self.dataset.schema();
922 let data_file_schema = data_file.schema(full_schema);
924 let projection = projection.unwrap_or(full_schema);
925 let schema_per_file = Arc::new(projection.intersection_ignore_types(&data_file_schema)?);
927
928 if data_file.is_legacy_file() {
929 let max_field_id = data_file.fields.iter().max().unwrap();
930 if !schema_per_file.fields.is_empty() {
931 let path = self
932 .dataset
933 .data_file_dir(data_file)?
934 .child(data_file.path.as_str());
935 let field_id_offset = Self::get_field_id_offset(data_file);
936 let reader = PreviousFileReader::try_new_with_fragment_id(
937 &self.dataset.object_store,
938 &path,
939 self.schema().clone(),
940 self.id() as u32,
941 field_id_offset as i32,
942 *max_field_id,
943 Some(&self.dataset.metadata_cache.file_metadata_cache(&path)),
944 )
945 .await?;
946 let initialized_schema = reader.schema().project_by_schema(
947 schema_per_file.as_ref(),
948 OnMissing::Error,
949 OnTypeMismatch::Error,
950 )?;
951 let reader = V1Reader::new(reader, Arc::new(initialized_schema));
952 Ok(Some(Box::new(reader)))
953 } else {
954 Ok(None)
955 }
956 } else if schema_per_file.fields.is_empty() {
957 Ok(None)
958 } else {
959 let path = self
960 .dataset
961 .data_file_dir(data_file)?
962 .child(data_file.path.as_str());
963 let (store_scheduler, reader_priority) = if let Some(base_id) = data_file.base_id {
964 let object_store = self.dataset.object_store_for_base(base_id).await?;
967 let config = SchedulerConfig::max_bandwidth(&object_store);
968 (
969 ScanScheduler::new(object_store, config),
970 read_config.reader_priority.unwrap_or(0),
971 )
972 } else if let Some(scan_scheduler) = read_config.scan_scheduler.as_ref() {
973 (
974 scan_scheduler.clone(),
975 read_config.reader_priority.unwrap_or(0),
976 )
977 } else {
978 (
979 ScanScheduler::new(
980 self.dataset.object_store.clone(),
981 SchedulerConfig::max_bandwidth(&self.dataset.object_store),
982 ),
983 0,
984 )
985 };
986 let file_scheduler = store_scheduler
987 .open_file_with_priority(&path, reader_priority as u64, &data_file.file_size_bytes)
988 .await?;
989 let file_metadata = self.get_file_metadata(&file_scheduler).await?;
990 let path = file_scheduler.reader().path().clone();
991 let metadata_cache = self.dataset.metadata_cache.file_metadata_cache(&path);
992 let reader = Arc::new(
993 lance_file::reader::FileReader::try_open_with_file_metadata(
994 Arc::new(LanceEncodingsIo::new(file_scheduler.clone())),
995 path,
996 None,
997 Arc::<DecoderPlugins>::default(),
998 file_metadata,
999 &metadata_cache,
1000 read_config
1001 .file_reader_options
1002 .clone()
1003 .or_else(|| self.dataset.file_reader_options.clone())
1004 .unwrap_or_default(),
1005 )
1006 .await?,
1007 );
1008 let field_id_to_column_idx = Arc::new(BTreeMap::from_iter(
1009 data_file
1010 .fields
1011 .iter()
1012 .copied()
1013 .zip(data_file.column_indices.iter().copied())
1014 .filter_map(|(field_id, column_index)| {
1015 if column_index < 0 {
1016 None
1017 } else {
1018 Some((field_id as u32, column_index as u32))
1019 }
1020 }),
1021 ));
1022 let reader = v2_adapter::Reader::new(
1023 reader,
1024 schema_per_file,
1025 field_id_to_column_idx,
1026 reader_priority,
1027 file_scheduler,
1028 );
1029 Ok(Some(Box::new(reader)))
1030 }
1031 }
1032
1033 async fn open_readers(
1034 &self,
1035 projection: &Schema,
1036 read_config: &FragReadConfig,
1037 ) -> Result<Vec<Box<dyn GenericFileReader>>> {
1038 let mut opened_files = vec![];
1039 for data_file in &self.metadata.files {
1040 if let Some(reader) = self
1041 .open_reader(data_file, Some(projection), read_config)
1042 .await?
1043 {
1044 opened_files.push(reader);
1045 }
1046 }
1047
1048 let num_rows = self.physical_rows().await?;
1051
1052 let field_ids_in_files = opened_files
1054 .iter()
1055 .flat_map(|r| r.projection().fields_pre_order().map(|f| f.id))
1056 .filter(|id| *id >= 0)
1057 .collect::<HashSet<_>>();
1058 let mut missing_fields = projection.field_ids();
1059 missing_fields.retain(|f| !field_ids_in_files.contains(f) && *f >= 0);
1060 if !missing_fields.is_empty() {
1061 let missing_projection = projection.project_by_ids(&missing_fields, true);
1062 let null_reader = NullReader::new(Arc::new(missing_projection), num_rows as u32);
1063 opened_files.push(Box::new(null_reader));
1064 }
1065
1066 Ok(opened_files)
1067 }
1068
1069 pub async fn count_rows(&self, filter: Option<String>) -> Result<usize> {
1071 match filter {
1072 Some(expr) => self
1073 .scan()
1074 .project(&Vec::<String>::default())
1075 .unwrap()
1076 .with_row_id()
1077 .filter(&expr)?
1078 .count_rows()
1079 .await
1080 .map(|v| v as usize),
1081 None => {
1082 let total_rows = self.physical_rows();
1083 let deletion_count = self.count_deletions();
1084
1085 let (total_rows, deletion_count) =
1086 futures::future::try_join(total_rows, deletion_count).await?;
1087
1088 Ok(total_rows - deletion_count)
1089 }
1090 }
1091 }
1092
1093 pub async fn count_deletions(&self) -> Result<usize> {
1095 match &self.metadata().deletion_file {
1096 Some(DeletionFile {
1097 num_deleted_rows: Some(num_deleted),
1098 ..
1099 }) => Ok(*num_deleted),
1100 _ => {
1101 let deleletion_vector = self.get_deletion_vector().await?;
1102 if let Some(deletion_vector) = deleletion_vector {
1103 Ok(deletion_vector.len())
1104 } else {
1105 Ok(0)
1106 }
1107 }
1108 }
1109 }
1110
1111 pub fn fast_physical_rows(&self) -> Result<usize> {
1116 if self.dataset.manifest.writer_version.is_some() && self.metadata.physical_rows.is_some() {
1117 Ok(self.metadata.physical_rows.unwrap())
1118 } else {
1119 Err(Error::internal(format!(
1120 "The method fast_physical_rows was called on a fragment that does not have the physical row count in the metadata. Fragment id: {}",
1121 self.id()
1122 )))
1123 }
1124 }
1125
1126 pub fn fast_num_deletions(&self) -> Result<usize> {
1131 match &self.metadata().deletion_file {
1132 Some(DeletionFile {
1133 num_deleted_rows: Some(num_deleted),
1134 ..
1135 }) => Ok(*num_deleted),
1136 None => Ok(0),
1137 _ => Err(Error::internal(format!(
1138 "The method fast_num_deletions was called on a fragment that does not have the deletion count in the metadata. Fragment id: {}",
1139 self.id()
1140 ))),
1141 }
1142 }
1143
1144 pub fn fast_logical_rows(&self) -> Result<usize> {
1149 let num_physical_rows = self.fast_physical_rows()?;
1150 let num_deleted_rows = self.fast_num_deletions()?;
1151 Ok(num_physical_rows - num_deleted_rows)
1152 }
1153
1154 pub async fn physical_rows(&self) -> Result<usize> {
1159 if self.metadata.files.is_empty() {
1160 return Err(Error::not_found(format!(
1161 "Fragment {} does not contain any data",
1162 self.id()
1163 )));
1164 };
1165
1166 if self.dataset.manifest.writer_version.is_some() && self.metadata.physical_rows.is_some() {
1172 return Ok(self.metadata.physical_rows.unwrap());
1173 }
1174
1175 let some_file = &self.metadata.files[0];
1177 let reader = self
1178 .open_reader(some_file, None, &FragReadConfig::default())
1179 .await?
1180 .ok_or_else(|| {
1181 Error::internal(format!(
1182 "The data file {} did not have any fields contained in the dataset schema",
1183 some_file.path
1184 ))
1185 })?;
1186
1187 Ok(reader.len() as usize)
1188 }
1189
1190 pub async fn validate(&self) -> Result<()> {
1201 let mut seen_fields = HashSet::new();
1202 for data_file in &self.metadata.files {
1203 let last = -1;
1204 for field_id in &data_file.fields {
1205 if *field_id <= last {
1206 return Err(Error::corrupt_file(
1207 self.dataset
1208 .data_file_dir(data_file)?
1209 .child(data_file.path.as_str()),
1210 format!(
1211 "Field id {} is not in increasing order in fragment {:#?}",
1212 field_id, self
1213 ),
1214 ));
1215 }
1216
1217 if !seen_fields.insert(field_id) {
1218 return Err(Error::corrupt_file(
1219 self.dataset
1220 .data_file_dir(data_file)?
1221 .child(data_file.path.as_str()),
1222 format!(
1223 "Field id {} is duplicated in fragment {:#?}",
1224 field_id, self
1225 ),
1226 ));
1227 }
1228 }
1229 }
1230
1231 if self.metadata.files.iter().any(|f| f.is_legacy_file())
1232 != self.metadata.files.iter().all(|f| f.is_legacy_file())
1233 {
1234 return Err(Error::corrupt_file(
1235 self.dataset
1236 .data_file_dir(&self.metadata.files[0])?
1237 .child(self.metadata.files[0].path.as_str()),
1238 "Fragment contains a mix of v1 and v2 data files".to_string(),
1239 ));
1240 }
1241
1242 for data_file in &self.metadata.files {
1243 data_file.validate(&self.dataset.data_file_dir(&self.metadata.files[0])?)?;
1244 }
1245
1246 let get_lengths = self.metadata.files.iter().map(|data_file| async move {
1247 let data_file_dir = self.dataset.data_file_dir(data_file)?;
1248 let reader = self
1249 .open_reader(data_file, None, &FragReadConfig::default())
1250 .await?
1251 .ok_or_else(|| {
1252 Error::corrupt_file(
1253 data_file_dir.child(data_file.path.as_str()),
1254 "did not have any fields in common with the dataset schema",
1255 )
1256 })?;
1257 Result::Ok(reader.len() as usize)
1258 });
1259 let get_lengths = try_join_all(get_lengths);
1260
1261 let deletion_vector = self.get_deletion_vector();
1262
1263 let (get_lengths, deletion_vector) = join!(get_lengths, deletion_vector);
1264
1265 let get_lengths = get_lengths?;
1266 let expected_length = get_lengths.first().unwrap_or(&0);
1267 for (length, data_file) in get_lengths.iter().zip(self.metadata.files.iter()) {
1268 if length != expected_length {
1269 let path = self
1270 .dataset
1271 .data_file_dir(data_file)?
1272 .child(data_file.path.as_str());
1273 return Err(Error::corrupt_file(
1274 path,
1275 format!(
1276 "data file has incorrect length. Expected: {} Got: {}",
1277 expected_length, length
1278 ),
1279 ));
1280 }
1281 }
1282 if let Some(physical_rows) = self.metadata.physical_rows
1283 && physical_rows != *expected_length
1284 {
1285 return Err(Error::corrupt_file(
1286 self.dataset
1287 .data_file_dir(&self.metadata.files[0])?
1288 .child(self.metadata.files[0].path.as_str()),
1289 format!(
1290 "Fragment metadata has incorrect physical_rows. Actual: {} Metadata: {}",
1291 expected_length, physical_rows
1292 ),
1293 ));
1294 }
1295
1296 if let Some(deletion_vector) = deletion_vector? {
1297 if let Some(num_deletions) = self
1298 .metadata
1299 .deletion_file
1300 .as_ref()
1301 .unwrap()
1302 .num_deleted_rows
1303 && num_deletions != deletion_vector.len()
1304 {
1305 return Err(Error::corrupt_file(
1306 deletion_file_path(
1307 &self.dataset.base,
1308 self.metadata.id,
1309 self.metadata.deletion_file.as_ref().unwrap(),
1310 ),
1311 format!(
1312 "deletion vector length does not match metadata. Metadata: {} Deletion vector: {}",
1313 num_deletions,
1314 deletion_vector.len()
1315 ),
1316 ));
1317 }
1318
1319 for offset in deletion_vector.iter() {
1320 if offset >= *expected_length as u32 {
1321 let deletion_file_meta = self.metadata.deletion_file.as_ref().unwrap();
1322 return Err(Error::corrupt_file(
1323 deletion_file_path(
1324 &self.dataset.base,
1325 self.metadata.id,
1326 deletion_file_meta,
1327 ),
1328 format!(
1329 "deletion vector contains an offset that is out of range. Offset: {} Fragment length: {}",
1330 offset, expected_length
1331 ),
1332 ));
1333 }
1334 }
1335 }
1336
1337 Ok(())
1338 }
1339
1340 pub async fn open_session(
1344 &self,
1345 projection: &Schema,
1346 with_row_address: bool,
1347 ) -> Result<FragmentSession> {
1348 FragmentSession::open(Arc::new(self.clone()), projection, with_row_address).await
1349 }
1350
1351 pub async fn take(&self, indices: &[u32], projection: &Schema) -> Result<RecordBatch> {
1356 let deletion_vector = self.get_deletion_vector().await?;
1358 let row_ids = if let Some(deletion_vector) = deletion_vector {
1359 let mut sorted_deleted_ids = deletion_vector
1363 .as_ref()
1364 .clone()
1365 .into_iter()
1366 .collect::<Vec<_>>();
1367 sorted_deleted_ids.sort();
1368
1369 Cow::Owned(resolve_actual_row_ids(indices, &sorted_deleted_ids))
1370 } else {
1371 Cow::Borrowed(indices)
1372 };
1373
1374 self.take_rows(&row_ids, projection, false, false, false, false)
1376 .await
1377 }
1378
1379 pub async fn get_deletion_vector(&self) -> Result<Option<Arc<DeletionVector>>> {
1381 let Some(deletion_file) = self.metadata.deletion_file.as_ref() else {
1382 return Ok(None);
1383 };
1384
1385 let deletion_vector =
1386 read_dataset_deletion_file(&self.dataset, self.id() as u64, deletion_file).await?;
1387
1388 Ok(Some(deletion_vector))
1389 }
1390
1391 async fn get_file_metadata(
1393 &self,
1394 file_scheduler: &FileScheduler,
1395 ) -> Result<Arc<CachedFileMetadata>> {
1396 let path = file_scheduler.reader().path();
1397 let cache = self.dataset.metadata_cache.file_metadata_cache(path);
1398
1399 let file_metadata = cache
1400 .get_or_insert_with_key(FileMetadataCacheKey, || async {
1401 let file_metadata: CachedFileMetadata =
1402 lance_file::reader::FileReader::read_all_metadata(file_scheduler).await?;
1403 Ok(file_metadata)
1404 })
1405 .await?;
1406 Ok(file_metadata)
1407 }
1408
1409 pub(crate) async fn take_rows(
1419 &self,
1420 row_offsets: &[u32],
1421 projection: &Schema,
1422 with_row_id: bool,
1423 with_row_address: bool,
1424 with_row_created_at_version: bool,
1425 with_row_last_updated_at_version: bool,
1426 ) -> Result<RecordBatch> {
1427 let reader = self
1428 .open(
1429 projection,
1430 FragReadConfig::default()
1431 .with_row_id(with_row_id)
1432 .with_row_address(with_row_address)
1433 .with_row_created_at_version(with_row_created_at_version)
1434 .with_row_last_updated_at_version(with_row_last_updated_at_version),
1435 )
1436 .await?;
1437
1438 if row_offsets.len() > 1 && Self::row_ids_contiguous(row_offsets) {
1439 let range =
1440 (row_offsets[0] as usize)..(row_offsets[row_offsets.len() - 1] as usize + 1);
1441 reader.legacy_read_range_as_batch(range).await
1442 } else {
1443 reader.take_as_batch(row_offsets, None).await
1445 }
1446 }
1447
1448 fn row_ids_contiguous(row_ids: &[u32]) -> bool {
1449 if row_ids.is_empty() {
1450 return false;
1451 }
1452
1453 let mut last_id = row_ids[0];
1454
1455 for id in row_ids.iter().skip(1) {
1456 if *id != last_id + 1 {
1457 return false;
1458 }
1459 last_id = *id;
1460 }
1461
1462 true
1463 }
1464
1465 pub fn scan(&self) -> Scanner {
1469 Scanner::from_fragment(self.dataset.clone(), self.metadata.clone())
1470 }
1471
1472 pub(crate) async fn updater<T: AsRef<str>>(
1490 &self,
1491 columns: Option<&[T]>,
1492 schemas: Option<(Schema, Schema)>,
1493 batch_size: Option<u32>,
1494 ) -> Result<Updater> {
1495 let mut schema = self.dataset.schema().clone();
1496
1497 let mut with_row_addr = false;
1498 let mut with_row_id = false;
1499 if let Some(columns) = columns {
1500 let mut projection = Vec::new();
1501 for column in columns {
1502 if column.as_ref() == ROW_ADDR {
1503 with_row_addr = true;
1504 } else if column.as_ref() == ROW_ID {
1505 with_row_id = true;
1506 } else {
1507 projection.push(column.as_ref());
1508 }
1509 }
1510 schema = schema.project(&projection)?;
1511 }
1512
1513 with_row_addr |= !with_row_id && schema.fields.is_empty();
1515
1516 let reader = self.open(
1517 &schema,
1518 FragReadConfig::default()
1519 .with_row_address(with_row_addr)
1520 .with_row_id(with_row_id),
1521 );
1522 let deletion_vector = self.get_deletion_vector();
1523 let (reader, deletion_vector) = join!(reader, deletion_vector);
1524 let reader = reader?;
1525 let deletion_vector = deletion_vector?.unwrap_or_default().as_ref().clone();
1526
1527 Updater::try_new(self.clone(), reader, deletion_vector, schemas, batch_size)
1528 }
1529
1530 pub async fn merge_columns(
1531 &mut self,
1532 stream: impl RecordBatchReader + Send + 'static,
1533 left_on: &str,
1534 right_on: &str,
1535 max_field_id: i32,
1536 ) -> Result<(Fragment, Schema)> {
1537 let stream = Box::new(stream);
1538 if self.schema().field(left_on).is_none() && left_on != ROW_ID && left_on != ROW_ADDR {
1539 return Err(Error::invalid_input(format!(
1540 "Column {} does not exist in the left side fragment",
1541 left_on
1542 )));
1543 };
1544 let right_schema = stream.schema();
1545 if right_schema.field_with_name(right_on).is_err() {
1546 return Err(Error::invalid_input(format!(
1547 "Column {} does not exist in the right side fragment",
1548 right_on
1549 )));
1550 };
1551
1552 for field in right_schema.fields() {
1553 if field.name() == right_on {
1554 continue;
1557 }
1558 if self.schema().field(field.name()).is_some() {
1559 return Err(Error::invalid_input(format!(
1560 "Column {} exists in left side fragment and right side dataset",
1561 field.name()
1562 )));
1563 }
1564 }
1565 let joiner = Arc::new(HashJoiner::try_new(stream, right_on).await?);
1567 let mut new_schema: Schema = self.schema().merge(joiner.out_schema().as_ref())?;
1570 new_schema.set_field_id(Some(max_field_id));
1571
1572 let new_fragment = self
1573 .clone()
1574 .merge(left_on, &joiner)
1575 .await
1576 .map(|f| f.metadata)?;
1577
1578 Ok((new_fragment, new_schema))
1579 }
1580
1581 pub(crate) async fn merge(mut self, join_column: &str, joiner: &HashJoiner) -> Result<Self> {
1582 let mut updater = self.updater(Some(&[join_column]), None, None).await?;
1583
1584 while let Some(batch) = updater.next().await? {
1585 let batch = joiner
1586 .collect(&self.dataset, batch[join_column].clone())
1587 .await?;
1588 updater.update(batch).await?;
1589 }
1590
1591 self.metadata = updater.finish().await?;
1592
1593 Ok(self)
1594 }
1595
1596 pub async fn update_columns(
1597 &mut self,
1598 right_stream: impl RecordBatchReader + Send + 'static,
1599 left_on: &str,
1600 right_on: &str,
1601 ) -> Result<(Fragment, Vec<u32>)> {
1602 if self.schema().field(left_on).is_none() && left_on != ROW_ID && left_on != ROW_ADDR {
1603 return Err(Error::invalid_input(format!(
1604 "Column {} does not exist in the left side fragment",
1605 left_on
1606 )));
1607 };
1608 let right_stream = Box::new(right_stream);
1609 let right_schema = right_stream.schema();
1610 if right_schema.field_with_name(right_on).is_err() {
1611 return Err(Error::invalid_input(format!(
1612 "Column {} does not exist in the right side fragment",
1613 right_on
1614 )));
1615 };
1616 let write_schema = right_schema.as_ref().without_column(right_on);
1617 for field in write_schema.fields() {
1618 if ROW_ID.eq(field.name()) || ROW_ADDR.eq(field.name()) {
1619 return Err(Error::invalid_input(format!(
1620 "Column {} is a reversed metadata column and cannot be updated",
1621 field.name()
1622 )));
1623 }
1624 if self.schema().field(field.name()).is_none() {
1625 return Err(Error::invalid_input(format!(
1626 "Column {} in right side fragment does not exist in left side fragment",
1627 field.name()
1628 )));
1629 }
1630 }
1631
1632 let write_schema = self.schema().project_by_schema(
1633 &write_schema,
1634 OnMissing::Error,
1635 OnTypeMismatch::Error,
1636 )?;
1637 let mut read_columns: Vec<String> =
1639 write_schema.fields.iter().map(|f| f.name.clone()).collect();
1640 read_columns.push(left_on.to_string());
1641 let mut updater = self
1642 .updater(
1643 Some(&read_columns),
1644 Some((write_schema.clone(), self.schema().clone())),
1645 None,
1646 )
1647 .await?;
1648 let joiner = Arc::new(HashJoiner::try_new(right_stream, right_on).await?);
1650 while let Some(batch) = updater.next().await? {
1651 let updated_batch = joiner
1652 .collect_with_fallback(batch, batch[left_on].clone(), self.dataset())
1653 .await?;
1654 updater.update(updated_batch).await?;
1655 }
1656
1657 let mut updated_fragment = updater.finish().await?;
1658 let updated_fields = updated_fragment.files.last().unwrap().fields.clone();
1660 for data_file in &mut updated_fragment.files.iter_mut().rev().skip(1) {
1661 for field in &mut data_file.fields {
1662 if updated_fields.contains(field) {
1663 *field = -2;
1665 }
1666 }
1667 }
1668 updated_fragment
1670 .files
1671 .retain(|data_file| data_file.fields.iter().any(|&field| field != -2));
1672 let updated_fields = updated_fields
1673 .iter()
1674 .filter_map(|&i| u32::try_from(i).ok())
1675 .collect();
1676 Ok((updated_fragment, updated_fields))
1678 }
1679
1680 pub async fn add_columns(
1684 &self,
1685 transforms: NewColumnTransform,
1686 read_columns: Option<Vec<String>>,
1687 batch_size: Option<u32>,
1688 ) -> Result<(Fragment, Schema)> {
1689 let (fragments, schema) = schema_evolution::add_columns_to_fragments(
1690 self.dataset.as_ref(),
1691 transforms,
1692 read_columns,
1693 std::slice::from_ref(self),
1694 batch_size,
1695 )
1696 .await?;
1697 assert_eq!(fragments.len(), 1);
1698 Ok((fragments.into_iter().next().unwrap(), schema))
1699 }
1700
1701 pub async fn delete(self, predicate: &str) -> Result<Option<Self>> {
1707 let mut deletion_vector = self
1709 .get_deletion_vector()
1710 .await?
1711 .unwrap_or_default()
1712 .as_ref()
1713 .clone();
1714
1715 let starting_length = deletion_vector.len();
1716
1717 let mut scanner = self.scan();
1719
1720 let predicate_lower = predicate.trim().to_lowercase();
1721 if predicate_lower == "true" {
1722 return Ok(None);
1723 } else if predicate_lower == "false" {
1724 return Ok(Some(self));
1725 }
1726
1727 scanner
1728 .with_row_address()
1729 .filter(predicate)?
1730 .project::<&str>(&[])?;
1731
1732 if let Some(predicate) = &scanner.get_expr_filter()? {
1737 if matches!(
1738 predicate,
1739 Expr::Literal(ScalarValue::Boolean(Some(false)), _)
1740 ) {
1741 return Ok(Some(self));
1742 }
1743 if matches!(
1744 predicate,
1745 Expr::Literal(ScalarValue::Boolean(Some(true)), _)
1746 ) {
1747 return Ok(None);
1748 }
1749 }
1750
1751 scanner
1753 .try_into_stream()
1754 .await?
1755 .try_for_each(|batch| {
1756 let array = batch[ROW_ADDR].clone();
1757 let int_array: &UInt64Array = as_primitive_array(array.as_ref());
1758
1759 let local_row_ids = int_array.values().iter().map(|v| *v as u32);
1763
1764 deletion_vector.extend(local_row_ids);
1765 futures::future::ready(Ok(()))
1766 })
1767 .await?;
1768
1769 if deletion_vector.len() == starting_length {
1771 return Ok(Some(self));
1772 }
1773
1774 self.write_deletions(deletion_vector).await
1775 }
1776
1777 pub async fn extend_deletions(
1778 self,
1779 new_deletions: impl IntoIterator<Item = u32>,
1780 ) -> Result<Option<Self>> {
1781 let mut deletion_vector = self
1782 .get_deletion_vector()
1783 .await?
1784 .unwrap_or_default()
1785 .as_ref()
1786 .clone();
1787
1788 deletion_vector.extend(new_deletions);
1789
1790 self.write_deletions(deletion_vector).await
1791 }
1792
1793 async fn write_deletions(mut self, deletion_vector: DeletionVector) -> Result<Option<Self>> {
1794 let physical_rows = self.physical_rows().await?;
1795 if deletion_vector.len() == physical_rows
1796 && deletion_vector.contains_range(0..physical_rows as u32)
1797 {
1798 return Ok(None);
1799 } else if deletion_vector.len() >= physical_rows {
1800 let dv_len = deletion_vector.len();
1801 let examples: Vec<u32> = deletion_vector
1802 .into_iter()
1803 .filter(|x| *x >= physical_rows as u32)
1804 .take(5)
1805 .collect();
1806 return Err(Error::internal(format!(
1807 "Deletion vector includes rows that aren't in the fragment. \
1808 Num physical rows {}; Deletion vector length: {}; \
1809 Examples: {:?}",
1810 physical_rows, dv_len, examples
1811 )));
1812 }
1813
1814 self.metadata.deletion_file = write_deletion_file(
1815 &self.dataset.base,
1816 self.metadata.id,
1817 self.dataset.version().version,
1818 &deletion_vector,
1819 self.dataset.object_store(),
1820 )
1821 .await?;
1822
1823 Ok(Some(self))
1824 }
1825}
1826
1827pub(crate) fn resolve_actual_row_ids(row_ids: &[u32], sorted_deleted_ids: &[u32]) -> Vec<u32> {
1829 let mut row_ids = row_ids.to_vec();
1830 for row_id in row_ids.iter_mut() {
1831 let mut new_row_id = *row_id;
1840 let offset = sorted_deleted_ids.partition_point(|v| *v <= new_row_id);
1841
1842 let mut deletion_i = offset;
1843 let mut i = 0;
1844 while i < offset {
1845 new_row_id += 1;
1847 while deletion_i < sorted_deleted_ids.len()
1848 && sorted_deleted_ids[deletion_i] == new_row_id
1849 {
1850 deletion_i += 1;
1853 new_row_id += 1;
1854 }
1855 i += 1;
1856 }
1857
1858 *row_id = new_row_id;
1859 }
1860
1861 row_ids
1862}
1863
1864#[derive(Debug, Clone)]
1866struct FileMetadataCacheKey;
1867
1868impl CacheKey for FileMetadataCacheKey {
1869 type ValueType = CachedFileMetadata;
1870
1871 fn key(&self) -> std::borrow::Cow<'_, str> {
1872 "".into()
1873 }
1874}
1875
1876impl From<FileFragment> for Fragment {
1877 fn from(fragment: FileFragment) -> Self {
1878 fragment.metadata
1879 }
1880}
1881
1882#[derive(Debug)]
1887pub struct FragmentReader {
1888 readers: Vec<Box<dyn GenericFileReader>>,
1890
1891 output_schema: ArrowSchema,
1893
1894 deletion_vec: Option<Arc<DeletionVector>>,
1896
1897 row_id_sequence: Option<Arc<RowIdSequence>>,
1901
1902 fragment_id: usize,
1904
1905 with_row_id: bool,
1907
1908 with_row_addr: bool,
1910
1911 with_row_last_updated_at_version: bool,
1913
1914 with_row_created_at_version: bool,
1916
1917 make_deletions_null: bool,
1920
1921 fragment: Arc<Fragment>,
1923
1924 last_updated_at_sequence: Option<Arc<lance_table::rowids::version::RowDatasetVersionSequence>>,
1926
1927 created_at_sequence: Option<Arc<lance_table::rowids::version::RowDatasetVersionSequence>>,
1929
1930 num_rows: usize,
1932
1933 num_physical_rows: usize,
1935}
1936
1937impl Clone for FragmentReader {
1942 fn clone(&self) -> Self {
1943 Self {
1944 readers: self
1945 .readers
1946 .iter()
1947 .map(|reader| reader.clone_box())
1948 .collect::<Vec<_>>(),
1949 output_schema: self.output_schema.clone(),
1950 deletion_vec: self.deletion_vec.clone(),
1951 row_id_sequence: self.row_id_sequence.clone(),
1952 fragment_id: self.fragment_id,
1953 with_row_id: self.with_row_id,
1954 with_row_addr: self.with_row_addr,
1955 with_row_last_updated_at_version: self.with_row_last_updated_at_version,
1956 with_row_created_at_version: self.with_row_created_at_version,
1957 make_deletions_null: self.make_deletions_null,
1958 fragment: self.fragment.clone(),
1959 last_updated_at_sequence: self.last_updated_at_sequence.clone(),
1960 created_at_sequence: self.created_at_sequence.clone(),
1961 num_rows: self.num_rows,
1962 num_physical_rows: self.num_physical_rows,
1963 }
1964 }
1965}
1966
1967impl std::fmt::Display for FragmentReader {
1968 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
1969 write!(f, "FragmentReader(id={})", self.fragment_id)
1970 }
1971}
1972
1973fn merge_batches(batches: &[RecordBatch]) -> Result<RecordBatch> {
1974 if batches.is_empty() {
1975 return Err(Error::invalid_input(
1976 "Cannot merge empty batches".to_string(),
1977 ));
1978 }
1979
1980 let mut merged = batches[0].clone();
1981 for batch in batches.iter().skip(1) {
1982 merged = merged.merge(batch)?;
1983 }
1984 Ok(merged)
1985}
1986
1987impl FragmentReader {
1988 #[allow(clippy::too_many_arguments)]
1989 fn try_new(
1990 fragment_id: usize,
1991 deletion_vec: Option<Arc<DeletionVector>>,
1992 row_id_sequence: Option<Arc<RowIdSequence>>,
1993 readers: Vec<Box<dyn GenericFileReader>>,
1994 output_schema: ArrowSchema,
1995 num_rows: usize,
1996 num_physical_rows: usize,
1997 fragment: Arc<Fragment>,
1998 ) -> Result<Self> {
1999 if let Some(legacy_reader) = readers.first().and_then(|reader| reader.as_legacy_opt()) {
2000 let num_batches = legacy_reader.num_batches();
2001 for reader in readers.iter().skip(1) {
2002 if let Some(other_legacy) = reader.as_legacy_opt() {
2003 if other_legacy.num_batches() != num_batches {
2004 return Err(Error::invalid_input("Cannot create FragmentReader from data files with different number of batches"
2005 .to_string()));
2006 }
2007 } else {
2008 return Err(Error::invalid_input(
2009 "Cannot mix legacy and non-legacy readers".to_string(),
2010 ));
2011 }
2012 }
2013 }
2014 Ok(Self {
2015 readers,
2016 output_schema,
2017 deletion_vec,
2018 row_id_sequence,
2019 fragment_id,
2020 with_row_id: false,
2021 with_row_addr: false,
2022 with_row_last_updated_at_version: false,
2023 with_row_created_at_version: false,
2024 make_deletions_null: false,
2025 fragment,
2026 last_updated_at_sequence: None,
2027 created_at_sequence: None,
2028 num_rows,
2029 num_physical_rows,
2030 })
2031 }
2032
2033 pub(crate) fn with_row_id(&mut self) -> &mut Self {
2034 self.with_row_id = true;
2035 self.output_schema = self
2036 .output_schema
2037 .try_with_column(ROW_ID_FIELD.clone())
2038 .expect("Table already has a column named _rowid");
2039 self
2040 }
2041
2042 pub(crate) fn with_row_address(&mut self) -> &mut Self {
2043 self.with_row_addr = true;
2044 self.output_schema = self
2045 .output_schema
2046 .try_with_column(ROW_ADDR_FIELD.clone())
2047 .expect("Table already has a column named _rowaddr");
2048 self
2049 }
2050
2051 pub(crate) fn with_make_deletions_null(&mut self) -> &mut Self {
2052 self.make_deletions_null = true;
2053 self
2054 }
2055
2056 pub(crate) fn with_row_last_updated_at_version(&mut self) -> &mut Self {
2057 self.with_row_last_updated_at_version = true;
2058
2059 if self.last_updated_at_sequence.is_none()
2061 && let Some(meta) = &self.fragment.last_updated_at_version_meta
2062 && let Ok(sequence) = meta.load_sequence()
2063 {
2064 self.last_updated_at_sequence = Some(Arc::new(sequence));
2065 }
2066 self.output_schema = self
2070 .output_schema
2071 .try_with_column(ROW_LAST_UPDATED_AT_VERSION_FIELD.clone())
2072 .expect("Table already has a column named _row_last_updated_at_version");
2073
2074 self
2075 }
2076
2077 pub(crate) fn with_row_created_at_version(&mut self) -> &mut Self {
2078 self.with_row_created_at_version = true;
2079
2080 if self.created_at_sequence.is_none()
2082 && let Some(meta) = &self.fragment.created_at_version_meta
2083 && let Ok(sequence) = meta.load_sequence()
2084 {
2085 self.created_at_sequence = Some(Arc::new(sequence));
2086 }
2087 self.output_schema = self
2091 .output_schema
2092 .try_with_column(ROW_CREATED_AT_VERSION_FIELD.clone())
2093 .expect("Table already has a column named _row_created_at_version");
2094
2095 self
2096 }
2097
2098 pub(crate) fn legacy_num_batches(&self) -> usize {
2102 let legacy_reader = self.readers[0].as_legacy();
2103 let num_batches = legacy_reader.num_batches();
2104 assert!(
2105 self.readers
2106 .iter()
2107 .all(|r| r.as_legacy().num_batches() == num_batches),
2108 "Data files have varying number of batches, which is not yet supported."
2109 );
2110 num_batches
2111 }
2112
2113 pub(crate) fn legacy_num_rows_in_batch(&self, batch_id: u32) -> Option<u32> {
2121 if let Some(legacy_reader) = self.readers.first().and_then(|r| r.as_legacy_opt()) {
2122 if batch_id < legacy_reader.num_batches() as u32 {
2123 Some(legacy_reader.num_rows_in_batch(batch_id as i32) as u32)
2124 } else {
2125 None
2126 }
2127 } else {
2128 None
2129 }
2130 }
2131
2132 pub(crate) async fn legacy_read_page_stats(
2138 &self,
2139 projection: Option<&Schema>,
2140 ) -> Result<Option<RecordBatch>> {
2141 let mut stats_batches = vec![];
2142 for reader in self.readers.iter() {
2143 let schema = match projection {
2144 Some(projection) => Arc::new(reader.projection().intersection(projection)?),
2145 None => reader.projection().clone(),
2146 };
2147 let reader = reader.as_legacy();
2148 if let Some(stats_batch) = reader.read_page_stats(&schema.field_ids()).await? {
2149 stats_batches.push(stats_batch);
2150 }
2151 }
2152
2153 if stats_batches.is_empty() {
2154 Ok(None)
2155 } else {
2156 Ok(Some(merge_batches(&stats_batches)?))
2157 }
2158 }
2159
2160 pub(crate) async fn legacy_read_batch_projected(
2169 &self,
2170 batch_id: usize,
2171 params: impl Into<ReadBatchParams> + Clone,
2172 projection: &Schema,
2173 ) -> Result<RecordBatch> {
2174 let first_reader = self.readers[0].as_legacy();
2175 let batch_offset = batch_id * first_reader.num_rows_in_batch(0);
2177 let rows_in_batch = first_reader.num_rows_in_batch(batch_id as i32);
2178
2179 let batches = if !projection.fields.is_empty() {
2180 let read_tasks = self.readers.iter().map(|reader| {
2181 let projection = reader.projection().intersection(projection);
2182 let params = params.clone();
2183
2184 let reader = reader.as_legacy();
2185
2186 async move {
2187 let projection = projection?;
2190 if projection.fields.is_empty() {
2191 Result::Ok(None)
2194 } else {
2195 Ok(Some(
2196 reader
2197 .read_batch(batch_id as i32, params, &projection)
2198 .await?,
2199 ))
2200 }
2201 }
2202 });
2203 let results = try_join_all(read_tasks).await?;
2204 results.into_iter().flatten().collect::<Vec<RecordBatch>>()
2205 } else {
2206 let expected_rows = params
2210 .clone()
2211 .into()
2212 .slice(0, rows_in_batch)
2213 .unwrap()
2214 .to_offsets()?
2215 .len();
2216 vec![RecordBatch::from(StructArray::new_empty_fields(
2217 expected_rows,
2218 None,
2219 ))]
2220 };
2221
2222 let params = params.into();
2223 let result = merge_batches(&batches)?;
2224
2225 let file_params = match params {
2229 ReadBatchParams::Indices(indices) => ReadBatchParams::Indices(
2230 indices
2231 .values()
2232 .iter()
2233 .map(|i| *i + batch_offset as u32)
2234 .collect(),
2235 ),
2236 ReadBatchParams::Ranges(_) => {
2237 return Err(Error::internal(
2238 "ReadBatchParams::Ranges should not be used in v1 files".to_string(),
2239 ));
2240 }
2241 ReadBatchParams::RangeFull => {
2242 ReadBatchParams::Range(batch_offset..(batch_offset + rows_in_batch))
2243 }
2244 ReadBatchParams::RangeFrom(start) => {
2245 ReadBatchParams::Range((start.start + batch_offset)..(batch_offset + rows_in_batch))
2246 }
2247 ReadBatchParams::RangeTo(end) => {
2248 ReadBatchParams::Range(batch_offset..(end.end + batch_offset))
2249 }
2250 ReadBatchParams::Range(range) => {
2251 ReadBatchParams::Range((range.start + batch_offset)..(range.end + batch_offset))
2252 }
2253 };
2254 let result = lance_table::utils::stream::apply_row_id_and_deletes(
2255 result,
2256 0,
2257 self.fragment_id as u32,
2258 &RowIdAndDeletesConfig {
2259 params: file_params,
2260 deletion_vector: self.deletion_vec.clone(),
2261 row_id_sequence: self.row_id_sequence.clone(),
2262 with_row_id: self.with_row_id,
2263 with_row_addr: self.with_row_addr,
2264 with_row_last_updated_at_version: self.with_row_last_updated_at_version,
2265 with_row_created_at_version: self.with_row_created_at_version,
2266 last_updated_at_sequence: self.last_updated_at_sequence.clone(),
2267 created_at_sequence: self.created_at_sequence.clone(),
2268 make_deletions_null: self.make_deletions_null,
2269 total_num_rows: first_reader.len() as u32,
2270 },
2271 )?;
2272
2273 let output_schema = {
2274 let mut output_schema = ArrowSchema::from(projection);
2275 if self.with_row_id {
2276 output_schema = output_schema.try_with_column(ROW_ID_FIELD.clone())?;
2277 }
2278 if self.with_row_addr {
2279 output_schema = output_schema.try_with_column(ROW_ADDR_FIELD.clone())?;
2280 }
2281 output_schema
2282 };
2283
2284 Ok(result.project_by_schema(&output_schema)?)
2285 }
2286
2287 fn new_read_impl(
2288 &self,
2289 params: ReadBatchParams,
2290 batch_size: u32,
2291 read_fn: impl Fn(&dyn GenericFileReader) -> Result<ReadBatchTaskStream>,
2292 ) -> Result<ReadBatchFutStream> {
2293 let total_num_rows = self.num_physical_rows as u32;
2294 if !params.valid_given_len(total_num_rows as usize) {
2298 return Err(Error::invalid_input(format!(
2299 "Invalid read params {} for fragment with {} addressable rows",
2300 params, total_num_rows
2301 )));
2302 }
2303 let merged = if self.num_system_cols() == self.output_schema.fields.len() {
2314 let selected_rows = params.to_offsets_total(total_num_rows).len();
2315 let tasks = (0..selected_rows)
2316 .step_by(batch_size as usize)
2317 .map(move |offset| {
2318 let num_rows = (batch_size as usize).min(selected_rows - offset);
2319 let batch = RecordBatch::from(StructArray::new_empty_fields(num_rows, None));
2320 ReadBatchTask {
2321 task: std::future::ready(Ok(batch)).boxed(),
2322 num_rows: num_rows as u32,
2323 }
2324 });
2325 stream::iter(tasks).boxed()
2326 } else {
2327 let read_streams = self
2331 .readers
2332 .iter()
2333 .filter_map(|reader| {
2334 if reader.projection().fields.is_empty() {
2339 None
2340 } else {
2341 Some(read_fn(reader.as_ref()))
2342 }
2343 })
2344 .collect::<Result<Vec<_>>>()?;
2345 lance_table::utils::stream::merge_streams(read_streams)
2347 };
2348
2349 let config = RowIdAndDeletesConfig {
2352 deletion_vector: self.deletion_vec.clone(),
2353 row_id_sequence: self.row_id_sequence.clone(),
2354 make_deletions_null: self.make_deletions_null,
2355 with_row_id: self.with_row_id,
2356 with_row_addr: self.with_row_addr,
2357 with_row_last_updated_at_version: self.with_row_last_updated_at_version,
2358 with_row_created_at_version: self.with_row_created_at_version,
2359 last_updated_at_sequence: self.last_updated_at_sequence.clone(),
2360 created_at_sequence: self.created_at_sequence.clone(),
2361 params,
2362 total_num_rows,
2363 };
2364 let output_schema = Arc::new(self.output_schema.clone());
2365 Ok(
2366 wrap_with_row_id_and_delete(merged, self.fragment_id as u32, config)
2367 .map(move |batch_fut| {
2369 let output_schema = output_schema.clone();
2370 batch_fut
2371 .map(move |batch| {
2372 batch?
2373 .project_by_schema(&output_schema)
2374 .map_err(Error::from)
2375 })
2376 .boxed()
2377 })
2378 .boxed(),
2379 )
2380 }
2381
2382 fn patch_range_for_deletions(&self, range: Range<u32>, dv: &DeletionVector) -> Range<u32> {
2383 let mut start = range.start;
2384 let mut end = range.end;
2385 for val in dv.to_sorted_iter() {
2386 if val <= start {
2387 start += 1;
2388 end += 1;
2389 } else if val < end {
2390 end += 1;
2391 } else {
2392 break;
2393 }
2394 }
2395 start..end
2396 }
2397
2398 fn do_read_range(
2399 &self,
2400 mut range: Range<u32>,
2401 batch_size: u32,
2402 skip_deleted_rows: bool,
2403 ) -> Result<ReadBatchFutStream> {
2404 if skip_deleted_rows && let Some(deletion_vector) = self.deletion_vec.as_ref() {
2405 range = self.patch_range_for_deletions(range, deletion_vector.as_ref());
2406 }
2407 self.new_read_impl(
2408 ReadBatchParams::Range(range.start as usize..range.end as usize),
2409 batch_size,
2410 move |reader| {
2411 reader.read_range_tasks(
2412 range.start as u64..range.end as u64,
2413 batch_size,
2414 reader.projection().clone(),
2415 )
2416 },
2417 )
2418 }
2419
2420 fn num_system_cols(&self) -> usize {
2421 self.with_row_id as usize
2422 + self.with_row_addr as usize
2423 + self.with_row_created_at_version as usize
2424 + self.with_row_last_updated_at_version as usize
2425 }
2426
2427 pub fn read_range(&self, range: Range<u32>, batch_size: u32) -> Result<ReadBatchFutStream> {
2432 self.do_read_range(range, batch_size, true)
2433 }
2434
2435 pub fn take_range(&self, range: Range<u32>, batch_size: u32) -> Result<ReadBatchFutStream> {
2440 self.do_read_range(range, batch_size, false)
2441 }
2442
2443 pub fn read_all(&self, batch_size: u32) -> Result<ReadBatchFutStream> {
2444 self.new_read_impl(ReadBatchParams::RangeFull, batch_size, move |reader| {
2445 reader.read_all_tasks(batch_size, reader.projection().clone())
2446 })
2447 }
2448
2449 pub fn read_ranges(
2453 &self,
2454 ranges: Arc<[Range<u64>]>,
2455 batch_size: u32,
2456 ) -> Result<ReadBatchFutStream> {
2457 let total_num_rows = self.num_physical_rows as u32;
2458 let mut num_requested_rows = 0;
2459 for range in ranges.as_ref() {
2461 if range.end > total_num_rows as u64 {
2462 return Err(Error::internal(format!(
2463 "Invalid read of range {:?} for fragment {} with {} addressable rows",
2464 range, self.fragment_id, total_num_rows
2465 )));
2466 }
2467 num_requested_rows += range.end - range.start;
2468 }
2469
2470 let merged_stream = if self.num_system_cols() == self.output_schema.fields.len() {
2471 let tasks = (0..num_requested_rows)
2472 .step_by(batch_size as usize)
2473 .map(move |offset| {
2474 let num_rows = (batch_size as u64).min(num_requested_rows - offset);
2475 let batch =
2476 RecordBatch::from(StructArray::new_empty_fields(num_rows as usize, None));
2477 ReadBatchTask {
2478 task: std::future::ready(Ok(batch)).boxed(),
2479 num_rows: num_rows as u32,
2480 }
2481 });
2482 stream::iter(tasks).boxed()
2483 } else {
2484 let read_streams = self
2488 .readers
2489 .iter()
2490 .map(|reader| {
2491 reader.read_ranges_tasks(
2492 ranges.clone(),
2493 batch_size,
2494 reader.projection().clone(),
2495 )
2496 })
2497 .collect::<Result<Vec<_>>>()?;
2498 lance_table::utils::stream::merge_streams(read_streams)
2500 };
2501
2502 let config = RowIdAndDeletesConfig {
2505 deletion_vector: self.deletion_vec.clone(),
2506 row_id_sequence: self.row_id_sequence.clone(),
2507 make_deletions_null: self.make_deletions_null,
2508 with_row_id: self.with_row_id,
2509 with_row_addr: self.with_row_addr,
2510 with_row_last_updated_at_version: self.with_row_last_updated_at_version,
2511 with_row_created_at_version: self.with_row_created_at_version,
2512 last_updated_at_sequence: self.last_updated_at_sequence.clone(),
2513 created_at_sequence: self.created_at_sequence.clone(),
2514 params: ReadBatchParams::Ranges(ranges),
2515 total_num_rows,
2516 };
2517 let output_schema = Arc::new(self.output_schema.clone());
2518 Ok(
2519 wrap_with_row_id_and_delete(merged_stream, self.fragment_id as u32, config)
2520 .map(move |batch_fut| {
2522 let output_schema = output_schema.clone();
2523 batch_fut
2524 .map(move |batch| {
2525 batch?
2526 .project_by_schema(&output_schema)
2527 .map_err(Error::from)
2528 })
2529 .boxed()
2530 })
2531 .boxed(),
2532 )
2533 }
2534
2535 pub async fn legacy_read_range_as_batch(&self, range: Range<usize>) -> Result<RecordBatch> {
2540 let batches = self
2541 .take_range(
2542 range.start as u32..range.end as u32,
2543 DEFAULT_BATCH_READ_SIZE,
2544 )?
2545 .buffered(get_num_compute_intensive_cpus())
2546 .try_collect::<Vec<_>>()
2547 .await?;
2548 concat_batches(&Arc::new(self.output_schema.clone()), batches.iter()).map_err(Error::from)
2549 }
2550
2551 pub async fn take(
2553 &self,
2554 indices: &[u32],
2555 batch_size: u32,
2556 take_priority: Option<u32>,
2557 ) -> Result<ReadBatchFutStream> {
2558 let indices_arr = UInt32Array::from(indices.to_vec());
2559 self.new_read_impl(
2560 ReadBatchParams::Indices(indices_arr),
2561 batch_size,
2562 move |reader| {
2563 reader.take_all_tasks(
2564 indices,
2565 batch_size,
2566 reader.projection().clone(),
2567 take_priority,
2568 )
2569 },
2570 )
2571 }
2572
2573 pub async fn take_as_batch(
2578 &self,
2579 indices: &[u32],
2580 take_priority: Option<u32>,
2581 ) -> Result<RecordBatch> {
2582 let has_duplicates = indices.windows(2).any(|w| w[0] == w[1]);
2585 let (unique_indices, expand_map) = if has_duplicates {
2586 let mut unique: Vec<u32> = Vec::with_capacity(indices.len());
2587 let mut mapping: Vec<u32> = Vec::with_capacity(indices.len());
2588 for &idx in indices {
2589 if unique.last() != Some(&idx) {
2590 unique.push(idx);
2591 }
2592 mapping.push((unique.len() - 1) as u32);
2593 }
2594 (Cow::Owned(unique), Some(UInt32Array::from(mapping)))
2595 } else {
2596 (Cow::Borrowed(indices), None)
2597 };
2598
2599 let batches = self
2600 .take(&unique_indices, u32::MAX, take_priority)
2601 .await?
2602 .buffered(get_num_compute_intensive_cpus())
2603 .try_collect::<Vec<_>>()
2604 .await?;
2605 let mut batch = concat_batches(&Arc::new(self.output_schema.clone()), batches.iter())?;
2606
2607 if let Some(expand_map) = expand_map {
2608 batch = arrow_select::take::take_record_batch(&batch, &expand_map)?;
2609 }
2610
2611 Ok(batch)
2612 }
2613}
2614
2615#[cfg(test)]
2616mod tests {
2617 use arrow_arith::numeric::mul;
2618 use arrow_array::{
2619 ArrayRef, BooleanArray, Int32Array, Int64Array, RecordBatchIterator, StringArray,
2620 };
2621 use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema};
2622 use lance_core::ROW_ID;
2623 use lance_core::utils::tempfile::TempStrDir;
2624 use lance_datagen::{RowCount, array, gen_batch};
2625 use lance_file::version::LanceFileVersion;
2626 use lance_file::writer::FileWriterOptions;
2627 use lance_io::{assert_io_eq, assert_io_lt, object_store::ObjectStore};
2628 use pretty_assertions::assert_eq;
2629 use rstest::rstest;
2630
2631 use super::*;
2632 use crate::{
2633 dataset::{
2634 InsertBuilder,
2635 transaction::{Operation, UpdateMode},
2636 },
2637 session::Session,
2638 utils::test::TestDatasetGenerator,
2639 };
2640
2641 async fn create_dataset(test_uri: &str, data_storage_version: LanceFileVersion) -> Dataset {
2642 let schema = Arc::new(ArrowSchema::new(vec![
2643 ArrowField::new("i", DataType::Int32, true),
2644 ArrowField::new("s", DataType::Utf8, true),
2645 ]));
2646
2647 let batches: Vec<RecordBatch> = (0..10)
2648 .map(|i| {
2649 RecordBatch::try_new(
2650 schema.clone(),
2651 vec![
2652 Arc::new(Int32Array::from_iter_values(i * 20..(i + 1) * 20)),
2653 Arc::new(StringArray::from_iter_values(
2654 (i * 20..(i + 1) * 20).map(|v| format!("s-{}", v)),
2655 )),
2656 ],
2657 )
2658 .unwrap()
2659 })
2660 .collect();
2661
2662 let write_params = WriteParams {
2663 max_rows_per_file: 40,
2664 max_rows_per_group: 10,
2665 data_storage_version: Some(data_storage_version),
2666 ..Default::default()
2667 };
2668 let batches = RecordBatchIterator::new(batches.into_iter().map(Ok), schema.clone());
2669 Dataset::write(batches, test_uri, Some(write_params))
2670 .await
2671 .unwrap();
2672
2673 Dataset::open(test_uri).await.unwrap()
2674 }
2675
2676 async fn create_dataset_v2(test_uri: &str) -> Dataset {
2677 let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new(
2678 "i",
2679 DataType::Int32,
2680 true,
2681 )]));
2682
2683 let batches: Vec<RecordBatch> = (0..10)
2684 .map(|i| {
2685 RecordBatch::try_new(
2686 schema.clone(),
2687 vec![Arc::new(Int32Array::from_iter_values(i * 20..(i + 1) * 20))],
2688 )
2689 .unwrap()
2690 })
2691 .collect();
2692
2693 let write_params = WriteParams {
2694 max_rows_per_file: 40,
2695 max_rows_per_group: 10,
2696 data_storage_version: Some(LanceFileVersion::Stable),
2697 ..Default::default()
2698 };
2699 let batches = RecordBatchIterator::new(batches.into_iter().map(Ok), schema.clone());
2700 Dataset::write(batches, test_uri, Some(write_params))
2701 .await
2702 .unwrap();
2703
2704 Dataset::open(test_uri).await.unwrap()
2705 }
2706
2707 #[rstest]
2708 #[tokio::test]
2709 async fn test_fragment_scan(
2710 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
2711 data_storage_version: LanceFileVersion,
2712 ) {
2713 let test_dir = TempStrDir::default();
2714 let test_uri = &test_dir;
2715 let dataset = create_dataset(test_uri, data_storage_version).await;
2716 let fragment = &dataset.get_fragments()[2];
2717 let mut scanner = fragment.scan();
2718 let batches = scanner
2719 .with_row_id()
2720 .filter(" i < 105")
2721 .unwrap()
2722 .try_into_stream()
2723 .await
2724 .unwrap()
2725 .try_collect::<Vec<_>>()
2726 .await
2727 .unwrap();
2728
2729 if data_storage_version == LanceFileVersion::Legacy {
2730 assert_eq!(batches.len(), 3);
2731
2732 assert_eq!(
2733 batches[0].column_by_name("i").unwrap().as_ref(),
2734 &Int32Array::from_iter_values(80..90)
2735 );
2736 assert_eq!(
2737 batches[1].column_by_name("i").unwrap().as_ref(),
2738 &Int32Array::from_iter_values(90..100)
2739 );
2740 assert_eq!(
2741 batches[2].column_by_name("i").unwrap().as_ref(),
2742 &Int32Array::from_iter_values(100..105)
2743 );
2744 } else {
2745 assert_eq!(batches.len(), 1);
2746
2747 assert_eq!(
2748 batches[0].column_by_name("i").unwrap().as_ref(),
2749 &Int32Array::from_iter_values(80..105)
2750 )
2751 }
2752 }
2753
2754 #[tokio::test]
2755 async fn test_fragment_scan_v2() {
2756 let test_dir = TempStrDir::default();
2757 let test_uri = &test_dir;
2758 let dataset = create_dataset_v2(test_uri).await;
2759 let fragment = &dataset.get_fragments()[2];
2760 let mut scanner = fragment.scan();
2761 let batches = scanner
2762 .with_row_id()
2763 .try_into_stream()
2764 .await
2765 .unwrap()
2766 .try_collect::<Vec<_>>()
2767 .await
2768 .unwrap();
2769
2770 assert_eq!(batches.len(), 1);
2771
2772 assert_eq!(
2773 batches[0].column_by_name("i").unwrap().as_ref(),
2774 &Int32Array::from_iter_values(80..120)
2775 );
2776
2777 let mut scanner = fragment.scan();
2778 let batches = scanner
2779 .with_row_id()
2780 .batch_size(20)
2781 .try_into_stream()
2782 .await
2783 .unwrap()
2784 .try_collect::<Vec<_>>()
2785 .await
2786 .unwrap();
2787
2788 assert_eq!(batches.len(), 2);
2789
2790 assert_eq!(
2791 batches[0].column_by_name("i").unwrap().as_ref(),
2792 &Int32Array::from_iter_values(80..100)
2793 );
2794 assert_eq!(
2795 batches[1].column_by_name("i").unwrap().as_ref(),
2796 &Int32Array::from_iter_values(100..120)
2797 );
2798 }
2799
2800 #[tokio::test]
2801 async fn test_fragment_update() {
2802 let test_dir = TempStrDir::default();
2803 let test_uri = &test_dir;
2804 let mut dataset = create_dataset_v2(test_uri).await;
2805
2806 let _ = dataset
2808 .add_columns(
2809 NewColumnTransform::SqlExpressions(vec![("col1".into(), "-1".into())]),
2810 None,
2811 None,
2812 )
2813 .await;
2814 let mut fragment1 = dataset.get_fragment(0).unwrap();
2815
2816 let schema1 = Arc::new(ArrowSchema::new(vec![
2817 ArrowField::new(ROW_ID, DataType::UInt64, false),
2818 ArrowField::new("col1", DataType::Int64, true),
2819 ]));
2820 let update_batch1 = RecordBatch::try_new(
2821 schema1.clone(),
2822 vec![
2823 Arc::new(UInt64Array::from(
2824 (0..40).filter(|&v| v != 0 && v != 3).collect::<Vec<_>>(),
2825 )),
2826 Arc::new(Int64Array::from(vec![2; 38])),
2827 ],
2828 )
2829 .unwrap();
2830 let right_stream1: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
2831 vec![Ok(update_batch1)].into_iter(),
2832 schema1,
2833 ));
2834 let (updated_fragment1, fields_modified1) = fragment1
2835 .update_columns(right_stream1, ROW_ID, ROW_ID)
2836 .await
2837 .unwrap();
2838 let op1 = Operation::Update {
2839 removed_fragment_ids: vec![],
2840 updated_fragments: vec![updated_fragment1],
2841 new_fragments: vec![],
2842 fields_modified: fields_modified1,
2843 merged_generations: Vec::new(),
2844 fields_for_preserving_frag_bitmap: vec![],
2845 update_mode: Some(UpdateMode::RewriteColumns),
2846 inserted_rows_filter: None,
2847 };
2848 let mut dataset1 = Dataset::commit(
2849 test_uri,
2850 op1,
2851 Some(dataset.version().version),
2852 None,
2853 None,
2854 Default::default(),
2855 true,
2856 )
2857 .await
2858 .unwrap();
2859 assert_eq!(dataset1.get_fragments().len(), 5);
2860 let scanner1 = dataset1.get_fragment(0).unwrap().scan();
2861 let batches1 = scanner1
2862 .try_into_stream()
2863 .await
2864 .unwrap()
2865 .try_collect::<Vec<_>>()
2866 .await
2867 .unwrap();
2868 assert_eq!(batches1.len(), 1);
2869 let mut expected_col1 = vec![2; 40];
2870 expected_col1[0] = -1;
2871 expected_col1[3] = -1;
2872 assert_eq!(
2873 batches1[0].column_by_name("col1").unwrap().as_ref(),
2874 &Int64Array::from(expected_col1)
2875 );
2876
2877 let _ = dataset1
2879 .add_columns(
2880 NewColumnTransform::SqlExpressions(vec![("col2".into(), "false".into())]),
2881 None,
2882 None,
2883 )
2884 .await;
2885 let mut fragment2 = dataset1.get_fragment(0).unwrap();
2886
2887 let schema2 = Arc::new(ArrowSchema::new(vec![
2888 ArrowField::new("i1", DataType::Int32, true),
2889 ArrowField::new("col2", DataType::Boolean, true),
2890 ArrowField::new("col1", DataType::Int64, true),
2891 ]));
2892 let update_batch2 = RecordBatch::try_new(
2893 schema2.clone(),
2894 vec![
2895 Arc::new(Int32Array::from(
2896 (0..40).filter(|&v| v != 0 && v != 3).collect::<Vec<_>>(),
2897 )),
2898 Arc::new(BooleanArray::from(vec![true; 38])),
2899 Arc::new(Int64Array::from(vec![3; 38])),
2900 ],
2901 )
2902 .unwrap();
2903 let right_stream2: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
2904 vec![Ok(update_batch2)].into_iter(),
2905 schema2,
2906 ));
2907 let (updated_fragment2, fields_modified2) = fragment2
2908 .update_columns(right_stream2, "i", "i1")
2909 .await
2910 .unwrap();
2911 let op = Operation::Update {
2912 removed_fragment_ids: vec![],
2913 updated_fragments: vec![updated_fragment2],
2914 new_fragments: vec![],
2915 fields_modified: fields_modified2,
2916 merged_generations: Vec::new(),
2917 fields_for_preserving_frag_bitmap: vec![],
2918 update_mode: Some(UpdateMode::RewriteColumns),
2919 inserted_rows_filter: None,
2920 };
2921 let dataset2 = Dataset::commit(
2922 test_uri,
2923 op,
2924 Some(dataset1.version().version),
2925 None,
2926 None,
2927 Default::default(),
2928 true,
2929 )
2930 .await
2931 .unwrap();
2932 assert_eq!(dataset2.get_fragments().len(), 5);
2933 let scanner2 = dataset2.get_fragment(0).unwrap().scan();
2934 let batches2 = scanner2
2935 .try_into_stream()
2936 .await
2937 .unwrap()
2938 .try_collect::<Vec<_>>()
2939 .await
2940 .unwrap();
2941 assert_eq!(batches2.len(), 1);
2942
2943 expected_col1 = vec![3; 40];
2944 expected_col1[0] = -1;
2945 expected_col1[3] = -1;
2946 assert_eq!(
2947 batches2[0].column_by_name("col1").unwrap().as_ref(),
2948 &Int64Array::from(expected_col1)
2949 );
2950 let mut expected_col2 = vec![true; 40];
2951 expected_col2[0] = false;
2952 expected_col2[3] = false;
2953 assert_eq!(
2954 batches2[0].column_by_name("col2").unwrap().as_ref(),
2955 &BooleanArray::from(expected_col2)
2956 );
2957 }
2958
2959 #[tokio::test]
2960 async fn test_out_of_range() {
2961 let test_dir = TempStrDir::default();
2962 let test_uri = &test_dir;
2963 let mut dataset = create_dataset(test_uri, LanceFileVersion::Legacy).await;
2965 dataset.delete("i >= 20").await.unwrap();
2967 let fragment = &dataset.get_fragments()[0];
2969 assert_eq!(fragment.metadata.num_rows().unwrap(), 20);
2970
2971 for with_row_id in [false, true] {
2973 let reader = fragment
2974 .open(
2975 fragment.schema(),
2976 FragReadConfig::default().with_row_id(with_row_id),
2977 )
2978 .await
2979 .unwrap();
2980 for valid_range in [0..40, 20..40] {
2981 reader
2982 .take_range(valid_range, 100)
2983 .unwrap()
2984 .buffered(1)
2985 .try_collect::<Vec<_>>()
2986 .await
2987 .unwrap();
2988 }
2989 for invalid_range in [0..41, 41..42] {
2990 assert!(reader.take_range(invalid_range, 100).is_err());
2991 }
2992 }
2993
2994 for with_row_id in [false, true] {
2996 let reader = fragment
2997 .open(
2998 fragment.schema(),
2999 FragReadConfig::default().with_row_id(with_row_id),
3000 )
3001 .await
3002 .unwrap();
3003 for valid_range in [0..20, 0..10, 10..20] {
3004 reader
3005 .read_range(valid_range, 100)
3006 .unwrap()
3007 .buffered(1)
3008 .try_collect::<Vec<_>>()
3009 .await
3010 .unwrap();
3011 }
3012 for invalid_range in [0..21, 21..22] {
3013 assert!(reader.read_range(invalid_range, 100).is_err());
3014 }
3015 }
3016 }
3017
3018 #[tokio::test]
3019 async fn test_rowid_rowaddr_only() {
3020 let test_dir = TempStrDir::default();
3021 let test_uri = &test_dir;
3022 let mut dataset = create_dataset(test_uri, LanceFileVersion::Legacy).await;
3024 dataset.delete("i >= 20").await.unwrap();
3026 let fragment = &dataset.get_fragments()[0];
3028 assert_eq!(fragment.metadata.num_rows().unwrap(), 20);
3029
3030 for (with_row_id, with_row_address) in [(false, true), (true, false), (true, true)] {
3032 let reader = fragment
3033 .open(
3034 &fragment.schema().project::<&str>(&[]).unwrap(),
3035 FragReadConfig::default()
3036 .with_row_id(with_row_id)
3037 .with_row_address(with_row_address),
3038 )
3039 .await
3040 .unwrap();
3041 for valid_range in [0..40, 20..40] {
3042 reader
3043 .take_range(valid_range, 100)
3044 .unwrap()
3045 .buffered(1)
3046 .try_collect::<Vec<_>>()
3047 .await
3048 .unwrap();
3049 }
3050 for invalid_range in [0..41, 41..42] {
3051 assert!(reader.take_range(invalid_range, 100).is_err());
3052 }
3053 }
3054
3055 for (with_row_id, with_row_address) in [(false, true), (true, false), (true, true)] {
3057 let reader = fragment
3058 .open(
3059 &fragment.schema().project::<&str>(&[]).unwrap(),
3060 FragReadConfig::default()
3061 .with_row_id(with_row_id)
3062 .with_row_address(with_row_address),
3063 )
3064 .await
3065 .unwrap();
3066 for valid_range in [0..20, 0..10, 10..20] {
3067 reader
3068 .read_range(valid_range, 100)
3069 .unwrap()
3070 .buffered(1)
3071 .try_collect::<Vec<_>>()
3072 .await
3073 .unwrap();
3074 }
3075 for invalid_range in [0..21, 21..22] {
3076 assert!(reader.read_range(invalid_range, 100).is_err());
3077 }
3078 }
3079 }
3080
3081 #[rstest]
3082 #[tokio::test]
3083 async fn test_fragment_take_range_deletions(
3084 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3085 data_storage_version: LanceFileVersion,
3086 ) {
3087 let test_dir = TempStrDir::default();
3088 let test_uri = &test_dir;
3089 let mut dataset = create_dataset(test_uri, data_storage_version).await;
3090 dataset.delete("i >= 0 and i < 15").await.unwrap();
3091
3092 let fragment = &dataset.get_fragments()[0];
3093 let mut reader = fragment
3094 .open(
3095 dataset.schema(),
3096 FragReadConfig::default().with_row_id(true),
3097 )
3098 .await
3099 .unwrap();
3100 reader.with_make_deletions_null();
3101
3102 if data_storage_version == LanceFileVersion::Legacy {
3103 let batch1 = reader
3105 .legacy_read_batch_projected(0, .., dataset.schema())
3106 .await
3107 .unwrap();
3108 assert_eq!(
3109 batch1.column_by_name(ROW_ID).unwrap().as_ref(),
3110 &UInt64Array::from_iter(std::iter::repeat_n(None, 10))
3111 );
3112
3113 let batch2 = reader
3116 .legacy_read_batch_projected(1, .., dataset.schema())
3117 .await
3118 .unwrap();
3119 assert_eq!(
3120 batch2.column_by_name(ROW_ID).unwrap().as_ref(),
3121 &UInt64Array::from_iter((10..20).map(|v| if v < 15 { None } else { Some(v) }))
3122 );
3123
3124 let batch3 = reader
3126 .legacy_read_batch_projected(2, .., dataset.schema())
3127 .await
3128 .unwrap();
3129 assert_eq!(
3130 batch3.column_by_name(ROW_ID).unwrap().as_ref(),
3131 &UInt64Array::from_iter_values(20..30)
3132 );
3133 } else {
3134 let to_batches = |range: Range<u32>| {
3135 let batch_size = range.len() as u32;
3136 reader
3137 .take_range(range, batch_size)
3138 .unwrap()
3139 .buffered(1)
3140 .try_collect::<Vec<_>>()
3141 };
3142
3143 let batches = to_batches(0..10).await.unwrap();
3145 assert_eq!(batches.len(), 1);
3146 let batch = batches.into_iter().next().unwrap();
3147 assert_eq!(
3148 batch.column_by_name(ROW_ID).unwrap().as_ref(),
3149 &UInt64Array::from_iter(std::iter::repeat_n(None, 10))
3150 );
3151
3152 let batches = to_batches(10..20).await.unwrap();
3153 assert_eq!(batches.len(), 1);
3154 let batch = batches.into_iter().next().unwrap();
3155 assert_eq!(
3158 batch.column_by_name(ROW_ID).unwrap().as_ref(),
3159 &UInt64Array::from_iter((10..20).map(|v| if v < 15 { None } else { Some(v) }))
3160 );
3161
3162 let batches = to_batches(20..30).await.unwrap();
3164 assert_eq!(batches.len(), 1);
3165 let batch = batches.into_iter().next().unwrap();
3166 assert_eq!(
3167 batch.column_by_name(ROW_ID).unwrap().as_ref(),
3168 &UInt64Array::from_iter_values(20..30)
3169 );
3170 }
3171 }
3172
3173 #[rstest]
3174 #[tokio::test]
3175 async fn test_range_scan_deletions(
3176 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3177 data_storage_version: LanceFileVersion,
3178 ) {
3179 let test_dir = TempStrDir::default();
3180 let test_uri = &test_dir;
3181 let dataset = create_dataset(test_uri, data_storage_version).await;
3182
3183 let version = dataset.version().version;
3184
3185 let check = |cond: &'static str, range: Range<u32>, expected: Vec<i32>| async {
3186 let mut dataset = dataset.checkout_version(version).await.unwrap();
3187 dataset.restore().await.unwrap();
3188 dataset.delete(cond).await.unwrap();
3189
3190 let fragment = &dataset.get_fragments()[0];
3191 let reader = fragment
3192 .open(
3193 dataset.schema(),
3194 FragReadConfig::default().with_row_id(true),
3195 )
3196 .await
3197 .unwrap();
3198
3199 let mut stream = reader.read_range(range, 20).unwrap();
3203 let mut batches = Vec::new();
3204 while let Some(next) = stream.next().await {
3205 batches.push(next.await.unwrap());
3206 }
3207 let schema = Arc::new(dataset.schema().into());
3208 let batch = arrow_select::concat::concat_batches(&schema, batches.iter()).unwrap();
3209
3210 assert_eq!(batch.num_rows(), expected.len());
3211 assert_eq!(
3212 batch.column_by_name("i").unwrap().as_ref(),
3213 &Int32Array::from(expected)
3214 );
3215 };
3216 check("i < 5", 0..2, vec![5, 6]).await;
3218 check("i < 5", 0..15, (5..20).collect()).await;
3219 check("i >= 5 and i < 15", 7..9, vec![17, 18]).await;
3221 check("i >= 5 and i < 15", 3..5, vec![3, 4]).await;
3222 check("i >= 5 and i < 15", 3..6, vec![3, 4, 15]).await;
3223 check("i >= 5 and i < 15", 5..6, vec![15]).await;
3224 check("i >= 5 and i < 15", 5..10, vec![15, 16, 17, 18, 19]).await;
3225 check(
3226 "i >= 5 and i < 15",
3227 0..10,
3228 vec![0, 1, 2, 3, 4, 15, 16, 17, 18, 19],
3229 )
3230 .await;
3231 check("i >= 15", 10..15, vec![10, 11, 12, 13, 14]).await;
3233 check("i >= 15", 0..15, (0..15).collect()).await;
3234 }
3235
3236 #[rstest]
3237 #[tokio::test]
3238 async fn test_fragment_take_indices(
3239 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3240 data_storage_version: LanceFileVersion,
3241 ) {
3242 let test_dir = TempStrDir::default();
3243 let test_uri = &test_dir;
3244 let mut dataset = create_dataset(test_uri, data_storage_version).await;
3245 let fragment = dataset
3246 .get_fragments()
3247 .into_iter()
3248 .find(|f| f.id() == 3)
3249 .unwrap();
3250
3251 let batch = fragment
3253 .take(&[1, 2, 4, 5, 5, 8], dataset.schema())
3254 .await
3255 .unwrap();
3256 assert_eq!(
3257 batch.column_by_name("i").unwrap().as_ref(),
3258 &Int32Array::from(vec![121, 122, 124, 125, 125, 128])
3259 );
3260
3261 dataset.delete("i in (122, 123, 125)").await.unwrap();
3262 dataset.validate().await.unwrap();
3263
3264 let fragment = dataset
3266 .get_fragments()
3267 .into_iter()
3268 .find(|f| f.id() == 3)
3269 .unwrap();
3270 assert!(fragment.metadata().deletion_file.is_some());
3271 let batch = fragment
3272 .take(&[1, 2, 4, 5, 8], dataset.schema())
3273 .await
3274 .unwrap();
3275 assert_eq!(
3276 batch.column_by_name("i").unwrap().as_ref(),
3277 &Int32Array::from(vec![121, 124, 127, 128, 131])
3278 );
3279
3280 let batch = fragment.take(&[], dataset.schema()).await.unwrap();
3282 assert_eq!(
3283 batch.column_by_name("i").unwrap().as_ref(),
3284 &Int32Array::from(Vec::<i32>::new())
3285 );
3286 }
3287
3288 #[rstest]
3289 #[tokio::test]
3290 async fn test_fragment_take_rows(
3291 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3292 data_storage_version: LanceFileVersion,
3293 ) {
3294 let test_dir = TempStrDir::default();
3295 let test_uri = &test_dir;
3296 let mut dataset = create_dataset(test_uri, data_storage_version).await;
3297 let fragment = dataset
3298 .get_fragments()
3299 .into_iter()
3300 .find(|f| f.id() == 3)
3301 .unwrap();
3302
3303 let batch = fragment
3305 .take_rows(
3306 &[1, 2, 4, 5, 5, 8],
3307 dataset.schema(),
3308 false,
3309 false,
3310 false,
3311 false,
3312 )
3313 .await
3314 .unwrap();
3315 assert_eq!(
3316 batch.column_by_name("i").unwrap().as_ref(),
3317 &Int32Array::from(vec![121, 122, 124, 125, 125, 128])
3318 );
3319
3320 dataset.delete("i in (122, 124)").await.unwrap();
3321 dataset.validate().await.unwrap();
3322
3323 let fragment = dataset
3325 .get_fragments()
3326 .into_iter()
3327 .find(|f| f.id() == 3)
3328 .unwrap();
3329 assert!(fragment.metadata().deletion_file.is_some());
3330 let batch = fragment
3331 .take_rows(
3332 &[1, 2, 4, 5, 8],
3333 dataset.schema(),
3334 false,
3335 false,
3336 false,
3337 false,
3338 )
3339 .await
3340 .unwrap();
3341 assert_eq!(
3342 batch.column_by_name("i").unwrap().as_ref(),
3343 &Int32Array::from(vec![121, 125, 128])
3344 );
3345
3346 let batch = fragment
3348 .take_rows(&[], dataset.schema(), false, false, false, false)
3349 .await
3350 .unwrap();
3351 assert_eq!(
3352 batch.column_by_name("i").unwrap().as_ref(),
3353 &Int32Array::from(Vec::<i32>::new())
3354 );
3355
3356 let batch = fragment
3358 .take_rows(
3359 &[1, 2, 4, 5, 8],
3360 dataset.schema(),
3361 false,
3362 true,
3363 false,
3364 false,
3365 )
3366 .await
3367 .unwrap();
3368 assert_eq!(
3369 batch.column_by_name("i").unwrap().as_ref(),
3370 &Int32Array::from(vec![121, 125, 128])
3371 );
3372 assert_eq!(
3373 batch.column_by_name(ROW_ADDR).unwrap().as_ref(),
3374 &UInt64Array::from(vec![(3 << 32) + 1, (3 << 32) + 5, (3 << 32) + 8])
3375 );
3376 }
3377
3378 #[tokio::test]
3379 async fn test_recommit_from_file() {
3380 let test_dir = TempStrDir::default();
3381 let test_uri = &test_dir;
3382 let dataset = create_dataset(test_uri, LanceFileVersion::Legacy).await;
3383 let schema = dataset.schema();
3384 let dataset_rows = dataset.count_rows(None).await.unwrap();
3385
3386 let mut paths: Vec<String> = Vec::new();
3387 for f in dataset.get_fragments() {
3388 for file in Fragment::from(f.clone()).files {
3389 let p = file.path.clone();
3390 paths.push(p);
3391 }
3392 }
3393
3394 let mut fragments: Vec<Fragment> = Vec::new();
3395 for (idx, path) in paths.iter().enumerate() {
3396 let f = FileFragment::create_from_file(path, &dataset, idx, None)
3397 .await
3398 .unwrap();
3399 fragments.push(f)
3400 }
3401
3402 let op = Operation::Overwrite {
3403 schema: schema.clone(),
3404 fragments,
3405 config_upsert_values: None,
3406 initial_bases: None,
3407 };
3408
3409 let new_dataset =
3410 Dataset::commit(test_uri, op, None, None, None, Default::default(), false)
3411 .await
3412 .unwrap();
3413
3414 assert_eq!(new_dataset.count_rows(None).await.unwrap(), dataset_rows);
3415
3416 let fragments = new_dataset.get_fragments();
3419 assert_eq!(fragments.len(), 5);
3420 for f in fragments {
3421 assert_eq!(f.metadata.num_rows(), Some(40));
3422 assert_eq!(f.count_rows(None).await.unwrap(), 40);
3423 assert_eq!(f.metadata().deletion_file, None);
3424 }
3425 }
3426
3427 #[rstest]
3428 #[tokio::test]
3429 async fn test_fragment_count(
3430 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3431 data_storage_version: LanceFileVersion,
3432 ) {
3433 let test_dir = TempStrDir::default();
3434 let test_uri = &test_dir;
3435 let dataset = create_dataset(test_uri, data_storage_version).await;
3436 let fragment = dataset.get_fragments().pop().unwrap();
3437
3438 assert_eq!(fragment.count_rows(None).await.unwrap(), 40);
3439 assert_eq!(fragment.physical_rows().await.unwrap(), 40);
3440 assert!(fragment.metadata.deletion_file.is_none());
3441
3442 assert_eq!(
3443 fragment
3444 .count_rows(Some("i < 170".to_string()))
3445 .await
3446 .unwrap(),
3447 10
3448 );
3449
3450 let fragment = fragment
3451 .delete("i >= 160 and i <= 172")
3452 .await
3453 .unwrap()
3454 .unwrap();
3455
3456 fragment.validate().await.unwrap();
3457
3458 assert_eq!(fragment.count_rows(None).await.unwrap(), 27);
3459 assert_eq!(fragment.physical_rows().await.unwrap(), 40);
3460 assert!(fragment.metadata.deletion_file.is_some());
3461 assert_eq!(
3462 fragment.metadata.deletion_file.unwrap().num_deleted_rows,
3463 Some(13)
3464 );
3465 }
3466
3467 #[rstest]
3468 #[tokio::test]
3469 async fn test_append_new_columns(
3470 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3471 data_storage_version: LanceFileVersion,
3472 ) {
3473 for with_delete in [true, false] {
3474 let test_dir = TempStrDir::default();
3475 let test_uri = &test_dir;
3476 let mut dataset = create_dataset(test_uri, data_storage_version).await;
3477 dataset.validate().await.unwrap();
3478 assert_eq!(dataset.count_rows(None).await.unwrap(), 200);
3479
3480 if with_delete {
3481 dataset.delete("i >= 15 and i < 20").await.unwrap();
3482 dataset.validate().await.unwrap();
3483 assert_eq!(dataset.count_rows(None).await.unwrap(), 195);
3484 }
3485
3486 let fragment = &mut dataset.get_fragment(0).unwrap();
3487 let mut updater = fragment.updater(Some(&["i"]), None, None).await.unwrap();
3488 let new_schema = Arc::new(ArrowSchema::new(vec![ArrowField::new(
3489 "double_i",
3490 DataType::Int32,
3491 true,
3492 )]));
3493 while let Some(batch) = updater.next().await.unwrap() {
3494 let input_col = batch.column_by_name("i").unwrap();
3495 let result_col = mul(input_col, &Int32Array::new_scalar(2)).unwrap();
3496 let batch = RecordBatch::try_new(
3497 new_schema.clone(),
3498 vec![Arc::new(result_col) as ArrayRef],
3499 )
3500 .unwrap();
3501 updater.update(batch).await.unwrap();
3502 }
3503 let new_fragment = updater.finish().await.unwrap();
3504
3505 assert_eq!(new_fragment.files.len(), 2);
3506
3507 let mut full_schema = dataset.schema().merge(new_schema.as_ref()).unwrap();
3509 full_schema.set_field_id(None);
3510 let before_version = dataset.version().version;
3511
3512 let op = Operation::Overwrite {
3513 fragments: vec![new_fragment],
3514 schema: full_schema.clone(),
3515 config_upsert_values: None,
3516 initial_bases: None,
3517 };
3518
3519 let dataset =
3520 Dataset::commit(test_uri, op, None, None, None, Default::default(), false)
3521 .await
3522 .unwrap();
3523
3524 assert_eq!(
3526 dataset.count_rows(None).await.unwrap(),
3527 if with_delete { 35 } else { 40 }
3528 );
3529 assert_eq!(dataset.version().version, before_version + 1);
3530 dataset.validate().await.unwrap();
3531 let new_projection = full_schema.project(&["i", "double_i"]).unwrap();
3532
3533 let stream = dataset
3534 .scan()
3535 .batch_size(10)
3536 .project(&["i", "double_i"])
3537 .unwrap()
3538 .try_into_stream()
3539 .await
3540 .unwrap();
3541 let batches = stream.try_collect::<Vec<_>>().await.unwrap();
3542
3543 assert_eq!(batches[1].schema().as_ref(), &(&new_projection).into());
3544 let expected_i = match (with_delete, data_storage_version) {
3545 (true, LanceFileVersion::Legacy) => vec![10, 11, 12, 13, 14],
3548 (true, _) => vec![10, 11, 12, 13, 14, 20, 21, 22, 23, 24],
3550 (false, _) => vec![10, 11, 12, 13, 14, 15, 16, 17, 18, 19],
3551 };
3552 let expected_batch = RecordBatch::try_new(
3553 Arc::new(ArrowSchema::new(vec![
3554 ArrowField::new("i", DataType::Int32, true),
3555 ArrowField::new("double_i", DataType::Int32, true),
3556 ])),
3557 vec![
3558 Arc::new(Int32Array::from_iter_values(expected_i.iter().copied())),
3559 Arc::new(Int32Array::from_iter_values(
3560 expected_i.iter().map(|i| 2 * i),
3561 )),
3562 ],
3563 )
3564 .unwrap();
3565 assert_eq!(batches[1], expected_batch);
3566 }
3567 }
3568
3569 #[rstest]
3570 #[tokio::test]
3571 async fn test_merge_fragment(
3572 #[values(LanceFileVersion::Legacy, LanceFileVersion::Stable)]
3573 data_storage_version: LanceFileVersion,
3574 ) {
3575 let test_dir = TempStrDir::default();
3576 let test_uri = &test_dir;
3577 let mut dataset = create_dataset(test_uri, data_storage_version).await;
3578 dataset.validate().await.unwrap();
3579 assert_eq!(dataset.count_rows(None).await.unwrap(), 200);
3580
3581 let deleted_range = 15..20;
3582 dataset.delete("i >= 15 and i < 20").await.unwrap();
3583 dataset.validate().await.unwrap();
3584 assert_eq!(dataset.count_rows(None).await.unwrap(), 195);
3585
3586 let schema = Arc::new(ArrowSchema::new(vec![
3588 ArrowField::new("i", DataType::Int32, true),
3589 ArrowField::new("double_i", DataType::Int32, true),
3590 ]));
3591 let to_merge = RecordBatch::try_new(
3592 schema.clone(),
3593 vec![
3594 Arc::new(Int32Array::from_iter_values(0..200)),
3595 Arc::new(Int32Array::from_iter_values((0..400).step_by(2))),
3596 ],
3597 )
3598 .unwrap();
3599
3600 let stream = RecordBatchIterator::new(vec![Ok(to_merge)], schema.clone());
3601 dataset.merge(stream, "i", "i").await.unwrap();
3602 dataset.validate().await.unwrap();
3603
3604 let batches = dataset
3606 .scan()
3607 .project(&["i", "double_i"])
3608 .unwrap()
3609 .try_into_stream()
3610 .await
3611 .unwrap()
3612 .try_collect::<Vec<_>>()
3613 .await
3614 .unwrap();
3615 let batch = concat_batches(&schema, &batches).unwrap();
3616
3617 let mut row_id: i32 = 0;
3618 let mut i: usize = 0;
3619 let array_i: &Int32Array = as_primitive_array(&batch["i"]);
3620 let array_double_i: &Int32Array = as_primitive_array(&batch["double_i"]);
3621 while row_id < 200 {
3622 if deleted_range.contains(&row_id) {
3623 row_id += 1;
3624 continue;
3625 }
3626 assert_eq!(array_i.value(i), row_id);
3627 assert_eq!(array_double_i.value(i), 2 * row_id);
3628 row_id += 1;
3629 i += 1;
3630 }
3631 }
3632
3633 #[tokio::test]
3634 async fn test_write_batch_size() {
3635 let test_dir = TempStrDir::default();
3641 let test_uri = &test_dir;
3642
3643 let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new(
3644 "i",
3645 DataType::Int32,
3646 true,
3647 )]));
3648
3649 let in_memory_batch = 1024;
3650 let batches: Vec<RecordBatch> = (0..10)
3651 .map(|i| {
3652 RecordBatch::try_new(
3653 schema.clone(),
3654 vec![Arc::new(Int32Array::from_iter_values(
3655 i * in_memory_batch..(i + 1) * in_memory_batch,
3656 ))],
3657 )
3658 .unwrap()
3659 })
3660 .collect();
3661
3662 let batch_iter = RecordBatchIterator::new(batches.into_iter().map(Ok), schema.clone());
3663
3664 let fragment = FileFragment::create(
3665 test_uri,
3666 10,
3667 batch_iter,
3668 Some(WriteParams {
3669 max_rows_per_group: 100,
3670 data_storage_version: Some(LanceFileVersion::Legacy),
3671 ..Default::default()
3672 }),
3673 )
3674 .await
3675 .unwrap();
3676
3677 let (object_store, base_path) = ObjectStore::from_uri(test_uri).await.unwrap();
3678 let file_reader = PreviousFileReader::try_new_with_fragment_id(
3679 &object_store,
3680 &base_path
3681 .child("data")
3682 .child(fragment.files[0].path.as_str()),
3683 schema.as_ref().try_into().unwrap(),
3684 10,
3685 0,
3686 1,
3687 None,
3688 )
3689 .await
3690 .unwrap();
3691
3692 for i in 0..file_reader.num_batches() - 1 {
3693 assert_eq!(file_reader.num_rows_in_batch(i as i32), 100);
3694 }
3695 assert_eq!(
3696 file_reader.num_rows_in_batch(file_reader.num_batches() as i32 - 1) as i32,
3697 in_memory_batch * 10 % 100
3698 );
3699 }
3700
3701 #[tokio::test]
3702 async fn test_shuffled_columns() -> Result<()> {
3703 let batch_i = RecordBatch::try_new(
3707 Arc::new(ArrowSchema::new(vec![ArrowField::new(
3708 "i",
3709 DataType::Int32,
3710 true,
3711 )])),
3712 vec![Arc::new(Int32Array::from_iter_values(0..20))],
3713 )?;
3714
3715 let batch_s = RecordBatch::try_new(
3716 Arc::new(ArrowSchema::new(vec![ArrowField::new(
3717 "s",
3718 DataType::Utf8,
3719 true,
3720 )])),
3721 vec![Arc::new(StringArray::from_iter_values(
3722 (0..20).map(|v| format!("s-{}", v)),
3723 ))],
3724 )?;
3725
3726 let test_dir = TempStrDir::default();
3728 let test_uri = &test_dir;
3729
3730 let dataset = Dataset::write(
3731 RecordBatchIterator::new(vec![Ok(batch_i.clone())], batch_i.schema().clone()),
3732 test_uri,
3733 None,
3734 )
3735 .await?;
3736
3737 let fragment = dataset.get_fragments().pop().unwrap();
3738
3739 let mut updater = fragment.updater(Some(&["i"]), None, None).await?;
3741 updater.next().await?;
3742 updater.update(batch_s.clone()).await?;
3743 let frag = updater.finish().await?;
3744
3745 let schema = updater.schema().unwrap().clone().project(&["s", "i"])?;
3747
3748 let dataset = Dataset::commit(
3749 test_uri,
3750 Operation::Merge {
3751 schema,
3752 fragments: vec![frag],
3753 },
3754 Some(dataset.manifest.version),
3755 None,
3756 None,
3757 Default::default(),
3758 false,
3759 )
3760 .await?;
3761
3762 let expected_data = batch_s.merge(&batch_i)?;
3763 let actual_data = dataset.scan().try_into_batch().await?;
3764 assert_eq!(expected_data, actual_data);
3765
3766 let reader = dataset
3768 .get_fragments()
3769 .first()
3770 .unwrap()
3771 .open(dataset.schema(), FragReadConfig::default())
3772 .await?;
3773 let actual_data = reader.take_as_batch(&[0, 1, 2], None).await?;
3774 assert_eq!(expected_data.slice(0, 3), actual_data);
3775
3776 let actual_data = reader
3777 .read_range(0..3, 3)
3778 .unwrap()
3779 .next()
3780 .await
3781 .unwrap()
3782 .await
3783 .unwrap();
3784 assert_eq!(expected_data.slice(0, 3), actual_data);
3785
3786 let expected_data = expected_data.try_with_column(
3788 ROW_ID_FIELD.clone(),
3789 Arc::new(UInt64Array::from_iter_values(0..20)),
3790 )?;
3791 let actual_data = dataset.scan().with_row_id().try_into_batch().await?;
3792 assert_eq!(expected_data, actual_data);
3793
3794 Ok(())
3795 }
3796
3797 #[tokio::test]
3798 async fn test_row_id_reader() -> Result<()> {
3799 let batch = RecordBatch::try_new(
3801 Arc::new(ArrowSchema::new(vec![ArrowField::new(
3802 "i",
3803 DataType::Int32,
3804 true,
3805 )])),
3806 vec![Arc::new(Int32Array::from_iter_values(0..20))],
3807 )?;
3808
3809 let test_dir = TempStrDir::default();
3810 let test_uri = &test_dir;
3811
3812 let dataset = Dataset::write(
3813 RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema().clone()),
3814 test_uri,
3815 None,
3816 )
3817 .await?;
3818
3819 let fragment = dataset.get_fragments().pop().unwrap();
3820
3821 let reader = fragment
3822 .open(
3823 &dataset.schema().project::<&str>(&[])?,
3824 FragReadConfig::default().with_row_id(true),
3825 )
3826 .await?;
3827 let batch = reader.legacy_read_range_as_batch(0..20).await?;
3828
3829 let expected_data = RecordBatch::try_new(
3830 Arc::new(ArrowSchema::new(vec![ROW_ID_FIELD.clone()])),
3831 vec![Arc::new(UInt64Array::from_iter_values(0..20))],
3832 )?;
3833 assert_eq!(expected_data, batch);
3834
3835 let res = fragment
3837 .open(
3838 &dataset.schema().project::<&str>(&[])?,
3839 FragReadConfig::default(),
3840 )
3841 .await;
3842 assert!(matches!(res, Err(Error::NotFound { .. })));
3843
3844 Ok(())
3845 }
3846
3847 #[tokio::test]
3848 async fn create_from_file_v2() {
3849 let test_dir = TempStrDir::default();
3850 let test_uri = &test_dir;
3851
3852 let make_gen = || {
3853 gen_batch()
3854 .col("str", array::rand_type(&DataType::Utf8))
3855 .col("int", array::rand_type(&DataType::Int32))
3856 };
3857
3858 let batch = make_gen().into_batch_rows(RowCount::from(128)).unwrap();
3859 let dataset = TestDatasetGenerator::new(vec![batch], LanceFileVersion::Stable)
3860 .make_hostile(test_uri)
3861 .await;
3862
3863 let new_data = make_gen().into_batch_rows(RowCount::from(128)).unwrap();
3864 let store = ObjectStore::local();
3865 let file_path = dataset.data_dir().child("some_file.lance");
3866 let object_writer = store.create(&file_path).await.unwrap();
3867 let mut file_writer =
3868 lance_file::writer::FileWriter::new_lazy(object_writer, FileWriterOptions::default());
3869 file_writer.write_batch(&new_data).await.unwrap();
3870 file_writer.finish().await.unwrap();
3871
3872 let frag = FileFragment::create_from_file("some_file.lance", &dataset, 0, Some(128))
3873 .await
3874 .unwrap();
3875
3876 assert_eq!(
3877 Fragment::try_infer_version(std::slice::from_ref(&frag))
3878 .unwrap()
3879 .unwrap(),
3880 LanceFileVersion::Stable.resolve()
3881 );
3882
3883 let op = Operation::Append {
3884 fragments: vec![frag],
3885 };
3886 let dataset = Dataset::commit(
3887 &dataset.uri,
3888 op,
3889 Some(dataset.version().version),
3890 None,
3891 None,
3892 Default::default(),
3893 false,
3894 )
3895 .await
3896 .unwrap();
3897
3898 assert_eq!(
3899 dataset
3900 .count_rows(Some("int IS NOT NULL".to_string()))
3901 .await
3902 .unwrap(),
3903 256
3904 );
3905 }
3906
3907 #[tokio::test]
3908 async fn test_iops_read_small() {
3909 let schema = Arc::new(ArrowSchema::new(
3911 (0..8)
3912 .map(|i| ArrowField::new(format!("col_{}", i), DataType::Int32, true))
3913 .collect::<Vec<_>>(),
3914 ));
3915
3916 let batch = RecordBatch::try_new(
3918 schema.clone(),
3919 (0..8)
3920 .map(|i| Arc::new(Int32Array::from(vec![i])) as ArrayRef)
3921 .collect(),
3922 )
3923 .unwrap();
3924 let session = Arc::new(Session::default());
3925 let write_params = WriteParams {
3926 session: Some(session.clone()),
3927 ..Default::default()
3928 };
3929 let dataset = InsertBuilder::new("memory://test")
3930 .with_params(&write_params)
3931 .execute(vec![batch])
3932 .await
3933 .unwrap();
3934 let fragment = dataset.get_fragments().pop().unwrap();
3935
3936 {
3938 let stats = dataset.object_store().io_stats_incremental();
3939 assert_io_eq!(stats, write_iops, 3);
3940 assert_io_lt!(stats, written_bytes, 4300);
3941 }
3942
3943 let projection = Schema::try_from(schema.as_ref())
3945 .unwrap()
3946 .project_by_ids(&[0, 1, 2, 3, 4, 6, 7, 8, 9], true);
3947 let reader = fragment
3948 .open(&projection, Default::default())
3949 .await
3950 .unwrap();
3951 let mut data = reader
3952 .read_all(1024)
3953 .unwrap()
3954 .buffered(1)
3955 .try_collect::<Vec<_>>()
3956 .await
3957 .unwrap();
3958 assert_eq!(data.len(), 1);
3959 let data = data.pop().unwrap();
3960 assert_eq!(data.num_rows(), 1);
3961 assert_eq!(data.num_columns(), 7);
3962
3963 let stats = dataset.object_store().io_stats_incremental();
3964 assert_io_eq!(stats, read_iops, 1);
3965 assert_io_lt!(stats, read_bytes, 4096);
3966 }
3967}