common_datasource/file_format/parquet/
packed_reader.rs1use 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
38pub 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 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 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
119pub 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 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}