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