Skip to main content

iceberg/arrow/
delete_file_loader.rs

1// Licensed to the Apache Software Foundation (ASF) under one
2// or more contributor license agreements.  See the NOTICE file
3// distributed with this work for additional information
4// regarding copyright ownership.  The ASF licenses this file
5// to you under the Apache License, Version 2.0 (the
6// "License"); you may not use this file except in compliance
7// with the License.  You may obtain a copy of the License at
8//
9//   http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing,
12// software distributed under the License is distributed on an
13// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14// KIND, either express or implied.  See the License for the
15// specific language governing permissions and limitations
16// under the License.
17
18use 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/// Delete File Loader
33#[allow(unused)]
34#[async_trait::async_trait]
35pub trait DeleteFileLoader {
36    /// Read the delete file referred to in the task
37    ///
38    /// Returns the contents of the delete file as a RecordBatch stream. Applies schema evolution.
39    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    /// Loads a RecordBatchStream for a given datafile.
66    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        /*
73           Essentially a super-cut-down ArrowReader. We can't use ArrowReader directly
74           as that introduces a circular dependency.
75        */
76        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    /// Evolves the schema of the RecordBatches from an equality delete file.
97    ///
98    /// Per the [Iceberg spec](https://iceberg.apache.org/spec/#equality-delete-files),
99    /// only evolves the specified `equality_ids` columns, not all table columns.
100    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        // For equality deletes, only evolve the equality_ids columns.
134        // For positional deletes (equality_ids is None), use all field IDs.
135        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}