Skip to main content

lance/dataset/
fragment.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4//! Wraps a Fragment of the dataset.
5
6pub 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/// A Fragment of a Lance [`Dataset`].
65///
66/// The interface is modeled after `pyarrow.dataset.Fragment`.
67#[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/// A trait for file readers to be implemented by both the v1 and v2 readers
77#[allow(clippy::len_without_is_empty)]
78pub trait GenericFileReader: std::fmt::Debug + Send + Sync {
79    /// Reads the requested range of rows from the file, returning as a stream
80    /// of tasks.
81    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    /// Reads the requested ranges of rows from the file, only supported by v2
88    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    /// Reads all rows from the file, returning as a stream of tasks
95    fn read_all_tasks(
96        &self,
97        batch_size: u32,
98        projection: Arc<lance_core::datatypes::Schema>,
99    ) -> Result<ReadBatchTaskStream>;
100    /// Take specific rows from the file, returning as a stream of tasks
101    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    /// Return the number of rows in the file
110    fn len(&self) -> u32;
111
112    /// Schema of the reader
113    fn projection(&self) -> &Arc<Schema>;
114
115    /// Get storage statistics for this file (ignored by v1 reader)
116    fn storage_stats(&self) -> Vec<(u32, u64)>;
117
118    // Helper functions to fallback to the legacy implementation while we
119    // slowly migrate functionality over to the generic reader
120
121    // Clone the reader, this is needed because Box<dyn Foo: Clone> doesn't
122    // implement Clone
123    fn clone_box(&self) -> Box<dyn GenericFileReader>;
124    // Return true if the reader is a v1 reader
125    fn is_legacy(&self) -> bool;
126    // Return a reference to the legacy reader, panics if called on a v2
127    // file.
128    fn as_legacy(&self) -> &PreviousFileReader {
129        self.as_legacy_opt()
130            .expect("legacy function called on v2 file")
131    }
132    // Return a reference to the legacy reader if this is a v1 reader and
133    // return None otherwise
134    fn as_legacy_opt(&self) -> Option<&PreviousFileReader>;
135    // Return a mutable reference to the legacy reader if this is a v1 reader
136    // and return None otherwise
137    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    /// Reads the requested range of rows from the file, returning as a stream
184    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        // In the new path the row id is added by the fragment and not the file
254        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    /// Return the number of rows in the file
268    fn len(&self) -> u32 {
269        self.reader.len() as u32
270    }
271
272    fn storage_stats(&self) -> Vec<(u32, u64)> {
273        // No-op for v1 files
274        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        /// Reads the requested range of rows from the file, returning as a stream
328        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            // Some fields span more than one column.  We assume a column that doesn't have an
454            // entry in the field_id_to_column_idx map is a continuation of the previous field.
455            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        /// Return the number of rows in the file
470        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/// A reader where all rows are null. Used when there are fields that have no
493/// data files in a fragment.
494#[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        // No-op for null reader
573        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    // Add the row id column
604    pub with_row_id: bool,
605    // Add the row address column
606    pub with_row_address: bool,
607    // Add the last updated at version column
608    pub with_row_last_updated_at_version: bool,
609    // Add the created at version column
610    pub with_row_created_at_version: bool,
611    /// The scan scheduler to use for reading data files.
612    ///
613    /// This should be specified if multiple readers are being used in
614    /// an operation
615    pub scan_scheduler: Option<Arc<ScanScheduler>>,
616    /// The default scan priority to use for reading data files
617    ///
618    /// Only used if `scan_scheduler` is provided
619    ///
620    /// The overall priority for reads will be
621    ///
622    /// operation_priority: u32 | reader_priority: u32 | file_position: u64
623    pub reader_priority: Option<u32>,
624    /// File reader options to use when reading data files.
625    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    /// Creates a new FileFragment.
674    pub fn new(dataset: Arc<Dataset>, metadata: Fragment) -> Self {
675        Self { dataset, metadata }
676    }
677
678    /// Create a new [`FileFragment`] from a [`StreamingWriteSource`].
679    ///
680    /// This method can be used before a `Dataset` is created. For example,
681    /// Fragments can be created distributed first, before a central machine to
682    /// commit the dataset with these fragments.
683    ///
684    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    /// Create a list of [`FileFragment`] from a [`StreamingWriteSource`].
700    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            // Load the file metadata, confirm the schema is compatible, and
742            // determine the column offsets
743            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            // If the schemas are not compatible we can't calculate field id offsets
760            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    /// Returns storage stats as `(field_id, bytes_on_disk)` pairs for this fragment.
789    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    /// Returns the fragment's metadata.
816    pub fn metadata(&self) -> &Fragment {
817        &self.metadata
818    }
819
820    /// The id of this [`FileFragment`].
821    pub fn id(&self) -> usize {
822        self.metadata.id as usize
823    }
824
825    /// The number of data files in this fragment.
826    pub fn num_data_files(&self) -> usize {
827        self.metadata.files.len()
828    }
829
830    /// Gets the data file for a given field
831    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    /// Open a FileFragment with a given default projection.
839    ///
840    /// All read operations (other than `read_projected`) will use the supplied
841    /// default projection. For `read_projected`, the projection must be a subset
842    /// of the default projection.
843    ///
844    /// Parameters
845    /// - `projection`: The projection schema.
846    /// - `read_config`: Controls what columns are included in the output.
847    /// - `scan_scheduler`: The scheduler to use for reading data files.  If not supplied
848    ///   and the data is v2 data then a new scheduler will be created
849    ///
850    /// `projection` may be an empty schema only if `with_row_id` is true. In that
851    /// case, the reader will only be generating row ids.
852    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        // The data file may contain fields that are not part of the dataset any longer, remove those
923        let data_file_schema = data_file.schema(full_schema);
924        let projection = projection.unwrap_or(full_schema);
925        // Also remove any fields that are not part of the user's provided projection
926        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                // TODO: make object stores for non-default bases reuse the same scan scheduler
965                //  currently we always create a new one
966                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        // This should return immediately on modern datasets.  Need to use physical_rows because
1049        // deletions will be applied later
1050        let num_rows = self.physical_rows().await?;
1051
1052        // Check if there are any fields that are not in any data files
1053        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    /// Count the rows in this fragment.
1070    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    /// Get the number of rows that have been deleted in this fragment.
1094    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    /// Get the number of physical rows in the fragment synchronously
1112    ///
1113    /// Fails if the fragment does not have the physical row count in the metadata.  This method should
1114    /// only be called in new workflows which are not run on old versions of Lance.
1115    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    /// Get the number of deleted rows in the fragment synchronously
1127    ///
1128    /// Fails if the fragment does not have deletion count in the metadata.  This method should only
1129    /// be called in new workflows which are not run on old versions of Lance.
1130    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    /// Get the number of logical rows (physical rows - deleted rows) in the fragment synchronously
1145    ///
1146    /// Fails if the fragment does not have the physical row count or deletion count in the metadata.  This method should only
1147    /// be called in new workflows which are not run on old versions of Lance.
1148    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    /// Get the number of physical rows in the fragment. This includes deleted rows.
1155    ///
1156    /// If there are no deleted rows, this is equal to the number of rows in the
1157    /// fragment.
1158    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        // Early versions that did not write the writer version also could write
1167        // incorrect `physical_row` values. So if we don't have a writer version,
1168        // we should not used the cached value. On write, we update the values
1169        // in the manifest, fixing the issue for future reads.
1170        // See: https://github.com/lance-format/lance/issues/1531
1171        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        // Just open any file. All of them should have same size.
1176        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    /// Validate the fragment
1191    ///
1192    /// Verifies:
1193    /// * All field ids in the fragment are distinct
1194    /// * Within each data file, field ids are in increasing order
1195    /// * All data files exist and have the same length
1196    /// * Field ids are distinct between data files.
1197    /// * Deletion file exists and has rowids in the correct range
1198    /// * `Fragment.physical_rows` matches length of file
1199    /// * `DeletionFile.num_deleted_rows` matches length of deletion vector
1200    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    /// Open a [`FragmentSession`], which manages a short-lived session of [`FileFragment`].
1341    ///
1342    /// This API works well for users making repeated requests over the same projected schema.
1343    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    /// Take rows from this fragment based on the offset in the file.
1352    ///
1353    /// This will always return the same number of rows as the input indices.
1354    /// If indices are out-of-bounds, this will return an error.
1355    pub async fn take(&self, indices: &[u32], projection: &Schema) -> Result<RecordBatch> {
1356        // Re-map the indices to row ids using the deletion vector
1357        let deletion_vector = self.get_deletion_vector().await?;
1358        let row_ids = if let Some(deletion_vector) = deletion_vector {
1359            // Naive case is O(N*M), where N = indices.len() and M = deletion_vector.len()
1360            // We can do better by sorting the deletion vector and using binary search
1361            // This is O(N * log M + M log M).
1362            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        // Then call take rows
1375        self.take_rows(&row_ids, projection, false, false, false, false)
1376            .await
1377    }
1378
1379    /// Get the deletion vector for this fragment, using the cache if available.
1380    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    /// Get the file metadata for this fragment, using the cache if available.
1392    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    /// Take rows based on internal local row offsets
1410    ///
1411    /// If the row offsets are out-of-bounds, this will return an error. But if the
1412    /// row offset is marked deleted, it will be ignored. Thus, the number of rows
1413    /// returned may be less than the number of row offsets provided.
1414    ///
1415    /// To recover the original row addresses from the returned RecordBatch, set the
1416    /// `with_row_address` parameter to true. This will add a column named `_rowaddr`
1417    /// to the RecordBatch at the end.
1418    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            // FIXME, change this method to streams
1444            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    /// Scan this [`FileFragment`].
1466    ///
1467    /// See [`Dataset::scan`].
1468    pub fn scan(&self) -> Scanner {
1469        Scanner::from_fragment(self.dataset.clone(), self.metadata.clone())
1470    }
1471
1472    /// Create an [`Updater`] to append new columns.
1473    ///
1474    /// The `columns` parameter is a list of existing columns to be read from
1475    /// the fragment. They can be used to derive new columns. This is allowed to
1476    /// be empty.
1477    ///
1478    /// The columns `_rowaddr` and `_rowid` can be used to load the row id or row address
1479    ///
1480    /// The `schemas` parameter is a tuple of the write schema (just the new fields)
1481    /// and the full schema (the target schema after the update). If the write
1482    /// schema is None, it is inferred from the first batch of results. The full
1483    /// schema is inferred by appending the write schema to the existing schema.
1484    ///
1485    /// The `batch_size` parameter can be used to influence how much data is processed
1486    /// at a time. This can be useful to control memory usage when processing very large
1487    /// fields. The batch_size will only be used if the dataset is a v2 dataset.  It will
1488    /// be ignored for v1 datasets.
1489    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        // If there is no projection, we at least need to read the row addresses
1514        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                // right_on is allowed to exist in the dataset, since it may be
1555                // the same as left_on.
1556                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        // Hash join
1566        let joiner = Arc::new(HashJoiner::try_new(stream, right_on).await?);
1567        // Final schema is union of current schema, plus the RHS schema without
1568        // the right_on key.
1569        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        // Prepare the read projection: align with the write_schema's columns and append the left_on column.
1638        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        // Hash join
1649        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        // Mark fields in updated data files as obsolete ("tombstone").
1659        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                    // Tombstone these fields
1664                    *field = -2;
1665                }
1666            }
1667        }
1668        // Remove data files that have become entirely tombstoned.
1669        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        // Note: updated field should be returned when committing, waiting to be done
1677        Ok((updated_fragment, updated_fields))
1678    }
1679
1680    /// Append new columns to the fragment
1681    ///
1682    /// This is the fragment-level version of [`Dataset::add_columns`].
1683    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    /// Delete rows from the fragment.
1702    ///
1703    /// If all rows are deleted, returns `Ok(None)`. Otherwise, returns a new
1704    /// fragment with the updated deletion vector. This must be persisted to
1705    /// the manifest.
1706    pub async fn delete(self, predicate: &str) -> Result<Option<Self>> {
1707        // Load existing deletion vector
1708        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        // scan with predicate and row addresses
1718        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 predicate is `true`, delete the whole fragment
1733        // else if predicate is `false`, filter the predicate
1734        // We do this on the expression level after expression optimization has
1735        // occurred so we also catch expressions that are equivalent to `true`
1736        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        // As we get row addrs, add them into our deletion vector
1752        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                // _rowaddr is global, not within fragment level. The high bits
1760                // are the fragment_id, the low bits are the row_id within the
1761                // fragment.
1762                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 we haven't deleted any additional rows, we can return the fragment as-is.
1770        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
1827/// Using deleted ids to remap row ids into actual row ids.
1828pub(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        // We find the number of deleted rows that are less than each row
1832        // index, and that becomes the initial offset. We increment the
1833        // index by that amount, plus the number of deleted row ids we
1834        // encounter along the way. So for example, if deleted rows are
1835        // [2, 3, 5] and we want row 4, we need to advanced by 2 (since
1836        // 2 and 3 are less than 4). That puts us at row 6, but since
1837        // we passed row 5, we need to advance by 1 more, giving a final
1838        // row id of 7.
1839        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            // Advance the row id
1846            new_row_id += 1;
1847            while deletion_i < sorted_deleted_ids.len()
1848                && sorted_deleted_ids[deletion_i] == new_row_id
1849            {
1850                // If we encounter a deleted row, we need to advance
1851                // again.
1852                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// Cache key for file metadata
1865#[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/// [`FragmentReader`] is an abstract reader for a [`FileFragment`].
1883///
1884/// It opens the data files that contains the columns of the projection schema, and
1885/// reconstruct the RecordBatch from columns read from each data file.
1886#[derive(Debug)]
1887pub struct FragmentReader {
1888    /// Readers and schema of each opened data file.
1889    readers: Vec<Box<dyn GenericFileReader>>,
1890
1891    /// The output schema. The defines the order in which the columns are returned.
1892    output_schema: ArrowSchema,
1893
1894    /// The deleted row IDs
1895    deletion_vec: Option<Arc<DeletionVector>>,
1896
1897    /// The row id sequence
1898    ///
1899    /// Only populated if the stable row id feature is enabled.
1900    row_id_sequence: Option<Arc<RowIdSequence>>,
1901
1902    /// ID of the fragment
1903    fragment_id: usize,
1904
1905    /// True if we should generate a row id for the output
1906    with_row_id: bool,
1907
1908    /// True if we should generate a row address column in output
1909    with_row_addr: bool,
1910
1911    /// True if we should generate a last updated at version column in output
1912    with_row_last_updated_at_version: bool,
1913
1914    /// True if we should generate a created at version column in output
1915    with_row_created_at_version: bool,
1916
1917    /// If true, deleted rows will be set to null, which is fast
1918    /// If false, deleted rows will be removed from the batch, requiring a copy
1919    make_deletions_null: bool,
1920
1921    /// The fragment metadata (needed for version columns)
1922    fragment: Arc<Fragment>,
1923
1924    /// The last_updated_at version sequence (loaded from fragment metadata)
1925    last_updated_at_sequence: Option<Arc<lance_table::rowids::version::RowDatasetVersionSequence>>,
1926
1927    /// The created_at version sequence (loaded from fragment metadata)
1928    created_at_sequence: Option<Arc<lance_table::rowids::version::RowDatasetVersionSequence>>,
1929
1930    // total number of real rows in the fragment (num_physical_rows - num_deleted_rows)
1931    num_rows: usize,
1932
1933    // total number of physical rows in the fragment (all rows, ignoring deletions)
1934    num_physical_rows: usize,
1935}
1936
1937// Custom clone impl needed because it is not easy to clone Box<dyn GenericFileReader>
1938//
1939// We currently need FragmentReader to be Clone because the pushdown scan clones it
1940// to reuse the fragment reader for both "scan with row id" and "scan without row id"
1941impl 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        // Load the version sequence if not already loaded
2060        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        // If no metadata or load fails, sequence remains None (will default to version 1)
2067
2068        // Add the version column to the output schema
2069        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        // Load the version sequence if not already loaded
2081        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        // If no metadata or load fails, sequence remains None (will default to version 1)
2088
2089        // Add the version column to the output schema
2090        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    /// TODO: This method is relied upon by the v1 pushdown mechanism and will need to stay
2099    /// in place until v1 is removed.  v2 uses a different mechanism for pushdown and so there
2100    /// is little benefit in updating the v1 pushdown node.
2101    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    /// TODO: This method is relied upon by the v1 pushdown mechanism and will need to stay
2114    /// in place until v1 is removed.  v2 uses a different mechanism for pushdown and so there
2115    /// is little benefit in updating the v1 pushdown node.
2116    ///
2117    /// This method is also used by the updater.  Even though the updater has been updated to
2118    /// use streams, the updater still needs to know the batch size in v1 so that it can create
2119    /// files with the same batch size.
2120    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    /// Read the page statistics of the fragment for the specified fields.
2133    ///
2134    /// TODO: This method is relied upon by the v1 pushdown mechanism and will need to stay
2135    /// in place until v1 is removed.  v2 uses a different mechanism for pushdown and so there
2136    /// is little benefit in updating the v1 pushdown node.
2137    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    /// Read a batch of rows from the fragment, with a subset of columns.
2161    ///
2162    /// Note: the projection must be a subset of the schema the reader was created with.
2163    /// Otherwise incorrect data will be returned.
2164    ///
2165    /// TODO: This method is relied upon by the v1 pushdown mechanism and will need to stay
2166    /// in place until v1 is removed.  v2 uses a different mechanism for pushdown and so there
2167    /// is little benefit in updating the v1 pushdown node.
2168    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        // All batches have the same size in v1, except for the last one.
2176        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                    // Apply ? inside the task to keep read_tasks a simple iter of futures
2188                    // for try_join_all
2189                    let projection = projection?;
2190                    if projection.fields.is_empty() {
2191                        // The projection caused one of the data files to become
2192                        // irrelevant and so we can skip it
2193                        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            // If we are selecting no columns, we can assume we are just getting
2207            // the row ids. If this is the case, we need to generate an empty
2208            // batch with the correct number of rows.
2209            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        // Need to apply deletions and row ids.
2226        // In order to apply deletions we need to change the parameters to be
2227        // relative to the file, not the batch.
2228        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        // Note that the fragment length might be considerably smaller if there are deleted rows.
2295        // E.g. if a fragment has 100 rows but rows 0..10 are deleted we still need to make
2296        // sure it is valid to read / take 0..100
2297        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        // If just the row id or address there is no need to actually read any data
2304        // and we don't need to involve the readers at all.
2305        //
2306        // The v1 reader does not support reading batches with zero columns, so
2307        // we need this as a separate code path.
2308        // In these cases, we can just emit batches with zero columns and rely
2309        // on `wrap_with_row_id_and_delete` to add the row id or address column.
2310        //
2311        // We could potentially delete the support for no-columns in the wrap function or
2312        // we can delete this path once we migrate away from any support of v1.
2313        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            // Read each data file, these reads should produce streams of equal sized
2328            // tasks.  In other words, if we get 3 tasks of 20 rows and then a task
2329            // of 10 rows from one data file we should get the same from the other.
2330            let read_streams = self
2331                .readers
2332                .iter()
2333                .filter_map(|reader| {
2334                    // Normally we filter out empty readers in the open_readers method
2335                    // However, we will keep the first empty reader to use for row id
2336                    // purposes on some legacy paths and so we need to filter that out
2337                    // here.
2338                    if reader.projection().fields.is_empty() {
2339                        None
2340                    } else {
2341                        Some(read_fn(reader.as_ref()))
2342                    }
2343                })
2344                .collect::<Result<Vec<_>>>()?;
2345            // Merge the streams, this merges the generated batches
2346            lance_table::utils::stream::merge_streams(read_streams)
2347        };
2348
2349        // Add the row id column (if needed) and delete rows (if a deletion
2350        // vector is present).
2351        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                // Finally, reorder the columns to match the order specified in the projection
2368                .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    /// Reads a range of rows from the fragment
2428    ///
2429    /// This function interprets the request as the Xth to the Nth row of the fragment (after deletions)
2430    /// and will always return range.len().min(self.num_rows()) rows.
2431    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    /// Takes a range of rows from the fragment
2436    ///
2437    /// Unlike [`Self::read_range`], this function will NOT skip deleted rows.  If rows are deleted they will
2438    /// be filtered or set to null.  This function may return less than range.len() rows as a result.
2439    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    // This method is a clone of new_read_impl but returns tasks instead of batches
2450    //
2451    // It also only supports v2 files
2452    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        // Note that row ranges at this point are physical and not logical.
2460        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            // Read each data file, these reads should produce streams of equal sized
2485            // tasks.  In other words, if we get 3 tasks of 20 rows and then a task
2486            // of 10 rows from one data file we should get the same from the other.
2487            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            // Merge the streams, this merges the generated batches
2499            lance_table::utils::stream::merge_streams(read_streams)
2500        };
2501
2502        // Add the row id column (if needed) and delete rows (if a deletion
2503        // vector is present).
2504        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                // Finally, reorder the columns to match the order specified in the projection
2521                .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    // Legacy function that reads a range of data and concatenates the results
2536    // into a single batch
2537    //
2538    // TODO: Move away from this by changing callers to support consuming a stream
2539    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    /// Take rows from this fragment.
2552    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    /// Take rows from this fragment, will perform a copy if the underlying reader returns multiple
2574    /// batches.  May return an error if the taken rows do not fit into a single batch.
2575    ///
2576    /// Duplicate indices are allowed and will produce duplicate rows in the output.
2577    pub async fn take_as_batch(
2578        &self,
2579        indices: &[u32],
2580        take_priority: Option<u32>,
2581    ) -> Result<RecordBatch> {
2582        // The v2 encoding layer requires strictly increasing indices. Deduplicate
2583        // here so callers (e.g. FTS with duplicate row matches) don't need to.
2584        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        // Test update with _rowid
2807        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        // Test update with user specified keys
2878        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        // Creates 400 rows in 10 fragments
2964        let mut dataset = create_dataset(test_uri, LanceFileVersion::Legacy).await;
2965        // Delete last 20 rows in first fragment
2966        dataset.delete("i >= 20").await.unwrap();
2967        // Last fragment has 20 rows but 40 addressable rows
2968        let fragment = &dataset.get_fragments()[0];
2969        assert_eq!(fragment.metadata.num_rows().unwrap(), 20);
2970
2971        // Test with take_range (all rows addressable)
2972        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        // Test with read_range (only non-deleted rows addressable)
2995        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        // Creates 400 rows in 10 fragments
3023        let mut dataset = create_dataset(test_uri, LanceFileVersion::Legacy).await;
3024        // Delete last 20 rows in first fragment
3025        dataset.delete("i >= 20").await.unwrap();
3026        // Last fragment has 20 rows but 40 addressable rows
3027        let fragment = &dataset.get_fragments()[0];
3028        assert_eq!(fragment.metadata.num_rows().unwrap(), 20);
3029
3030        // Test with take_range (all rows addressable)
3031        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        // Test with read_range (only non-deleted rows addressable)
3056        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            // The first batch is entirely deleted, deleted rows will be marked null with null row ids.
3104            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            // The second batch is partially deleted, so the deleted rows will be
3114            // marked null with null row ids.
3115            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            // The final batch is not deleted, so it will be returned as-is.
3125            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            // Since the first batch is all deleted, it will return all nulls row ids.
3144            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            // The second batch is partially deleted, so the deleted rows will be
3156            // marked null with null row ids.
3157            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            // The final batch is not deleted, so it will be returned as-is.
3163            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            // Using batch_size=20 here.  If we use batch_size=range.len() we get
3200            // multiple batches because we might have to read from a larger range
3201            // to satisfy the request
3202            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        // Deleting from the start
3217        check("i < 5", 0..2, vec![5, 6]).await;
3218        check("i < 5", 0..15, (5..20).collect()).await;
3219        // Deleting from the middle
3220        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        // Deleting from the end
3232        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        // Repeated indices are repeated in result.
3252        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        // Deleted rows are skipped
3265        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        // Empty indices gives empty result
3281        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        // Repeated indices are repeated in result.
3304        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        // Cannot get rows 2 and 4 anymore
3324        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        // Empty indices gives empty result
3347        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        // Can get row ids
3357        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        // Fragments will have number of rows recorded in metadata, even though
3417        // we passed `None` when constructing the `FileFragment`.
3418        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            // Scan again
3508            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            // We only kept the first fragment of 40 rows
3525            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                // Legacy format uses old scan node which deletes after read and
3546                // so the batch is truncated
3547                (true, LanceFileVersion::Legacy) => vec![10, 11, 12, 13, 14],
3548                // Newer formats delete before read and so we get a full batch of 10
3549                (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        // Create data to merge: merge in double the data
3587        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        // Validate the resulting data
3605        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        // V1 ONLY
3636        //
3637        // This test is only for the legacy version of the file format.
3638        // It ensures that the `max_rows_per_group` property is respected
3639        // and this property does not exist in V2.
3640        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        // Validates we can handle datasets where the order of columns is not
3704        // aligned with the order of the data files. This can happen when replacing
3705        // columns in a dataset.
3706        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        // Write batch_i as a fragment
3727        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        // Write batch_s using add_columns
3740        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        // Rearrange schema so it's `s` then `i`.
3746        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        // Also take, read_range, and read_batch_projected
3767        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        // Also check case of row_id.
3787        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        // Make sure we can create a fragment reader that only captures the row_id.
3800        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        // We should get error if we pass empty schema and with_row_id false
3836        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        // Create a file that has 8 columns.
3910        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        // Single row batch
3917        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        // Assert file is small (< 4300 bytes)
3937        {
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        // Measure IOPS needed to scan all data first time.
3944        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}