Skip to main content

common_datasource/file_format/parquet/
packed_reader.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Request-owned, bounded backing buffers for independent Parquet streams.
16
17use std::ops::Range;
18use std::sync::Arc;
19
20use bytes::Bytes;
21use futures::future::BoxFuture;
22use object_store::ObjectStore;
23use parquet::arrow::arrow_reader::ArrowReaderOptions;
24use parquet::arrow::async_reader::AsyncFileReader;
25use parquet::errors::{ParquetError, Result};
26use parquet::file::metadata::{ParquetMetaData, ParquetMetaDataReader};
27use tokio::sync::Mutex;
28
29pub const WINDOW_SIZE: usize = 8 * 1024 * 1024;
30const WINDOW_BUDGET: usize = 2 * WINDOW_SIZE;
31
32struct Window {
33    path: String,
34    offset: u64,
35    bytes: Bytes,
36}
37
38/// Shared only by table decoders belonging to one authorized COPY request.
39/// The lock also coalesces concurrent fetches for the same bytes.
40pub struct PackReadWindows {
41    store: ObjectStore,
42    windows: Mutex<Vec<Window>>,
43}
44
45impl PackReadWindows {
46    pub fn new(store: ObjectStore) -> Arc<Self> {
47        Arc::new(Self {
48            store,
49            windows: Mutex::new(Vec::new()),
50        })
51    }
52
53    async fn read(&self, path: &str, object_length: u64, range: Range<u64>) -> Result<Bytes> {
54        if range.start > range.end || range.end > object_length {
55            return Err(ParquetError::General("packed read outside object".into()));
56        }
57        if range.is_empty() {
58            return Ok(Bytes::new());
59        }
60        let mut windows = self.windows.lock().await;
61        for window in windows.iter() {
62            if window.path == path
63                && range.start >= window.offset
64                && range.end <= window.offset + window.bytes.len() as u64
65            {
66                return Ok(window.bytes.slice(
67                    (range.start - window.offset) as usize..(range.end - window.offset) as usize,
68                ));
69            }
70        }
71        let required = usize::try_from(range.end - range.start)
72            .map_err(|_| ParquetError::General("packed read length overflow".into()))?;
73        if required > WINDOW_SIZE {
74            return Err(ParquetError::General(
75                "packed cache range exceeds window".into(),
76            ));
77        }
78        let fetch_len = (object_length - range.start).min(WINDOW_SIZE as u64) as usize;
79        // Evict only unpinned allocations. Bytes slices held by a decoder still
80        // count against the request budget, including their entire backing window.
81        let mut retained: usize = windows.iter().map(|w| w.bytes.len()).sum();
82        while retained.saturating_add(fetch_len) > WINDOW_BUDGET || windows.len() >= 2 {
83            let Some(i) = windows.iter().position(|w| w.bytes.is_unique()) else {
84                return Err(ParquetError::General(
85                    "packed read backing-buffer budget exhausted".into(),
86                ));
87            };
88            retained -= windows.remove(i).bytes.len();
89        }
90        let bytes = self
91            .store
92            .read_with(path)
93            .range(range.start..range.start + fetch_len as u64)
94            .await
95            .map_err(|e| ParquetError::External(Box::new(e)))?
96            .to_bytes();
97        if bytes.len() != fetch_len {
98            return Err(ParquetError::General("short packed object read".into()));
99        }
100        // Backend buffers can be slices of larger allocations. Retain an exact
101        // window allocation; transport staging is separate from cached backing.
102        let bytes = Bytes::copy_from_slice(&bytes);
103        common_telemetry::debug!(
104            path,
105            offset = range.start,
106            length = fetch_len,
107            "Fetched packed read window"
108        );
109        let result = bytes.slice(..required);
110        windows.push(Window {
111            path: path.into(),
112            offset: range.start,
113            bytes,
114        });
115        Ok(result)
116    }
117}
118
119/// Pins one complete small stream; larger streams use the streaming reader.
120pub struct PackedParquetReader {
121    windows: Arc<PackReadWindows>,
122    path: String,
123    object_length: u64,
124    offset: u64,
125    length: u64,
126    bytes: Option<Bytes>,
127}
128
129impl PackedParquetReader {
130    pub fn new(
131        windows: Arc<PackReadWindows>,
132        path: String,
133        object_length: u64,
134        offset: u64,
135        length: u64,
136    ) -> Result<Self> {
137        if length < 12
138            || length > WINDOW_SIZE as u64
139            || offset
140                .checked_add(length)
141                .is_none_or(|end| end > object_length)
142        {
143            return Err(ParquetError::General("invalid packed Parquet range".into()));
144        }
145        Ok(Self {
146            windows,
147            path,
148            object_length,
149            offset,
150            length,
151            bytes: None,
152        })
153    }
154}
155
156impl AsyncFileReader for PackedParquetReader {
157    fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, Result<Bytes>> {
158        Box::pin(async move {
159            if range.start > range.end || range.end > self.length {
160                return Err(ParquetError::General(
161                    "read outside indexed Parquet stream".into(),
162                ));
163            }
164            if self.bytes.is_none() {
165                self.bytes = Some(
166                    self.windows
167                        .read(
168                            &self.path,
169                            self.object_length,
170                            self.offset..self.offset + self.length,
171                        )
172                        .await?,
173                );
174            }
175            Ok(self
176                .bytes
177                .as_ref()
178                .unwrap()
179                .slice(range.start as usize..range.end as usize))
180        })
181    }
182
183    fn get_metadata<'a>(
184        &'a mut self,
185        options: Option<&'a ArrowReaderOptions>,
186    ) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
187        Box::pin(async move {
188            let length = self.length;
189            let metadata = ParquetMetaDataReader::new()
190                .with_metadata_options(options.map(|o| o.metadata_options().clone()))
191                .load_and_finish(self, length)
192                .await?;
193            Ok(Arc::new(metadata))
194        })
195    }
196}
197
198#[cfg(test)]
199mod tests {
200    use arrow::array::{ArrayRef, Int64Array, StringArray};
201    use arrow::record_batch::RecordBatch;
202    use arrow_schema::{Field, Schema};
203    use futures::TryStreamExt;
204    use parquet::arrow::{ArrowWriter, ParquetRecordBatchStreamBuilder};
205
206    use super::*;
207
208    fn parquet(array: ArrayRef) -> (Vec<u8>, RecordBatch) {
209        let schema = Arc::new(Schema::new(vec![Field::new(
210            "value",
211            array.data_type().clone(),
212            false,
213        )]));
214        let batch = RecordBatch::try_new(schema.clone(), vec![array]).unwrap();
215        let mut writer = ArrowWriter::try_new(Vec::new(), schema, None).unwrap();
216        writer.write(&batch).unwrap();
217        (writer.into_inner().unwrap(), batch)
218    }
219
220    #[tokio::test]
221    async fn independent_schemas_empty_streams_and_shared_ranges() {
222        let dir = common_test_util::temp_dir::create_temp_dir("packed-reader");
223        let store = crate::test_util::test_store(dir.path().to_str().unwrap());
224        let fixtures = [
225            parquet(Arc::new(Int64Array::from(vec![42]))),
226            parquet(Arc::new(StringArray::from(vec!["different type"]))),
227            parquet(Arc::new(Int64Array::from(Vec::<i64>::new()))),
228        ];
229        let packed: Vec<u8> = fixtures
230            .iter()
231            .flat_map(|(b, _)| b.iter().copied())
232            .collect();
233        store.write("pack-0.bin", packed.clone()).await.unwrap();
234        let windows = PackReadWindows::new(store.clone());
235        // A pinned shared window serves all the independent stream readers.
236        let pinned = windows
237            .read("pack-0.bin", packed.len() as u64, 0..packed.len() as u64)
238            .await
239            .unwrap();
240        store.delete("pack-0.bin").await.unwrap();
241        let mut offset = 0;
242        for (bytes, expected) in fixtures {
243            let reader = PackedParquetReader::new(
244                windows.clone(),
245                "pack-0.bin".into(),
246                packed.len() as u64,
247                offset,
248                bytes.len() as u64,
249            )
250            .unwrap();
251            let builder = ParquetRecordBatchStreamBuilder::new(reader).await.unwrap();
252            assert_eq!(builder.schema(), &expected.schema());
253            let batches: Vec<_> = builder.build().unwrap().try_collect().await.unwrap();
254            assert_eq!(
255                batches.iter().map(|b| b.num_rows()).sum::<usize>(),
256                expected.num_rows()
257            );
258            if expected.num_rows() > 0 {
259                assert_eq!(batches[0], expected);
260            }
261            offset += bytes.len() as u64;
262        }
263        assert_eq!(windows.windows.lock().await.len(), 1);
264        assert_eq!(pinned.len(), packed.len());
265        let isolated = PackReadWindows::new(store);
266        assert!(
267            isolated
268                .read("pack-0.bin", packed.len() as u64, 0..1)
269                .await
270                .is_err()
271        );
272    }
273
274    #[tokio::test]
275    async fn small_stream_pins_one_window_for_multiple_columns() {
276        use parquet::file::properties::WriterProperties;
277        let schema = Arc::new(Schema::new(
278            (0..3)
279                .map(|i| Field::new(format!("v{i}"), arrow_schema::DataType::Utf8, false))
280                .collect::<Vec<_>>(),
281        ));
282        let value = "x".repeat(1024 * 1024);
283        let arrays = (0..3)
284            .map(|_| Arc::new(StringArray::from(vec![value.as_str()])) as ArrayRef)
285            .collect();
286        let batch = RecordBatch::try_new(schema.clone(), arrays).unwrap();
287        let mut writer = ArrowWriter::try_new(
288            Vec::new(),
289            schema,
290            Some(
291                WriterProperties::builder()
292                    .set_dictionary_enabled(false)
293                    .build(),
294            ),
295        )
296        .unwrap();
297        writer.write(&batch).unwrap();
298        let bytes = writer.into_inner().unwrap();
299        assert!(bytes.len() < WINDOW_SIZE);
300        let dir = common_test_util::temp_dir::create_temp_dir("packed-columns");
301        let store = crate::test_util::test_store(dir.path().to_str().unwrap());
302        let mut object = vec![0; WINDOW_SIZE - 1024];
303        let offset = object.len() as u64;
304        object.extend_from_slice(&bytes);
305        let length = object.len() as u64;
306        store.write("pack-0.bin", object).await.unwrap();
307        let windows = PackReadWindows::new(store);
308        let pinned = windows.read("pack-0.bin", length, 0..1).await.unwrap();
309        let mut reader = PackedParquetReader::new(
310            windows.clone(),
311            "pack-0.bin".into(),
312            length,
313            offset,
314            bytes.len() as u64,
315        )
316        .unwrap();
317        let metadata = reader.get_metadata(None).await.unwrap();
318        let mut columns = Vec::new();
319        for column in metadata.row_group(0).columns() {
320            let (start, len) = column.byte_range();
321            columns.push(reader.get_bytes(start..start + len).await.unwrap());
322        }
323        let batches: Vec<_> = ParquetRecordBatchStreamBuilder::new(reader)
324            .await
325            .unwrap()
326            .build()
327            .unwrap()
328            .try_collect()
329            .await
330            .unwrap();
331        assert_eq!(batches, vec![batch]);
332        let retained = windows.windows.lock().await;
333        assert_eq!(retained.len(), 2);
334        assert!(retained.iter().all(|w| w.bytes.len() <= WINDOW_SIZE));
335        assert!(retained.iter().map(|w| w.bytes.len()).sum::<usize>() <= WINDOW_BUDGET);
336        assert_eq!(columns.len(), 3);
337        assert_eq!(pinned[0], 0);
338    }
339
340    #[tokio::test]
341    async fn corrupt_metadata_short_reads_and_relative_bounds() {
342        let dir = common_test_util::temp_dir::create_temp_dir("packed-corrupt");
343        let store = crate::test_util::test_store(dir.path().to_str().unwrap());
344        store.write("pack-0.bin", vec![0u8; 12]).await.unwrap();
345        let windows = PackReadWindows::new(store);
346        assert!(windows.read("pack-0.bin", 20, 0..20).await.is_err());
347        let mut reader = PackedParquetReader::new(windows, "pack-0.bin".into(), 12, 0, 12).unwrap();
348        assert!(reader.get_bytes(0..13).await.is_err());
349        assert!(reader.get_metadata(None).await.is_err());
350    }
351}