Skip to main content

lance_table/io/
deletion.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use std::{collections::HashSet, sync::Arc};
5
6use arrow_array::{RecordBatch, UInt32Array};
7use arrow_ipc::CompressionType;
8use arrow_ipc::reader::FileReader as ArrowFileReader;
9use arrow_ipc::writer::{FileWriter as ArrowFileWriter, IpcWriteOptions};
10use arrow_schema::{ArrowError, DataType, Field, Schema};
11use bytes::Buf;
12use lance_core::error::{CorruptFileSnafu, box_error};
13use lance_core::utils::deletion::DeletionVector;
14use lance_core::utils::tracing::{AUDIT_MODE_CREATE, AUDIT_TYPE_DELETION, TRACE_FILE_AUDIT};
15use lance_core::{Error, Result};
16use lance_io::object_store::ObjectStore;
17use object_store::path::Path;
18use rand::Rng;
19use roaring::bitmap::RoaringBitmap;
20use snafu::ResultExt;
21use tracing::{info, instrument};
22
23use crate::format::{DeletionFile, DeletionFileType};
24
25pub const DELETIONS_DIR: &str = "_deletions";
26
27/// Get the Arrow schema for an Arrow deletion file.
28fn deletion_arrow_schema() -> Arc<Schema> {
29    Arc::new(Schema::new(vec![Field::new(
30        "row_id",
31        DataType::UInt32,
32        false,
33    )]))
34}
35
36/// Get the file path for a deletion file. This is relative to the dataset root.
37pub fn deletion_file_path(base: &Path, fragment_id: u64, deletion_file: &DeletionFile) -> Path {
38    let DeletionFile {
39        read_version,
40        id,
41        file_type,
42        ..
43    } = deletion_file;
44    let suffix = file_type.suffix();
45    base.child(DELETIONS_DIR)
46        .child(format!("{fragment_id}-{read_version}-{id}.{suffix}"))
47}
48
49pub fn relative_deletion_file_path(fragment_id: u64, deletion_file: &DeletionFile) -> String {
50    let DeletionFile {
51        read_version,
52        id,
53        file_type,
54        ..
55    } = deletion_file;
56    let suffix = file_type.suffix();
57    format!("{DELETIONS_DIR}/{fragment_id}-{read_version}-{id}.{suffix}")
58}
59
60/// Write a deletion file for a fragment for a given deletion vector.
61///
62/// Returns the deletion file if one was written. If no deletions were present,
63/// returns `Ok(None)`.
64pub async fn write_deletion_file(
65    base: &Path,
66    fragment_id: u64,
67    read_version: u64,
68    removed_rows: &DeletionVector,
69    object_store: &ObjectStore,
70) -> Result<Option<DeletionFile>> {
71    let deletion_file = match removed_rows {
72        DeletionVector::NoDeletions => None,
73        DeletionVector::Set(set) => {
74            let id = rand::rng().random::<u64>();
75            let deletion_file = DeletionFile {
76                read_version,
77                id,
78                file_type: DeletionFileType::Array,
79                num_deleted_rows: Some(set.len()),
80                base_id: None,
81            };
82            let path = deletion_file_path(base, fragment_id, &deletion_file);
83
84            let array = UInt32Array::from_iter(set.iter().copied());
85            let array = Arc::new(array);
86
87            let schema = deletion_arrow_schema();
88            let batch = RecordBatch::try_new(schema.clone(), vec![array])?;
89
90            let mut out: Vec<u8> = Vec::new();
91            let write_options =
92                IpcWriteOptions::default().try_with_compression(Some(CompressionType::ZSTD))?;
93            {
94                let mut writer = ArrowFileWriter::try_new_with_options(
95                    &mut out,
96                    schema.as_ref(),
97                    write_options,
98                )?;
99                writer.write(&batch)?;
100                writer.finish()?;
101                // Drop writer so out is no longer borrowed.
102            }
103
104            object_store.put(&path, &out).await?;
105
106            info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_DELETION, path = path.to_string());
107
108            Some(deletion_file)
109        }
110        DeletionVector::Bitmap(bitmap) => {
111            let id = rand::rng().random::<u64>();
112            let deletion_file = DeletionFile {
113                read_version,
114                id,
115                file_type: DeletionFileType::Bitmap,
116                num_deleted_rows: Some(bitmap.len() as usize),
117                base_id: None,
118            };
119            let path = deletion_file_path(base, fragment_id, &deletion_file);
120
121            let mut out: Vec<u8> = Vec::new();
122            bitmap.serialize_into(&mut out)?;
123
124            object_store.put(&path, &out).await?;
125
126            info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_DELETION, path = path.to_string());
127
128            Some(deletion_file)
129        }
130    };
131    Ok(deletion_file)
132}
133
134#[instrument(
135    level = "debug",
136    skip(base, object_store),
137    fields(
138        base = base.as_ref(),
139        bytes_read = tracing::field::Empty
140    )
141)]
142pub async fn read_deletion_file(
143    fragment_id: u64,
144    deletion_file: &DeletionFile,
145    base: &Path,
146    object_store: &ObjectStore,
147) -> Result<DeletionVector> {
148    let span = tracing::Span::current();
149    match deletion_file.file_type {
150        DeletionFileType::Array => {
151            let path = deletion_file_path(base, fragment_id, deletion_file);
152
153            let data = object_store.read_one_all(&path).await?;
154            span.record("bytes_read", data.len());
155            let data = std::io::Cursor::new(data);
156            let mut batches: Vec<RecordBatch> = ArrowFileReader::try_new(data, None)?
157                .collect::<std::result::Result<_, ArrowError>>()
158                .map_err(box_error)
159                .context(CorruptFileSnafu { path: path.clone() })?;
160
161            if batches.len() != 1 {
162                return Err(Error::corrupt_file(
163                    path,
164                    format!(
165                        "Expected exactly one batch in deletion file, got {}",
166                        batches.len()
167                    ),
168                ));
169            }
170
171            let batch = batches.pop().unwrap();
172            if batch.schema() != deletion_arrow_schema() {
173                return Err(Error::corrupt_file(
174                    path,
175                    format!(
176                        "Expected schema {:?} in deletion file, got {:?}",
177                        deletion_arrow_schema(),
178                        batch.schema()
179                    ),
180                ));
181            }
182
183            let array = batch.columns()[0]
184                .as_any()
185                .downcast_ref::<UInt32Array>()
186                .unwrap();
187
188            let mut set = HashSet::with_capacity(array.len());
189            for val in array.iter() {
190                if let Some(val) = val {
191                    set.insert(val);
192                } else {
193                    return Err(Error::corrupt_file(
194                        path,
195                        "Null values are not allowed in deletion files",
196                    ));
197                }
198            }
199
200            Ok(DeletionVector::Set(set))
201        }
202        DeletionFileType::Bitmap => {
203            let path = deletion_file_path(base, fragment_id, deletion_file);
204
205            let data = object_store.read_one_all(&path).await?;
206            span.record("bytes_read", data.len());
207            let reader = data.reader();
208            let bitmap = RoaringBitmap::deserialize_from(reader)
209                .map_err(box_error)
210                .context(CorruptFileSnafu { path })?;
211
212            Ok(DeletionVector::Bitmap(bitmap))
213        }
214    }
215}
216
217#[cfg(test)]
218mod test {
219
220    use super::*;
221
222    #[tokio::test]
223    async fn test_write_no_deletions() {
224        let dv = DeletionVector::NoDeletions;
225
226        let (object_store, path) = ObjectStore::from_uri("memory:///no_deletion")
227            .await
228            .unwrap();
229        let file = write_deletion_file(&path, 0, 0, &dv, &object_store)
230            .await
231            .unwrap();
232        assert!(file.is_none());
233    }
234
235    #[tokio::test]
236    async fn test_write_array() {
237        let dv = DeletionVector::Set(HashSet::from_iter(0..100));
238
239        let fragment_id = 21;
240        let read_version = 12;
241
242        let object_store = ObjectStore::memory();
243        let path = Path::from("/write");
244        let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
245            .await
246            .unwrap();
247
248        assert!(matches!(
249            file,
250            Some(DeletionFile {
251                file_type: DeletionFileType::Array,
252                ..
253            })
254        ));
255
256        let file = file.unwrap();
257        assert_eq!(file.read_version, read_version);
258        let path = deletion_file_path(&path, fragment_id, &file);
259        assert_eq!(
260            path,
261            Path::from(format!("/write/_deletions/21-12-{}.arrow", file.id))
262        );
263
264        let data = object_store
265            .inner
266            .get(&path)
267            .await
268            .unwrap()
269            .bytes()
270            .await
271            .unwrap();
272        let data = std::io::Cursor::new(data);
273        let mut batches: Vec<RecordBatch> = ArrowFileReader::try_new(data, None)
274            .unwrap()
275            .collect::<std::result::Result<_, ArrowError>>()
276            .unwrap();
277
278        assert_eq!(batches.len(), 1);
279        let batch = batches.pop().unwrap();
280        assert_eq!(batch.schema(), deletion_arrow_schema());
281        let array = batch["row_id"]
282            .as_any()
283            .downcast_ref::<UInt32Array>()
284            .unwrap();
285        let read_dv = DeletionVector::from_iter(array.iter().map(|v| v.unwrap()));
286        assert_eq!(read_dv, dv);
287    }
288
289    #[tokio::test]
290    async fn test_write_bitmap() {
291        let dv = DeletionVector::Bitmap(RoaringBitmap::from_iter(0..100));
292
293        let fragment_id = 21;
294        let read_version = 12;
295
296        let object_store = ObjectStore::memory();
297        let path = Path::from("/bitmap");
298        let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
299            .await
300            .unwrap();
301
302        assert!(matches!(
303            file,
304            Some(DeletionFile {
305                file_type: DeletionFileType::Bitmap,
306                ..
307            })
308        ));
309
310        let file = file.unwrap();
311        assert_eq!(file.read_version, read_version);
312        let path = deletion_file_path(&path, fragment_id, &file);
313        assert_eq!(
314            path,
315            Path::from(format!("/bitmap/_deletions/21-12-{}.bin", file.id))
316        );
317
318        let data = object_store
319            .inner
320            .get(&path)
321            .await
322            .unwrap()
323            .bytes()
324            .await
325            .unwrap();
326        let reader = data.reader();
327        let read_bitmap = RoaringBitmap::deserialize_from(reader).unwrap();
328        assert_eq!(read_bitmap, dv.into_iter().collect::<RoaringBitmap>());
329    }
330
331    #[tokio::test]
332    async fn test_roundtrip_array() {
333        let dv = DeletionVector::Set(HashSet::from_iter(0..100));
334
335        let fragment_id = 21;
336        let read_version = 12;
337
338        let object_store = ObjectStore::memory();
339        let path = Path::from("/roundtrip");
340        let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
341            .await
342            .unwrap();
343
344        let read_dv = read_deletion_file(fragment_id, &file.unwrap(), &path, &object_store)
345            .await
346            .unwrap();
347        assert_eq!(read_dv, dv);
348    }
349
350    #[tokio::test]
351    async fn test_roundtrip_bitmap() {
352        let dv = DeletionVector::Bitmap(RoaringBitmap::from_iter(0..100));
353
354        let fragment_id = 21;
355        let read_version = 12;
356
357        let object_store = ObjectStore::memory();
358        let path = Path::from("/bitmap");
359        let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
360            .await
361            .unwrap();
362
363        let read_dv = read_deletion_file(fragment_id, &file.unwrap(), &path, &object_store)
364            .await
365            .unwrap();
366        assert_eq!(read_dv, dv);
367    }
368}