Skip to main content

lance_file/
io.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use std::sync::Arc;
5
6use futures::{FutureExt, future::BoxFuture};
7use lance_encoding::EncodingsIo;
8use lance_io::scheduler::FileScheduler;
9
10use super::reader::DEFAULT_READ_CHUNK_SIZE;
11
12#[derive(Debug)]
13pub struct LanceEncodingsIo {
14    scheduler: FileScheduler,
15    /// Size of chunks when reading large pages
16    read_chunk_size: u64,
17}
18
19impl LanceEncodingsIo {
20    pub fn new(scheduler: FileScheduler) -> Self {
21        Self {
22            scheduler,
23            read_chunk_size: DEFAULT_READ_CHUNK_SIZE,
24        }
25    }
26
27    pub fn with_read_chunk_size(mut self, read_chunk_size: u64) -> Self {
28        self.read_chunk_size = read_chunk_size;
29        self
30    }
31}
32
33impl EncodingsIo for LanceEncodingsIo {
34    fn with_bypass_backpressure(&self) -> Option<Arc<dyn EncodingsIo>> {
35        Some(Arc::new(Self {
36            scheduler: self.scheduler.with_bypass_backpressure(),
37            read_chunk_size: self.read_chunk_size,
38        }))
39    }
40
41    fn with_io_stats(
42        &self,
43        stats: Arc<dyn lance_core::utils::io_stats::IoStatsRecorder>,
44    ) -> Option<Arc<dyn EncodingsIo>> {
45        Some(Arc::new(Self {
46            scheduler: self.scheduler.with_io_stats(stats),
47            read_chunk_size: self.read_chunk_size,
48        }))
49    }
50
51    fn submit_request(
52        &self,
53        ranges: Vec<std::ops::Range<u64>>,
54        priority: u64,
55    ) -> BoxFuture<'static, lance_core::Result<Vec<bytes::Bytes>>> {
56        let mut split_ranges = Vec::new();
57        let mut split_indices = Vec::new(); // Track which original range each split came from
58        // Large ranges (above read_chunk_size) will be split into
59        // multiple reads.  Empty ranges will skip the I/O layer
60        // entirely.  If we have either of these we will need to
61        // reassemble our results, inserting empties and merging parts
62        let mut needs_reassembly = false;
63
64        // Split large ranges into smaller chunks
65        //
66        // TODO: consider read_chunk_size before submitting requests.
67        for (idx, range) in ranges.iter().enumerate() {
68            if range.start == range.end {
69                // EncodingsIo requires one result per input range. Zero-length
70                // ranges schedule no I/O, so their empty results are restored
71                // after the non-empty requests complete.
72                needs_reassembly = true;
73                continue;
74            }
75            let range_size = range.end - range.start;
76
77            if range_size > self.read_chunk_size {
78                needs_reassembly = true;
79                let num_chunks = range_size.div_ceil(self.read_chunk_size);
80                let chunk_size = range_size / num_chunks;
81
82                for i in 0..num_chunks {
83                    let start = range.start + i * chunk_size;
84                    let end = if i == num_chunks - 1 {
85                        range.end // Last chunk gets any remaining bytes
86                    } else {
87                        start + chunk_size
88                    };
89                    split_ranges.push(start..end);
90                    split_indices.push(idx);
91                }
92            } else {
93                split_ranges.push(range.clone());
94                split_indices.push(idx);
95            }
96        }
97
98        let fut = self.scheduler.submit_request(split_ranges, priority);
99
100        async move {
101            let split_results = fut.await?;
102
103            if split_results.len() != split_indices.len() {
104                return Err(lance_core::Error::internal(format!(
105                    "Encoding I/O returned {} results for {} requested range chunks",
106                    split_results.len(),
107                    split_indices.len()
108                )));
109            }
110            if !needs_reassembly {
111                return Ok(split_results);
112            }
113
114            let mut results = vec![Vec::new(); ranges.len()];
115
116            for (split_result, orig_idx) in split_results.into_iter().zip(split_indices) {
117                results[orig_idx].push(split_result);
118            }
119
120            let mut reassembled = Vec::with_capacity(ranges.len());
121            for (range, chunks) in ranges.iter().zip(results) {
122                if chunks.is_empty() {
123                    if range.start == range.end {
124                        reassembled.push(bytes::Bytes::new());
125                        continue;
126                    }
127                    return Err(lance_core::Error::internal(format!(
128                        "Encoding I/O returned no data for non-empty range {}..{}",
129                        range.start, range.end
130                    )));
131                }
132                if chunks.len() == 1 {
133                    reassembled.push(chunks[0].clone());
134                    continue;
135                }
136
137                let total_size: usize = chunks.iter().map(|c| c.len()).sum();
138                let mut combined = Vec::with_capacity(total_size);
139                for chunk in chunks {
140                    combined.extend_from_slice(&chunk);
141                }
142                reassembled.push(bytes::Bytes::from(combined));
143            }
144            Ok(reassembled)
145        }
146        .boxed()
147    }
148}