1use std::sync::Arc;
19
20use futures::{StreamExt, TryStreamExt};
21use parquet::arrow::ParquetRecordBatchStreamBuilder;
22
23use crate::arrow::ArrowReader;
24use crate::arrow::reader::ParquetReadOptions;
25use crate::arrow::record_batch_transformer::RecordBatchTransformerBuilder;
26use crate::arrow::scan_metrics::ScanMetrics;
27use crate::io::FileIO;
28use crate::scan::{ArrowRecordBatchStream, FileScanTaskDeleteFile};
29use crate::spec::{Schema, SchemaRef};
30use crate::{Error, ErrorKind, Result};
31
32#[allow(unused)]
34#[async_trait::async_trait]
35pub trait DeleteFileLoader {
36 async fn read_delete_file(
40 &self,
41 task: &FileScanTaskDeleteFile,
42 schema: SchemaRef,
43 ) -> Result<ArrowRecordBatchStream>;
44}
45
46#[derive(Clone, Debug)]
47pub(crate) struct BasicDeleteFileLoader {
48 file_io: FileIO,
49 scan_metrics: ScanMetrics,
50}
51
52#[allow(unused_variables)]
53impl BasicDeleteFileLoader {
54 pub fn new(file_io: FileIO, scan_metrics: ScanMetrics) -> Self {
55 BasicDeleteFileLoader {
56 file_io,
57 scan_metrics,
58 }
59 }
60
61 pub(crate) fn file_io(&self) -> &FileIO {
62 &self.file_io
63 }
64
65 pub(crate) async fn parquet_to_batch_stream(
67 &self,
68 data_file_path: &str,
69 file_size_in_bytes: u64,
70 key_metadata: Option<&[u8]>,
71 ) -> Result<ArrowRecordBatchStream> {
72 let parquet_read_options = ParquetReadOptions::builder().build();
77
78 let (parquet_file_reader, arrow_metadata) = ArrowReader::open_parquet_file(
79 data_file_path,
80 &self.file_io,
81 file_size_in_bytes,
82 parquet_read_options,
83 self.scan_metrics.bytes_read_counter(),
84 key_metadata,
85 )
86 .await?;
87
88 let record_batch_stream =
89 ParquetRecordBatchStreamBuilder::new_with_metadata(parquet_file_reader, arrow_metadata)
90 .build()?
91 .map_err(|e| Error::new(ErrorKind::Unexpected, format!("{e}")));
92
93 Ok(Box::pin(record_batch_stream) as ArrowRecordBatchStream)
94 }
95
96 pub(crate) async fn evolve_schema(
101 record_batch_stream: ArrowRecordBatchStream,
102 target_schema: Arc<Schema>,
103 equality_ids: &[i32],
104 ) -> Result<ArrowRecordBatchStream> {
105 let mut record_batch_transformer =
106 RecordBatchTransformerBuilder::new(target_schema.clone(), equality_ids).build();
107
108 let record_batch_stream = record_batch_stream.map(move |record_batch| {
109 record_batch.and_then(|record_batch| {
110 record_batch_transformer.process_record_batch(record_batch)
111 })
112 });
113
114 Ok(Box::pin(record_batch_stream) as ArrowRecordBatchStream)
115 }
116}
117
118#[async_trait::async_trait]
119impl DeleteFileLoader for BasicDeleteFileLoader {
120 async fn read_delete_file(
121 &self,
122 task: &FileScanTaskDeleteFile,
123 schema: SchemaRef,
124 ) -> Result<ArrowRecordBatchStream> {
125 let raw_batch_stream = self
126 .parquet_to_batch_stream(
127 &task.file_path,
128 task.file_size_in_bytes,
129 task.key_metadata.as_deref(),
130 )
131 .await?;
132
133 let field_ids = match &task.equality_ids {
136 Some(ids) => ids.clone(),
137 None => schema.field_id_to_name_map().keys().cloned().collect(),
138 };
139
140 Self::evolve_schema(raw_batch_stream, schema, &field_ids).await
141 }
142}
143
144#[cfg(test)]
145mod tests {
146 use tempfile::TempDir;
147
148 use super::*;
149 use crate::arrow::delete_filter::tests::setup;
150 use crate::arrow::test_utils::write_encrypted_parquet;
151
152 #[tokio::test]
153 async fn test_basic_delete_file_loader_read_delete_file() {
154 let tmp_dir = TempDir::new().unwrap();
155 let table_location = tmp_dir.path();
156 let file_io = FileIO::new_with_fs();
157
158 let scan_metrics = ScanMetrics::new();
159 let delete_file_loader = BasicDeleteFileLoader::new(file_io.clone(), scan_metrics);
160
161 let file_scan_tasks = setup(table_location);
162
163 let result = delete_file_loader
164 .read_delete_file(
165 &file_scan_tasks[0].deletes()[0],
166 file_scan_tasks[0].schema_ref(),
167 )
168 .await
169 .unwrap();
170
171 let result = result.try_collect::<Vec<_>>().await.unwrap();
172
173 assert_eq!(result.len(), 1);
174 }
175
176 #[tokio::test]
177 async fn test_read_encrypted_positional_delete_file() {
178 use std::sync::Arc;
179
180 use arrow_array::{Int64Array, RecordBatch, StringArray};
181
182 use crate::arrow::delete_filter::tests::create_pos_del_schema;
183 use crate::encryption::StandardKeyMetadata;
184 use crate::scan::FileScanTaskDeleteFile;
185 use crate::spec::{DataContentType, DataFileFormat};
186
187 let encryption_key = b"0123456789abcdef";
188 let aad_prefix = b"aad_prefix";
189
190 let tmp_dir = TempDir::new().unwrap();
191 let table_location = tmp_dir.path().to_str().unwrap();
192 let file_io = FileIO::new_with_fs();
193
194 let positional_delete_schema = create_pos_del_schema();
195 let file_path_col = Arc::new(StringArray::from_iter_values(vec!["data.parquet"; 4]));
196 let pos_col = Arc::new(Int64Array::from(vec![0i64, 1, 5, 10]));
197 let batch = RecordBatch::try_new(positional_delete_schema.clone(), vec![
198 file_path_col,
199 pos_col,
200 ])
201 .unwrap();
202
203 let del_path = format!("{table_location}/encrypted-pos-del.parquet");
204 write_encrypted_parquet(&del_path, &batch, encryption_key, Some(aad_prefix));
205
206 let key_metadata = StandardKeyMetadata::try_new(encryption_key)
207 .unwrap()
208 .with_aad_prefix(aad_prefix)
209 .encode()
210 .unwrap();
211
212 let schema = Arc::new(
213 Schema::builder()
214 .with_schema_id(1)
215 .with_fields(vec![
216 crate::spec::NestedField::required(
217 2147483546,
218 "file_path",
219 crate::spec::Type::Primitive(crate::spec::PrimitiveType::String),
220 )
221 .into(),
222 crate::spec::NestedField::required(
223 2147483545,
224 "pos",
225 crate::spec::Type::Primitive(crate::spec::PrimitiveType::Long),
226 )
227 .into(),
228 ])
229 .build()
230 .unwrap(),
231 );
232
233 let task = FileScanTaskDeleteFile {
234 file_path: del_path.clone(),
235 file_size_in_bytes: std::fs::metadata(&del_path).unwrap().len(),
236 file_type: DataContentType::PositionDeletes,
237 file_format: DataFileFormat::Parquet,
238 partition_spec_id: 0,
239 equality_ids: None,
240 key_metadata: Some(Box::from(key_metadata.as_ref())),
241 referenced_data_file: None,
242 content_offset: None,
243 content_size_in_bytes: None,
244 record_count: None,
245 };
246
247 let scan_metrics = ScanMetrics::new();
248 let delete_file_loader = BasicDeleteFileLoader::new(file_io, scan_metrics);
249
250 let result = delete_file_loader
251 .read_delete_file(&task, schema)
252 .await
253 .unwrap();
254
255 let batches: Vec<_> = result.try_collect().await.unwrap();
256 assert_eq!(batches.len(), 1);
257 assert_eq!(batches[0].num_rows(), 4);
258 }
259
260 #[tokio::test]
261 async fn test_read_encrypted_equality_delete_file() {
262 use std::collections::HashMap;
263 use std::sync::Arc;
264
265 use arrow_array::{Int64Array, RecordBatch};
266 use parquet::arrow::PARQUET_FIELD_ID_META_KEY;
267
268 use crate::encryption::StandardKeyMetadata;
269 use crate::scan::FileScanTaskDeleteFile;
270 use crate::spec::{DataContentType, DataFileFormat};
271
272 let encryption_key = b"0123456789abcdef";
273 let aad_prefix = b"my-table-uuid!!";
274
275 let tmp_dir = TempDir::new().unwrap();
276 let table_location = tmp_dir.path().to_str().unwrap();
277 let file_io = FileIO::new_with_fs();
278
279 let arrow_schema = Arc::new(arrow_schema::Schema::new(vec![
280 arrow_schema::Field::new("id", arrow_schema::DataType::Int64, false).with_metadata(
281 HashMap::from([(PARQUET_FIELD_ID_META_KEY.to_string(), "1".to_string())]),
282 ),
283 ]));
284
285 let id_col = Arc::new(Int64Array::from(vec![100i64, 200, 300]));
286 let batch = RecordBatch::try_new(arrow_schema.clone(), vec![id_col]).unwrap();
287
288 let del_path = format!("{table_location}/encrypted-eq-del.parquet");
289 write_encrypted_parquet(&del_path, &batch, encryption_key, Some(aad_prefix));
290
291 let key_metadata = StandardKeyMetadata::try_new(encryption_key)
292 .unwrap()
293 .with_aad_prefix(aad_prefix)
294 .encode()
295 .unwrap();
296
297 let schema = Arc::new(
298 Schema::builder()
299 .with_schema_id(1)
300 .with_fields(vec![
301 crate::spec::NestedField::required(
302 1,
303 "id",
304 crate::spec::Type::Primitive(crate::spec::PrimitiveType::Long),
305 )
306 .into(),
307 ])
308 .build()
309 .unwrap(),
310 );
311
312 let task = FileScanTaskDeleteFile {
313 file_path: del_path.clone(),
314 file_size_in_bytes: std::fs::metadata(&del_path).unwrap().len(),
315 file_type: DataContentType::EqualityDeletes,
316 file_format: DataFileFormat::Parquet,
317 partition_spec_id: 0,
318 equality_ids: Some(vec![1]),
319 key_metadata: Some(Box::from(key_metadata.as_ref())),
320 referenced_data_file: None,
321 content_offset: None,
322 content_size_in_bytes: None,
323 record_count: None,
324 };
325
326 let scan_metrics = ScanMetrics::new();
327 let delete_file_loader = BasicDeleteFileLoader::new(file_io, scan_metrics);
328
329 let result = delete_file_loader
330 .read_delete_file(&task, schema)
331 .await
332 .unwrap();
333
334 let batches: Vec<_> = result.try_collect().await.unwrap();
335 assert_eq!(batches.len(), 1);
336 assert_eq!(batches[0].num_rows(), 3);
337 }
338}